stacksaga-kafka-worker — 端点设计手册 (Endpoints)

stacksaga-kafka-worker 是为 Worker 微服务提供命令执行能力的核心模块,用于执行编排器通过常规主题发送的业务命令。它负责消费来自常规主题的命令消息、执行具体业务逻辑,并通过编排器托管的基于领域的专用回调主题将应答结果回传给编排器。

在高层架构中,stacksaga-kafka-worker 承担以下核心职责:

  1. 在目标客户端服务中,按照 StackSaga Kafka 引擎规范提供创建 Kafka Saga 端点的核心抽象。

  2. 在每次端点调用完成后,自动向编排器服务派发应答消息。

  3. 针对可重试的临时异常,提供开箱即用的进程内即时重试能力。

端点主题 (Endpoint Topics)

每个 Saga 端点均从属于*托管该端点的 Worker 服务*的 Kafka 主题中消费命令消息。 该主题是*在 Worker 端通过 @SagaEndpoint 的 topicNameSuffix 属性声明并定义的* — 这是该端点主题的权威法定声明。 编排器在其 AbstractTopic 类中声明*相同*的主题名称,纯粹是为了在路由命令消息时进行*网络寻址*;编排器向该主题发布消息,但既不拥有也不消费它。

端点主题命名规范

因为主题是 Worker 服务的业务能力(且完全相同的端点可被多个不同的 Saga 领域复用),所以请以其*针对的服务和操作动作*命名,采用 {service-name}.{action} 形式(例如 user-service.fetch-user-details、payment-service.make-payment),而非以任何单一 Saga 领域命名。 例如,一个 user-service.fetch-user-details 查询端点可以同时被下单 Saga、取消订阅 Saga 以及其他业务流程复用。因此,基于服务限定的命名具有以下优势:

  • 保持主题名称*全局唯一* — 所有端点主题共享统一的 Kafka 命名空间,服务限定命名可彻底避免不同微服务之间的端点主题冲突;

  • 使主题归属*清晰自解释* — 一目了然哪个服务拥有并消费该主题;

  • 支持*跨 Saga 领域高度复用* — 每个业务领域均通过完全相同的标准主题名称引用该端点。

命名规则: 使用英文句点 . 作为分层分段分隔符;连字符 - 允许在分段*内部*使用(例如 user-service、fetch-user-details),但不能*作为*分隔符本身;严禁在任何位置使用下划线 _。

saga. 前缀机制

一个端点消费单个 Kafka 主题。对于 CommandEndpoint,主执行命令与补偿命令均在同一个主题中接收,并分别对应执行 doProcess() 或 undoProcess()。

框架会自动为主题名称追加 saga. 前缀:

  • 如果您的 topicNameSuffix 未以 saga. 开头,框架会自动在前端追加 — 即 topicNameSuffix = "payment-service.make-payment" 将生成真实的 Kafka 主题 saga.payment-service.make-payment。

  • 如果您在声明时已自行包含了 saga. 前缀(例如 topicNameSuffix = "saga.payment-service.make-payment"),框架会原样保留,不会重复追加。

无论采用哪种方式,实际生成的 Kafka 主题名称格式必须为 saga.<name>,其中 <name> 全小写,且分段之间使用 . 或 - 连接(禁止大写字母与下划线)。

编排器端在 AbstractTopic 中为相同跨度定义的主题应用了*完全相同*的 saga. 前缀规则,因此两端必须使用*完全一致*的主题后缀 — 编排器与 Worker 必须精准解析到完全相同的真实 Kafka 主题。 详见 编排器端主题声明 与 主题键规范。
StackSaga 不会编程式自动创建这些主题。是否在首次使用时自动建表取决于 Broker 配置 (auto.create.topics.enable);生产环境下推荐在启动应用前在 Kafka 集群中显式预先创建好主题。详见 端点主题监听器模型。

StackSaga Kafka 端点类型 (Endpoint Types)

StackSaga Kafka 端点代表按照 StackSaga 引擎规范构建的执行节点。根据业务属性,可以创建两类端点:

异常处理参考手册 (Exception Handling Reference)

在端点方法内部选择正确的异常类型,是 StackSaga-Kafka 中最具决定性的开发考量之一。 抛出错误的异常类型可能导致 Saga 无限期停滞、触发不必要的全局补偿回滚,或是丢失补偿所需的珍贵故障上下文。

异常类型速查表

异常类型 适用方法 框架响应机制 使用场景

NonRetryableExecutorException

doProcess(), undoProcess()

在 doProcess() 中:将 Saga 状态流转为 FAILED,并立即触发逆向补偿序列。
在 undoProcess() 中:立即终止补偿流程。

故障为永久性不可恢复 — 业务规则校验不通过、输入参数非法或任何重试绝不可能解决的错误。 使用 put() 附加用于补偿流程的错误元数据。

RetryableExecutorException

doProcess(), undoProcess()

通过环形协调器将 Saga 调度为延期异步重试。 该步骤将从事件存储库中最后持久化的状态重新执行。

故障为暂时性瞬态波动 — 网络超时、下游依赖临时不可用或预计可自我愈合的基础设施抖动。

JustRetryableExecutorException

doProcess(), undoProcess() (仅限非响应式)

通过 Spring RetryTemplate 在当前 Worker 进程内立即触发重试,不向编排器发送通知。 当即时重试次数耗尽后,回退抛出 RetryableExecutorException 以触发延期调度重试。

在升级为全局调度重试之前,值得在本地立即快速重试的场景。 始终作为非响应式端点防御瞬态故障的第一道防线。

任何未捕获的常规异常

doProcess(), undoProcess(), onNext(), onNextRevert()

被框架等同于 NonRetryableExecutorException 处理。

在端点中绝非刻意为之 — 参见下文警告。 而在 onNext() 中,刻意抛出异常是被官方支持的强制回滚模式。

切勿任由未捕获的常规异常从端点向外逃逸。 任何未包装在框架专用异常中的 Throwable 都会被当作不可重试致命错误。 补偿序列将被迫启动(如果在 undoProcess() 中则直接终止),且不携带任何错误上下文,因为未捕获异常无法通过 put() 传递键值对供补偿步骤使用。 务必在端点内显式捕获底层异常并包装为合规的框架异常类型。

异常决策树 (Exception Decision Flow)

在实现端点方法时,依据以下逻辑判定抛出何种异常:

  1. 故障是否为瞬态网络/资源波动?(超时、临时资源不可用、基础设施抖动)

    1. 是 — 且为非响应式端点? → 优先抛出 JustRetryableExecutorException。若即时重试上限耗尽 → 抛出 RetryableExecutorException。

    2. 是 — 且为响应式端点? → 在响应式管道中使用 retryWhen()。若重试耗尽 → 返回 Mono.error(RetryableExecutorException.of(…​))。

  2. 是否为永久性业务校验失败?(无效数据、业务冲突、不可逆错误) → 抛出 NonRetryableExecutorException。使用 put() 附加供补偿步骤使用的上下文。

  3. 属于上述之外的其他非预期异常? → 包装为包含清晰描述信息的 NonRetryableExecutorException。绝不允许作为裸异常向外冒泡。

在 onNext() 与 onNextRevert() 中抛出异常的行为

异常处理不仅限于端点方法。EventManager 中的 onNext() 与 onNextRevert() 同样参与 Saga 故障生命周期:

  • onNext() — 在此处抛出的任何异常(或通过 actionUtil.error(…​) 返回的异常)均被视为不可重试故障。 框架立即将 Saga 状态流转为 FAILED 并启动逆向补偿序列。 这意味着您可以*有意识地*利用 onNext() 基于编排器端业务规则强制发起补偿 — 例如在收到 Worker 应答后,规则校验判定该订单不应继续推进。

  • onNextRevert() — 在此处抛出的任何异常都会导致补偿序列立即终止。 事务被标记为补偿失败 (Compensation Failed) 且不再重试。

抛出 / 返回自 即时效应 产生的状态流转

doProcess() — NonRetryableExecutorException 或未捕获异常

启动补偿回滚

IN_PROGRESS → FAILED → COMPENSATING

onNext() — 任何异常(或 actionUtil.error(…​))

启动补偿回滚

IN_PROGRESS → FAILED → COMPENSATING

undoProcess() — 未捕获异常

终止补偿流程

COMPENSATING → 补偿失败 (Compensation Failed)

onNextRevert() — 任何异常

终止补偿流程

COMPENSATING → 补偿失败 (Compensation Failed)

查询端点 (Query Endpoint)

如果某个原子执行在执行期间不对任何数据库造成状态修改,此类执行应在*查询端点 (Query Endpoint)* 中实现。由于它不改变数据库状态,因此无需任何逆向补偿操作撤销。
这些主题在 StackSaga 主题 中定义为 SagaEventType.QUERY_DO_ACTION 类型。

例如在下单示例中,*获取用户详情并校验*就是典型的查询执行,因为它不修改数据库,仅查询数据并返回。

查询端点在 StackSaga 中支持两种编程风格:

非响应式(命令式)查询端点

非响应式查询端点采用传统且直观的阻塞式编写模式:

(2)
@SagaEndpoint(
        //真实 Kafka 主题为 saga.user-service.get-user
        topicNameSuffix = "user-service.get-user", (3)
        listenerScope = WorkerListenerScope.SHARED_GLOBAL, (4)
        groupType = GroupType.SHARE
)
public class UserValidateEndpoint extends QueryEndpoint { (1)

    (5)
    @Override
    public void doProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException, NonRetryableExecutorException {

        log.info("Message Key (Transaction Id): {}", consumerRecord.key()); (6)
        log.info("IdempotencyKey: {}", consumerRecord.value().getIdempotencyKey()); (7)

        consumerRecord.value().getCurrentDomainEntityStateForUpdate().ifPresent(currentDomainEntityState -> { (8)
            log.info("Received payload for before user validation: {}", currentDomainEntityState);
            try {
                {
                    //用户校验业务逻辑 (9)
                    String username = currentDomainEntityState.get("username").asText();
                    log.info("Validating user with username: {}", username);
                }
                {(10)
                    // 将处理结果存入领域实体状态,供后续主题消费
                    currentDomainEntityState.put("is_user_validated", true);
                    currentDomainEntityState.put("validation_note", "user validated at " + LocalDateTime.now());
                }
            } catch (SomeRetryableException e) {
                (11)
                throw RetryableExecutorException.of(e);
            } catch (SomeNonRetryableException e) {
                (12)
                throw NonRetryableExecutorException
                        .buildWith(e)
                        .put("error_code", "USER_VALIDATION_FAILED")
                        .put("key-1", "value-1")
                        .put("key-2", "value-2")
                        .build();
            }
        });
    }
}
1 继承 QueryEndpoint 类以创建自定义查询端点。
2 在类上标注 @SagaEndpoint 注解,声明其为 Spring Bean 与 Saga 执行节点。
3 topicNameSuffix:指定该端点的主题名称后缀。框架通过在该后缀前加上 saga. 前缀构建实际的 Kafka 主题 — 即 user-service.get-user 解析为真实主题 saga.user-service.get-user。(若已包含 saga. 前缀则原样保留)。EventManager 正是向该主题派发命令以触发本端点。
4 listenerScope 与 groupType:配置该端点主题的消费容器。
listenerScope 选择容器分配策略 — WorkerListenerScope.SHARED_GLOBAL(汇入单个全局容器,默认之选)、WorkerListenerScope.SHARED_GROUP(按 containerId 分组共享容器)或 WorkerListenerScope.ISOLATED(独立专属容器)。
groupType 选择 Kafka 消费组协议 — GroupType.CONSUMER(传统消费者组)或 GroupType.SHARE(Kafka 4 共享组 / KIP-932,并行度不受分区数限制);参见 stacksaga-kafka-implementation/worker/worker-configuration.adoc#worker_group_type。
5 重写 doProcess() 方法以实现执行查询的业务逻辑。当对应主题收到消息时框架调用此方法。可通过传入的 ConsumerRecord 在方法内获取消息键、载荷及其他上下文元数据。
6 以字符串格式从 ConsumerRecord 中获取消息键。该键即为编排器随消息传递的全局事务 ID。
将事务 ID 设为消息键的核心原因在于保证同一事务的所有消息均被哈希派发到 Kafka 的同一个分区中,确保单个事务内消息处理的严格顺序性。
7 从 ConsumerRecord 的载荷中提取幂等键 (idempotencyKey)。该键由编排器端框架生成并随消息传输,确保即使因网络重试或重复投递而多次处理同一消息,也始终携带完全一致的幂等键,以便在端点内部实现安全幂等控制。
8 通过载荷的 getCurrentDomainEntityStateForUpdate() 方法获取领域实体的当前状态。该方法返回一个包含当前状态的 Optional。在主执行流程中领域实体始终存在。
在目标 Worker 服务端,领域实体以 ObjectNode(Jackson 提供的 JSON 树形结构)的形式提供,而非原生的 Java 实体对象(因为原生的强类型 CustomDomainEntity 仅存在于编排器端)。可使用 ObjectNode 的 get("fieldName") 提取属性,并使用 put("fieldName", value) 追加或修改属性。
9 执行具体的查询与校验业务逻辑。
10 处理完毕后,将业务计算结果追加回领域实体中,供后续主题读取。
Worker 端向领域实体添加的新字段,应事先在编排器端的 CustomDomainEntity 中完成定义。若在 Worker 端添加了编排器端未声明的字段,该字段会自动归档至领域实体的 missingProperties 映射中,编排器端可通过 getMissingProperties() 访问它。
如架构所述,若从方法中抛出了任何异常(RetryableExecutorException 或 NonRetryableExecutorException),方法内部对领域实体状态所做的任何修改都不会被应用。因为重试时必须基于原始输入状态重新调用;而在不可重试失败时,系统将开启逆向补偿,状态回退到上一次成功快照。若需传递故障相关的上下文数据,必须使用 NonRetryableExecutorException.put() 以键值对形式附加,供补偿阶段提取。
11 正确包装并抛出异常:若为可重试瞬态异常,使用 RetryableExecutorException.of(e) 包装抛出;若为不可重试业务失败,包装为 NonRetryableExecutorException。
12 抛出 NonRetryableExecutorException 时,可通过 put() 链式附加丰富的上下文元数据。这些元数据会随消息序列化传递回编排器,并在补偿执行期间供相关步骤使用。

响应式(非阻塞)查询端点

@Slf4j
@SagaEndpoint(
        //真实主题名称为 saga.user-service.get-user
        topicNameSuffix = "user-service.get-user",
        listenerScope = WorkerListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE
)
public class ReactiveUserValidateEndpoint extends ReactiveQueryEndpoint {(1)

    (2)
    @Override
    public Mono<Void> doProcess(ConsumerRecord<String, SagaPayload> consumerRecord) {
        log.info("Message Key (Transaction Id): {}", consumerRecord.key());
        log.info("IdempotencyKey: {}", consumerRecord.value().getIdempotencyKey());
        final ObjectNode currentDomainEntityState = consumerRecord
                .value()
                .getCurrentDomainEntityStateForUpdate()
                .orElseThrow();
        log.info("Received payload for before user validation: {}", currentDomainEntityState);
        String username = currentDomainEntityState.get("username").asText();
        return this.internalUserService
                .getUserDetails(username)
                .flatMap(userDetails -> {
                    {
                        //用户校验业务逻辑
                    }
                    {
                        // 将结果存入领域实体状态供后续主题使用
                        currentDomainEntityState.put("is_user_validated", true);
                        currentDomainEntityState.put("validation_note", "user validated at " + LocalDateTime.now());
                    }
                    return Mono.empty();
                })
                .onErrorResume(SomeRetryableException.class, e -> {
                    return Mono.error(RetryableExecutorException.of(e));
                })
                .onErrorResume(SomeNonRetryableException.class, e -> {
                    return Mono.error(NonRetryableExecutorException
                            .buildWith(e)
                            .put("error_code", "USER_VALIDATION_FAILED")
                            .build());
                });
    }
}

在响应式查询端点中,doProcess() 方法返回 Mono<Void>,并利用 Project Reactor 的响应式管道实现非阻塞调用。在 onErrorResume 中将底层异常转换为 Mono.error(RetryableExecutorException…​) 或 Mono.error(NonRetryableExecutorException…​)。

命令端点 (Command Endpoint)

如果某个原子执行会在数据库中引起状态变更,此类执行必须在*命令端点 (Command Endpoint)* 中实现。由于它更改了数据库状态,因此必须提供逆向补偿操作以撤销变更。
这些端点对应在 StackSaga 主题 中声明的 SagaEventType.COMMAND_DO_ACTION 与 SagaEventType.COMMAND_UNDO_ACTION 主题。

例如在下单示例中,*执行支付扣款*就是命令执行,因为它扣减了用户余额并创建了支付凭证。若后续步骤失败,必须触发补偿动作执行退款。

在命令端点中,需要实现两个核心方法:

  • doProcess():负责执行正向主业务命令。

  • undoProcess():负责执行逆向补偿逻辑,撤销 doProcess() 所造成的数据库副作用。

命令端点同样支持两种模式:

非响应式(命令式)命令端点

@Slf4j
@SagaEndpoint(
        topicNameSuffix = "payment-service.make-payment",
        listenerScope = WorkerListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE
)
class MakePaymentEndpoint extends CommandEndpoint { (1)

    (2)
    @Override
    public void doProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException, NonRetryableExecutorException {
        //正向处理原理与查询端点完全一致
    }

    (2)
    @Override
    public void undoProcess(ConsumerRecord<String, SagaPayload> consumerRecord) throws JustRetryableExecutorException, RetryableExecutorException {
        log.info("Message Key (Transaction Id): {}", consumerRecord.key());
        log.info("IdempotencyKey: {}", consumerRecord.value().getIdempotencyKey());
        (3)
        final JsonNode lastDomainEntityState = consumerRecord.value().getDomainEntityState();
        log.info("Last domain entity state {}", lastDomainEntityState);

        try {
            {//退款补偿业务逻辑
                final double totalAmount = lastDomainEntityState.get("total_amount").asDouble();
                final String paymentRef = lastDomainEntityState.get("payment_reference_id").asText();
                this.internalPaymentSerivce.refundPayment(paymentRef, totalAmount);
                //若需要,为后续补偿执行更新提示信息
                consumerRecord.value().getHintStore().ifPresent(hintStore -> {(4)
                    hintStore.put("last_refund_time", LocalDateTime.now().toString());
                });
            }
        } catch (SomeRetryableException e) {(5)
            throw RetryableExecutorException.of(e);
        }
    }
}
1 继承 CommandEndpoint 类创建自定义命令端点。
2 分别重写 doProcess() 实现正向命令,重写 undoProcess() 实现逆向补偿回滚。
3 undoProcess() 中 ConsumerRecord 载荷包含正向命令执行前已知的领域实体快照(可通过 getDomainEntityState() 访问)。此状态在补偿期间为严格只读不可修改。
4 若需向后续补偿步骤传递元数据,可通过 getHintStore() 访问提示存储库并通过 put() 存入键值对。
提示存储库中的更新仅在 undoProcess() 正常成功返回时才会被持久化并传递给后续步骤。若抛出重试异常或非预期异常,当前所写的提示不会生效。
5 识别补偿逻辑中可能出现的瞬态异常,并使用 RetryableExecutorException.of(e) 包装抛出以启动重试。在 undoProcess() 中仅允许抛出 JustRetryableExecutorException 与 RetryableExecutorException。抛出任何未捕获异常均会导致全局补偿流程彻底终止。

基于 JustRetryableExecutorException 的进程内即时重试

基于 JustRetryableExecutorException 的即时重试专用于非响应式端点(响应式端点可在管道中直接使用 retryWhen())。 抛出 JustRetryableExecutorException 会立即在当前线程通过 Spring RetryTemplate 触发本地即时重试,无需等待调度器唤醒,亦不向编排器发送冗余网络消息。当即时重试耗尽后,再回退抛出 RetryableExecutorException 以触发延期调度重试。

@Slf4j
@SagaEndpoint(
        topicNameSuffix = "payment-service.make-payment",
        listenerScope = WorkerListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE
)
public class MakePaymentEndpoint extends CommandEndpoint {

    @Override
    public void doProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException, NonRetryableExecutorException {
        consumerRecord
                .value()
                .getCurrentDomainEntityStateForUpdate()
                .ifPresent(currentDomainEntityState -> {
                    final String paymentReferenceId = currentDomainEntityState.get("payment_reference_id").asText();
                    final double totalAmount = currentDomainEntityState.get("total_amount").asDouble();
                    try {
                        this.internalPaymentService.makePayment(paymentReferenceId, totalAmount);
                    } catch (SomeRetryableException e) {
                        RetryContext retryContext = consumerRecord.value().getRetryContext().orElseThrow();
                        if (retryContext.getRetryCount() >= 3) {
                            //达到重试阈值,抛出不可重试异常终止并触发全局回滚
                            throw NonRetryableExecutorException
                                    .buildWith(e)
                                    .put("error_code", "PAYMENT_FAILED_AFTER_RETRIES")
                                    .build();
                        } else {
                            throw JustRetryableExecutorException.of("Payment failed, retrying... attempt " + (retryContext.getRetryCount() + 1));
                        }
                    } catch (SomeNonRetryableException e) {
                        throw NonRetryableExecutorException
                                .buildWith(e)
                                .put("error_code", "PAYMENT_FAILED")
                                .build();
                    }
                });
    }

    @Override
    public void undoProcess(ConsumerRecord<String, SagaPayload> consumerRecord) throws
            JustRetryableExecutorException, RetryableExecutorException {
    }
}

为非响应式端点配置 RetryTemplate

非响应式端点利用 RetryTemplate 处理本地重试。您可以通过定义继承自 AbstractRetryTemplateProvider 的 Spring Bean 来自定义该模板:

@Component
public class MakePaymentRetryTemplateProvider extends AbstractRetryTemplateProvider {
    @Override
    protected RetryTemplateBuilder retryTemplateBuilder() {
        SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy();
        simpleRetryPolicy.setMaxAttempts(3);
        return RetryTemplate.builder().customPolicy(simpleRetryPolicy);
    }
}

在端点上通过 Bean 名称引用该 Provider:

@Slf4j
@SagaEndpoint(
        topicNameSuffix = "payment-service.make-payment",
        listenerScope = WorkerListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE,
        primaryExecutionRetryTemplate = "makePaymentRetryTemplateProvider", (1)
        revertExecutionRetryTemplate = "makePaymentRetryTemplateProvider" (2)
)
public class MakePaymentEndpoint extends CommandEndpoint {
    //...
}

若省略这些属性,默认使用框架内置的全局共享 Provider(通过 stacksaga.kafka.worker.retry.* 统一调优)。

响应式(非阻塞)命令端点

@Slf4j
@SagaEndpoint( (2)
        topicNameSuffix = "payment-service.make-payment",
        listenerScope = WorkerListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE
)
class ReactiveMakePaymentEndpoint extends ReactiveCommandEndpoint { (1)

    @Override
    public Mono<Void> doProcess(ConsumerRecord<String, SagaPayload> consumerRecord) {
        //正向处理逻辑与响应式查询端点一致
    }

    @Override
    public Mono<Void> undoProcess(ConsumerRecord<String, SagaPayload> consumerRecord) {
        log.info("Message Key (Transaction Id): {}", consumerRecord.key());
        log.info("IdempotencyKey: {}", consumerRecord.value().getIdempotencyKey());
        final JsonNode lastDomainEntityState = consumerRecord.value().getDomainEntityState();
        log.info("Last domain entity state {}", lastDomainEntityState);
        final double totalAmount = lastDomainEntityState.get("total_amount").asDouble();
        final String paymentRef = lastDomainEntityState.get("payment_reference_id").asText();
        return this.internalPaymentSerivce
                .refundPayment(paymentRef, totalAmount)
                .flatMap(response -> {
                        //若需要,更新补偿提示信息
                        consumerRecord.value().getHintStore().ifPresent(hintStore -> {
                            hintStore.put("last_refund_time", LocalDateTime.now().toString());
                        });
                        return Mono.empty();
                })
                .onErrorResume(SomeRetryableException.class, e -> {(3)
                    return Mono.error(RetryableExecutorException.of(e));
                });
    }
}
1 继承 ReactiveCommandEndpoint。
2 标注 @SagaEndpoint 注解。
3 在 onErrorResume 中捕获重试异常并转换为 Mono.error(RetryableExecutorException.of(e))。

@SagaEndpoint 注解属性参考手册

@SagaEndpoint 是在 StackSaga Kafka Worker 应用中标记自定义端点类的核心元注解。其属性参考如下:

  1. value

    • 描述:指定端点的全局唯一标识名称(选填)。若不指定,默认使用类名生成 Spring Bean。

  2. topicNameSuffix

    • 描述:端点的主题名称后缀。框架通过在该后缀前加上 saga. 构建实际的 Kafka 主题。必须与编排器 AbstractTopic 中声明的名称完全匹配。

      Table 1. 主题命名规范速查表
      规则 详细说明

      命名约定(最佳实践)

      按照 {service-name}.{action} 形式命名(如 user-service.fetch-user-details),切勿以 Saga 领域命名,以便跨领域复用。

      分段分隔符

      使用英文句点 . 分隔各层级段(如 order-service.initialize-order)。

      连字符

      允许在分段内部使用(如 user-service),但不能充当分段分隔符。

      下划线

      严禁在主题名称中的任何位置使用。

      自动追加前缀

      框架会自动在前端追加 saga. 生成真实的实际主题(前提是后缀未以 saga. 开头)。例如 payment-service.make-payment 生成 saga.payment-service.make-payment。

      主题键

      每个主题在领域内必须拥有唯一的 float 键。正整数(1, 2, 3…)用于主执行;对应的负值(-1, -2, -3…)用于其补偿。小数保留给未来子执行特性,暂请勿用。

  3. listenerScope

  4. groupType

  5. containerId

    • 描述:仅在 listenerScope = WorkerListenerScope.SHARED_GROUP 时生效。作为将多个端点绑定到同一个共享容器的分组标识键。其他作用域下被忽略。

  6. concurrency

    • 描述:端点监听器容器的并发度(默认 1)。 在 ISOLATED 模式下直接设定该专用容器的并发数;在 SHARED_GROUP 模式下,容器采用该组所有成员中声明的*最大并发度*。对于 GroupType.CONSUMER,有效并行度受限于分区数;对于 GroupType.SHARE 则不受限制。

  7. primaryExecutionRetryTemplate

    • 描述:用于非响应式 CommandEndpoint 与 QueryEndpoint 主执行 (doProcess()) 进程内即时重试的 AbstractRetryTemplateProvider 的 Spring Bean 名称(字符串)。默认为 "sharedDefaultPrimaryAbstractRetryTemplateProvider"。

  8. revertExecutionRetryTemplate

    • 描述:用于非响应式 CommandEndpoint 补偿执行 (undoProcess()) 进程内即时重试的 AbstractRetryTemplateProvider 的 Spring Bean 名称(字符串)。默认为 "sharedDefaultRevertAbstractRetryTemplateProvider"。