快速入门示例:StackSaga-Kafka:MySQL

本快速入门示例将引导您在 StackSaga-Kafka(异步)编排引擎上构建一个完整的*下单 (order-placement)* Saga,并使用 MySQL 作为事件存储 (Event Store)。

与 同步快速入门示例(其中每个步骤都在单个编排器进程内运行)不同,在 Kafka 架构下,编排器完全通过 Kafka 主题与独立的*工作节点服务 (Worker Services)* 进行异步通信,因此业务执行分布在多个独立的 Spring Boot 应用程序之间。 通过本指南,您将构建一个跨三个服务运行的分布式多微服务 Saga,并在后续步骤发生失败时自动*补偿 (Compensate / 回滚)* 此前已完成的步骤。

概述 (Overview)

我们编排的业务场景是一个简化的*下单 (place-order)* 操作,建模为由三个原子步骤(跨度 / Spans)组成的单个长事务 (Long-Running Transaction, LRT)。 由于这是 Kafka(异步)引擎,每个跨度都由编排器通过 Kafka 调用对应的*工作节点服务 (Worker Service)* 执行——绝不通过直接的进程内调用或阻塞式 HTTP 请求。

# 跨度 (步骤) 动作类型 工作节点服务 补偿对应项

1

验证用户

查询 (Query / 只读)

user-service

— (查询操作从不补偿)

2

预留订单

命令 (Command / 状态变更)

order-service (编排器自身)

UNDO_ORDER_RESERVED

3

执行支付

命令 (Command / 状态变更)

payment-service

UNDO_MAKE_PAYMENT

如果任何步骤失败,StackSaga 会自动按倒序*补偿 (Compensate)* 此前已完成的*命令*步骤——例如在支付失败时释放已预留的订单——从而确保系统最终始终处于一致状态。

服务及其职责角色 (Services and Their Roles)

此示例涵盖三个 Spring Boot 应用程序。 编排器/工作节点标签描述的是服务*在当前 Saga 事务流中的参与角色*,而非物理部署形态——有关完整说明,请参阅 服务分类 (Service Classification)。

服务 StackSaga 角色 核心依赖项

order-service

编排器 (Orchestrator)(同时兼任“预留订单”跨度的*工作节点 (Worker)*)

stacksaga-kafka-orchestrator-spring-boot-starter + stacksaga-mysql-reactive-support

user-service

工作节点 (Worker)

stacksaga-kafka-worker-spring-boot-starter

payment-service

工作节点 (Worker)

stacksaga-kafka-worker-spring-boot-starter

order-service 同时扮演双重角色——它既*驱动*整个 Saga(编排器),又*执行*“预留订单”跨度(工作节点)。您无需向其单独引入 worker starter:orchestrator starter 已经通过传递依赖包含了 worker 的全部能力。只有编排器拥有 MySQL 事件存储 (Event Store);工作节点服务相对于 Saga 而言是无状态的。

Saga 执行流程 (How the Saga Executes)

  1. 客户端向*编排器* (order-service) 发起 POST /api/v1/order 请求。控制器初始化 OrderDomainEntity 并通过 StackSagaKafkaTemplate 将其移交给引擎。

  2. 引擎向 PlaceOrderEventManager 请求首个主题 (DO_USER_VALIDATED),并将*命令消息 (Command Message)* 发布至该目标工作节点的 Kafka 主题。

  3. 目标*工作节点* (user-service) 通过其 @SagaEndpoint 消费该命令,执行本地业务逻辑,并将*应答 (Reply)* 发布回框架统一管理的领域专属回调主题。

  4. 编排器消费该应答并调用 EventManager.onNext(…​) 判定下一步走向——依次路由至 DO_ORDER_RESERVED、DO_MAKE_PAYMENT,最后调用 complete()。

  5. 每次状态流转都会持久化至 MySQL 事件存储,因此运行中的在途事务随时均可恢复。

  6. 如果工作节点上报不可重试的失败(或 onNext() 返回错误),引擎将切换至补偿模式,并按相反顺序向相关工作节点分发 UNDO_* 撤销命令。

有关整体架构全貌,请参阅 StackSaga-Kafka 架构。

本指南涵盖的内容 (What This Guide Covers)

本指南实现了*最小可用拓扑*:编排器、三个工作节点跨度、故障补偿以及 MySQL 持久化。 一旦运行就绪,您可以在*无需修改以下任何业务代码*的前提下,叠加两层生产加固能力:

阶段一 [最小实现] (Stage-1 [Minimal Implementation])

本阶段演示以 MySQL 作为事件存储的 StackSaga Kafka 编排的最小端到端实现。我们将按以下顺序进行构建:

  1. 前置准备

  2. 创建订单服务(编排器服务) —— 编排器:领域实体、主题定义、EventManager、处理器、控制器以及配置。

  3. 配置用户服务作为工作节点 —— 第一个工作节点服务。

  4. 配置订单服务作为工作节点 —— 编排器同时承担“预留订单”跨度的工作节点角色。

  5. 配置支付服务作为工作节点 —— 第二个外部工作节点服务。

  6. 运行应用并测试事务流程 —— 启动所有组件并观察正常流与补偿流的实际运行。

前置准备 (Prerequisites)

在开始之前,请确保已安装并运行以下组件:

  1. Java 21+

  2. Spring Boot 4.x

  3. Apache Kafka 4.x —— 启用了共享消费组 (Share Groups) 的运行中 Broker(编排器与工作节点之间完全通过 Kafka 通信;默认 groupType = SHARE 使用 Kafka 4 共享消费组协议 / KIP-932)。

  4. MySQL 8+ —— 由编排器用作事件存储 (Event Store)。

  5. Maven

启动本地 Broker 和数据库的最快捷方式是 Docker——例如针对 Kafka 运行 bitnami/kafka(或 Redpanda)容器,针对 MySQL 运行 mysql:8 容器。

创建订单服务(编排器服务)(Creating Order Service (Orchestrator service))

第一步,我们将创建订单服务(编排器服务),负责处理订单提交并编排其他微服务。

起步 - 初始项目搭建 (Getting Started - Initial project setup)

在本示例中,我们使用 Spring MVC 作为 Web 框架,并使用 MySQL 作为事件存储的主数据库。

首先创建一个新的 Spring Boot 应用程序,并在 pom.xml 中添加以下依赖项与插件配置:

建议使用 StackSaga Initializer 获取相关依赖配置代码,以确保版本兼容性与初始配置正确无误。
<dependencyManagement> (1)
    <dependencies>
        <dependency>
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>
<dependencies>
    <!-- Spring Boot Starter Web -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>
    <dependency> (2)
        <groupId>org.stacksaga</groupId>
        <artifactId>stacksaga-kafka-orchestrator-spring-boot-starter</artifactId>
    </dependency>
    <dependency> (3)
        <groupId>org.stacksaga</groupId>
        <artifactId>stacksaga-mysql-reactive-support</artifactId>
    </dependency>
</dependencies>
<build>
    <plugins>
        <plugin> (4)
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-maven-plugin</artifactId>
            <executions>
                <execution>
                    <goals>
                        <goal>build-info</goal>
                    </goals>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>
1 此部分导入 StackSaga BOM (物料清单),统一管理所有 StackSaga 依赖项的版本,确保兼容性并简化版本控制。
2 这是核心 StackSaga Kafka 编排器 Starter 依赖项,包含编排引擎、事件处理和事务管理等核心功能,提供在应用中实现 Saga 编排所需的所有组件。
3 该依赖项提供以响应式 (Reactive) 方式将 MySQL 作为事件存储的支持。它包含将 MySQL 与 StackSaga 集成所需的组件和配置,允许您在 MySQL 数据库中存储和管理事件,同时利用响应式编程范式。StackSaga MySQL 响应式支持 (Reactive Support)
4 为 spring-boot-maven-plugin 配置 build-info 目标会在构建时生成 META-INF/build-info.properties。StackSaga 利用此元数据记录每个事务启动时所处的确切应用程序版本。
添加 build-info 目标有助于追踪事务是在哪个应用版本下发起的。若未配置该插件目标,StackSaga 将无法识别应用程序版本,并在事件存储中将事务的应用版本记录为 unknown。

创建领域实体 (Creating the Domain-Entity)

领域实体 (Domain-Entity) 是事务流程中流转的核心实体,它代表在整个事务期间被操作和处理的主数据结构。 在 StackSaga-Kafka 中,它也是通过 Kafka 在编排器与工作节点之间网络传输(序列化)的载荷数据 (Payload)。了解更多

@Getter
@Setter
@SagaDomainEntity(
        version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
        name = "OrderDomainEntity"
)(1)
public class OrderDomainEntity extends DomainEntity { (2)

    // Domain-specific fields
    @JsonProperty("username")
    private String username;
    @JsonProperty("user_validated")
    private boolean userValidated;
    @JsonProperty("total_amount")
    private double totalAmount;
    @JsonProperty("payment_reference_id")
    private String paymentReferenceId;
    @JsonProperty("product_items")
    private List<String> productItems;

    (3)
    public OrderDomainEntity() {
        super(OrderDomainEntity.class);
    }
}
1 @SagaDomainEntity 注解将该类标记为 StackSaga 的领域实体,表示它将用于承载事务流经 Saga 各步骤时的状态与数据。版本号属性使您能够随着时间推移规范管理领域实体的版本演进。
2 OrderDomainEntity 类继承自 StackSaga 提供的 DomainEntity 基类,该基类包含了领域实体所需的通用功能与内置属性。
3 构造函数通过传入自身类类型调用父类构造函数,这是 StackSaga 在事务流转期间正确管理与序列化领域实体所必需的。

创建主题 (Creating Topics)

定义主题 (Topics) 是 StackSaga-Kafka 实现中最核心的步骤之一,因为主题*即* Saga 的结构骨架:每个主题代表路由至特定工作节点服务的一个原子执行点(跨度 / Span)。 在 StackSaga-Kafka 中,主题不仅是普通字符串名称——它们携带了丰富的元数据,例如主题名称、唯一主题键 (Topic Key)、动作类型(正向 DO 或逆向补偿 UNDO)以及目标服务。 通过继承 AbstractTopic<T> 并将每个主题暴露为静态常量来声明它们,如下所示。了解更多

public class PlaceOrderTopic extends AbstractTopic<PlaceOrderTopic> { (1)

    (2)
    protected PlaceOrderTopic(String topicName, float topicKey, SagaEventType sagaEventType, String targetService) {
        super(topicName, topicKey, sagaEventType, targetService);
    }
    (2)
    protected PlaceOrderTopic(String topicName, float topicKey, SagaEventType sagaEventType, String targetService, PlaceOrderTopic parent) {
        super(topicName, topicKey, sagaEventType, targetService, parent);
    }


     (3)
    //do topics
    public static final PlaceOrderTopic DO_USER_VALIDATED = new PlaceOrderTopic("user-service.validate-user", 1.0f, SagaEventType.QUERY_DO_ACTION, "user-service");
    public static final PlaceOrderTopic DO_ORDER_RESERVED = new PlaceOrderTopic("order-service.reserve-order", 2.0f, SagaEventType.COMMAND_DO_ACTION, "order-service");
    public static final PlaceOrderTopic DO_MAKE_PAYMENT = new PlaceOrderTopic("payment-service.make-payment", 3.0f, SagaEventType.COMMAND_DO_ACTION, "payment-service");

    //undo topics
    public static final PlaceOrderTopic UNDO_ORDER_RESERVED = new PlaceOrderTopic("order-service.reserve-order", -2.0f, SagaEventType.COMMAND_UNDO_ACTION, "order-service", DO_ORDER_RESERVED);
    public static final PlaceOrderTopic UNDO_MAKE_PAYMENT = new PlaceOrderTopic("payment-service.make-payment", -3.0f, SagaEventType.COMMAND_UNDO_ACTION, "payment-service", DO_MAKE_PAYMENT);
}
1 为您的业务领域创建一个自定义主题类并继承 AbstractTopic<T>,其中 <T> 为主题类自身(即奇异递归模板模式 / Curiously Recurring Template Pattern)。
2 重写 AbstractTopic<T> 所要求的两个构造函数。4 参数构造函数用于声明*正向* (DO) 主题;5 参数构造函数增加了 parent 参数,用于声明*逆向补偿* (UNDO) 主题,并将其关联回所逆转的正向主题。除了在下方的常量中调用外,您绝不需要自行调用这些构造函数——只需重写它们并通过 super(…​) 委托即可。
3 将长事务 (LRT) 的每个主题声明为 public static final 常量。本示例包含 3 个正向 (DO) 主题 与 2 个补偿 (UNDO) 主题——注意 DO_USER_VALIDATED 为只读查询,因而没有补偿主题。每个主题指定其名称、唯一的 float 类型主题键(补偿主题为负数)、SagaEventType 以及目标服务。有关命名与主题键规则,请参阅 主题规范。
自定义主题类仅根据 StackSaga-Kafka 规范*声明*主题元数据——它是纯 Java 类,而非 Spring Bean,因此无需添加任何 Spring 组件注解。

创建主题 EventManager (Creating the Topic EventManager)

要使用定义的主题,必须根据 StackSaga 规范创建自定义事件管理器 (EventManager)。该自定义事件管理器会声明用于接收来自目标服务响应结果的领域回调主题。了解更多

(2)
@SagaEventManager(
        value = "placeOrderEventManager",
        listenerScope = OrchestratorListenerScope.SHARED_GLOBAL,
        groupType = GroupType.SHARE,
        domainCallbackTopicSuffix = "place-order"
)
public class PlaceOrderEventManager extends AbstractEventManager<OrderDomainEntity, PlaceOrderTopic> { (1)

    (3)
    @Override
    public Supplier<List<PlaceOrderTopic>> registerTopics() {
        return () -> List.of(
                PlaceOrderTopic.DO_USER_VALIDATED,
                PlaceOrderTopic.DO_ORDER_RESERVED,
                PlaceOrderTopic.UNDO_ORDER_RESERVED,
                PlaceOrderTopic.DO_MAKE_PAYMENT,
                PlaceOrderTopic.UNDO_MAKE_PAYMENT
        );
    }

    @NonNull
    @Override
    (4)
    public SagaPrimaryEventAction<PlaceOrderTopic> onNext(
            PlaceOrderTopic recentTopic,
            OrderDomainEntity currentDomainEntityState,
            SagaPrimaryEventActionUtil<PlaceOrderTopic> actionUtil
    ) {
        (5)
        if (recentTopic.equals(PlaceOrderTopic.DO_USER_VALIDATED)) {
            return actionUtil.next(PlaceOrderTopic.DO_ORDER_RESERVED);
        }
        (6)
        if (recentTopic.equals(PlaceOrderTopic.DO_ORDER_RESERVED)) {
            return actionUtil.next(PlaceOrderTopic.DO_MAKE_PAYMENT);
        }
        (7)
        if (recentTopic.equals(PlaceOrderTopic.DO_MAKE_PAYMENT)) {
            return actionUtil.complete();
        }
        (8)
        return actionUtil.error(new IllegalStateException("Unexpected topic: " + recentTopic));
    }
}
1 通过继承 AbstractEventManager<DE,T> 类创建自定义事件管理器。DE 为领域实体类型,<T> 为上文创建的主题类。
2 使用 @SagaEventManager 注解标记该类并配置必要属性:
  • value:Bean 的名称。

  • listenerScope:消费回调(应答)主题的监听器容器作用域。此处为 OrchestratorListenerScope.SHARED_GLOBAL,表示该事件管理器共享由所有默认事件管理器共用的单个全局回调监听容器,而非创建专用容器。

  • groupType:该容器使用的 Kafka 消费协议。此处为 GroupType.SHARE(Kafka 4 共享消费组协议 / KIP-932,队列式消费);使用 GroupType.CONSUMER 则对应经典的传统消费者组协议。

  • domainCallbackTopicSuffix:用于创建该领域专属回调主题的后缀。
    了解更多

    此处未展示,因为 SHARED_GLOBAL 不需要它,但 SHARED_GROUP/ISOLATED 事件管理器还可以通过 sharedGroupExecutionListener/isolatedExecutionListener 属性配置 autoStart(容器是否在应用启动时自动启动,默认 true)——参见 implementations:stacksaga-kafka-implementation/orchestrator/properties.adoc#orchestrator_auto_start。
3 重写 registerTopics 方法,返回该事件管理器使用的完整主题列表(包含 DO 和 UNDO)。在此处向 StackSaga 注册所有主题,仅在应用启动时调用一次。
4 重写 onNext 方法以定义事件管理器的路由调度逻辑。在每个正向执行主题(跨度/原子执行)成功完成后,框架会依次调用此方法直至事务结束(若未发生 NonRetryableExecutorException 失败)。使用传入的参数和 actionUtil 辅助工具返回下一步导航操作(actionUtil.next(…​)、actionUtil.complete() 或 actionUtil.error(…​))。
5 若最近完成的主题是 DO_USER_VALIDATED,通过 actionUtil.next(…​) 将引擎路由至 DO_ORDER_RESERVED。
6 若最近完成的主题是 DO_ORDER_RESERVED,通过 actionUtil.next(…​) 将引擎路由至 DO_MAKE_PAYMENT。
7 DO_MAKE_PAYMENT 是*最后*一个跨度,因此一旦其完成,后续无需再路由任何步骤——返回 actionUtil.complete() 顺利完成 Saga 并将其置为终态 COMPLETED。
IMPORTANT: 切勿遗漏此分支。若缺少该分支,*成功*的支付也会穿透到下方的 error(…​),导致引擎错误地触发补偿。Saga 总是以在此调用 complete() 成功结束,或者因故障触发补偿而终结。
8 针对非预期主题的防御性兜底分支——在上述路由逻辑下永远不应被触发。从 onNext() 返回 actionUtil.error(…​)(或抛出任何异常)会将 Saga 置为 FAILED 并触发逆向补偿。有关完整的 onNext()/onNextRevert() 契约,请参阅 EventManager 参考。

创建事件处理器 (Creating the Handler)

遵循 StackSaga 团队的 推荐 Handler 模式,处理器 (Handler) 是负责启动 Saga、查询其状态以及接收执行期间引擎分发的事务事件的单一职责类。 在本快速示例中,我们保持 Handler 最小化,仅用于接收事务状态变更事件;在实际应用中,它也是承载调用 StackSagaKafkaTemplate 发起 Saga 的理想场所。

@Slf4log
@Component(1)
public class PlaceOrderHandler implements KafkaTransactionEventListener<OrderDomainEntity> {(2)
    @Override(3)
    public void onStateChanged(TransactionState<OrderDomainEntity, AsyncExecutionEvent> transactionState) {
        log.info("onStateChanged : {}", transactionState.getCurrentStatus());
    }
}
1 为类添加 @Component 注解,声明为 Spring Bean,以便 Spring 容器自动检测并注册。
2 PlaceOrderHandler 类实现了 KafkaTransactionEventListener 接口,用于监听事务执行流中由 StackSaga 引擎发出的状态变更事件。了解更多
NOTE: 监听器提供了阻塞与非阻塞版本——在 Servlet (MVC) 应用中使用 KafkaTransactionEventListener,在响应式 (WebFlux) 应用中使用 ReactiveKafkaTransactionEventListener。详情参阅 事务状态变更监听器。
3 重写 onStateChanged 方法处理事务状态变更。它接收事务的当前状态作为参数,其中包含当前事务状态、执行历史、当前领域实体快照以及各阶段时间戳。本示例中仅在状态发生变化时打印日志。您可以根据实际业务场景扩展此方法,例如向用户推送通知或触发下游系统同步。

创建控制器以访问 StackSagaKafkaTemplate 触发事务 (Creating the Controller to trigger the transaction accessing StackSagaKafkaTemplate)

在此我们创建一个简单的 REST 控制器,通过暴露端点注入并调用 StackSagaKafkaTemplate(与 StackSaga-Kafka 引擎交互的核心 API)来启动 Saga。了解更多

这是一个普通的 Spring REST 控制器,包含一个用于下单的 POST 端点。在端点内部,我们使用 StackSagaKafkaTemplate 启动流程——提供初始领域实体状态,并指定首个主题及其所属的 EventManager。

@Slf4j
@RestController
@RequestMapping("/api/v1/order")
@RequiredArgsConstructor
public class PlaceOrderController {

    private final StackSagaKafkaTemplate<OrderDomainEntity, PlaceOrderTopic> stackSagaKafkaTemplate; (1)

    @PostMapping
    public String placeOrder(@RequestBody PlaceOrderRequest placeOrderRequest) {
        final String transactionId = this
                .stackSagaKafkaTemplate
                (2)
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    orderDomainEntity.setUsername(placeOrderRequest.getUsername());
                    orderDomainEntity.setTotalAmount(placeOrderRequest.getTotalAmount());
                    orderDomainEntity.setProductItems(Arrays.asList(placeOrderRequest.getItems()));
                    return orderDomainEntity;
                })
                (3)
                .peek(orderDomainEntity -> {
                    //you can do some local operation for the orderDomainEntity before the saga process.
                    //and also here you can access the unique id for the transaction that provide by the engine.
                    log.info("Transaction Id : {}", orderDomainEntity.getTransactionId());
                })
                (4)
                .startWith(PlaceOrderTopic.DO_USER_VALIDATED, PlaceOrderEventManager.class)
                (5)
                .execute();
        (6)
        return "Order placed successfully with transaction id: " + transactionId;
    }

    (7)
    @Data
    public static class PlaceOrderRequest {
        private String username;
        private double totalAmount;
        private String[] items;
    }
}
1 自动装配 StackSagaKafkaTemplate<DE, T>,这是与 StackSaga-Kafka 引擎交互的核心 API。 它提供了流畅的链式调用 API 来定义和启动事务。 两个类型参数分别为领域实体 DE (OrderDomainEntity) 与主题类 T (PlaceOrderTopic)——与参数化 EventManager 的类型对完全一致。
2 调用 init 方法,使用领域实体的初始状态初始化事务流程。它接收一个返回初始领域实体状态的 Supplier 函数,在本示例中根据请求体数据进行构建。
3 peek 方法是一个可选步骤,允许您在 Saga 流程正式启动前对领域实体执行本地准备工作。同时也可以在此访问由引擎生成的全局唯一事务 ID,便于日志记录与链路追踪。
4 startWith 方法指定两项内容:首个分发的主题 (PlaceOrderTopic.DO_USER_VALIDATED) 以及负责该 Saga 领域路由调度的 EventManager 类 (PlaceOrderEventManager)。从此时起,EventManager.onNext(…​) 逻辑将逐个主题驱动后续流程。
5 调用 execute 方法以既定配置启动事务流程。它返回由引擎为该事务生成的唯一事务 ID,可用于后续的追踪与状态查询。
6 端点返回响应,包含下单成功的提示信息与该事务的唯一 ID。
7 PlaceOrderRequest 是一个简单的 DTO 类,用于承载下单请求体,包含构建初始领域实体所需的 username、totalAmount 和 items 字段。不建议直接将领域实体类暴露为 HTTP 请求体,以避免 API 协议层与内部领域层产生强耦合,并确保 API 契约的清晰稳定。

调整配置 (Tune Configuration)

要运行该应用,请在 application.properties 或 application.yml 文件中配置以下属性。

#spring properties
spring.application.name=order-service
#stacksaga-instance properties (1)
stacksaga.instance.cluster=local-cluster
stacksaga.instance.region=local-region
stacksaga.instance.zone=local-zone
#stacksaga-starter properties
stacksaga.kafka.orchestrator.domain-entity-scan=org.example.orderservice.domain (2)
#stacksaga-mysql-database-support properties (3)
stacksaga.mysql.r2dbc.url=r2dbc:mysql://localhost:3306
stacksaga.mysql.r2dbc.database=order_service_event_store
stacksaga.mysql.r2dbc.username=${MYSQL_USER:root}
stacksaga.mysql.r2dbc.password=${MYSQL_PASSWORD:password}
(4)
stacksaga.mysql.jdbc.url=jdbc:mysql://localhost:3306/order_service_event_store
stacksaga.mysql.jdbc.username=${MYSQL_USER:root}
stacksaga.mysql.jdbc.password=${MYSQL_PASSWORD:password}
(5)
spring.kafka.bootstrap-servers=${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
1 stacksaga.instance 属性用于配置 StackSaga 引擎的实例拓扑信息,例如集群 (cluster)、区域 (region) 与可用区 (zone)。这些属性用于在分布式环境中精准标识事务,以便进行故障重试与宕机恢复。请注意,如果您指定了除 default 之外的自定义 region,则必须声明自定义 SagaRegionResolver Spring Bean。在此示例中配置为 local-cluster、local-region 和 local-zone 仅供演示。
2 stacksaga.kafka.orchestrator.domain-entity-scan 属性用于指定领域实体类所在的包路径。StackSaga 引擎将扫描该包并注册其中的领域实体类供事务流程使用。如果领域实体分布在多个包中,可以通过逗号分隔列出。
3 stacksaga.mysql.r2dbc 属性配置 R2DBC 连接,供 StackSaga 引擎连接 MySQL 事件存储数据库。提供连接 URL、数据库名称、用户名和密码。事件存储应为专用于存储事务事件和状态的独立 Schema,避免与业务表混用以获得良好的关注点分离。
4 stacksaga.mysql.jdbc 属性用于配置 JDBC 连接。内部由 Liquibase 使用该连接执行数据库版本迁移与 Schema 表结构初始化。
5 配置 Kafka 连接信息。

配置用户服务作为工作节点 (Configuring User-Service As Worker)

接下来是将 StackSaga Kafka 功能引入目标支撑服务。让我们对 user-service 进行配置。 向 user-service 添加以下依赖项:

依赖管理 (Dependency Management)

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>
<dependency>
    <groupId>org.stacksaga</groupId>
    <artifactId>stacksaga-kafka-worker-spring-boot-starter</artifactId>
</dependency>

接下来,创建自定义 UserValidationEndpoint 类,在工作节点端处理“验证用户”跨度。

自定义 UserValidationEndpoint

@Slf4j
@SagaEndpoint(topicNameSuffix = "user-service.validate-user", listenerScope = WorkerListenerScope.SHARED_GLOBAL, groupType = GroupType.SHARE) (2)
public class UserValidationEndpoint extends QueryEndpoint { (1)

    @Override (3)
    public void doProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException, NonRetryableExecutorException {
        log.info("Transaction Id: {},Idempotency Key: {}", consumerRecord.value().getTransactionId(), consumerRecord.value().getIdempotencyKey());
        consumerRecord.value().getCurrentDomainEntityStateForUpdate().ifPresent(domainEntity -> {
            String username = domainEntity.get("username").asText();
            log.info("username: {}", username);
            //Validate username with your logic here and if the user is not valid you can throw NonRetryableExecutorException
            //If there is a transient error, you can throw RetryableExecutorException or JustRetryableExecutorException

    try {
                //simulate some delay for make the request
                Thread.sleep(new Random().nextInt(1000, 3000));
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }

        });
    }
}
1 继承 QueryEndpoint 抽象类创建自定义 UserValidationEndpoint。
2 为类添加 @SagaEndpoint 注解并配置必要属性:
  • topicNameSuffix:提供该端点监听的主题名称后缀。框架会自动添加 saga. 前缀构建实际的主题名称——因此 user-service.validate-user 会对应真实主题 saga.user-service.validate-user(若您的后缀已包含 saga.,则保持原样)。这是 EventManager 发送命令消息以触发该端点执行的目标主题。

  • listenerScope:消费该端点主题的监听器容器作用域。此处使用 WorkerListenerScope.SHARED_GLOBAL,表示该端点共享由所有默认端点共用的单个全局监听器容器,而非创建专用容器。

  • groupType:该容器使用的 Kafka 组协议。此处为 GroupType.SHARE(Kafka 4 共享消费组协议 / KIP-932,队列式消费);使用 GroupType.CONSUMER 则对应经典的传统消费者组协议。

  • autoStart:此处未展示,默认值为 true;控制该端点的监听容器是否在应用启动时自动启动。所有共享同一个 SHARED_GROUP containerId 的 @SagaEndpoint 必须声明*完全相同*的 autoStart,否则启动将失败并抛出 ValidationException——参见 implementations:stacksaga-kafka-implementation/worker/worker-configuration.adoc#worker_auto_start。

3 重写 doProcess 方法实现用户校验逻辑。consumerRecord 携带了与该事务相关的全部上下文——当前领域实体状态、事务 ID、幂等键以及 更多详细信息。

有关更多详细信息,请参阅 工作节点端点说明。

配置属性 (Configuration Properties)

spring.application.name=user-service
#kafka connection (1)
spring.kafka.bootstrap-servers=${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
#stacksaga-instance properties (2)
stacksaga.instance.cluster=local-cluster
stacksaga.instance.region=local-region
stacksaga.instance.zone=local-zone
1 工作节点连接至与编排器*相同*的 Kafka 集群——该 Broker 是其接收命令与回发应答的唯一通道。工作节点无需配置数据库或事件存储;它相对于 Saga 是无状态的。
2 声明 StackSaga 实例坐标。所有 Saga 参与方——无论是编排器还是工作节点——都必须共享相同的 cluster / region / zone,以便引擎能够正确关联它们。

配置订单服务作为工作节点 (Configuring Order-Service As a Worker)

“预留订单”跨度在 order-service 自身内部执行——即同时承担编排器角色的同一个应用程序。 因此,order-service 在此 Saga 中*既是*编排器,*又是*工作节点。

其 pom.xml 中无需追加任何依赖:先前添加的 stacksaga-kafka-orchestrator-spring-boot-starter 已经通过传递依赖引入了完整的 Worker 功能。我们只需创建负责预留已购商品的端点类即可。

自定义 ReserverOrderEndpoint

@Slf4j
@SagaEndpoint(topicNameSuffix = "order-service.reserve-order", listenerScope = WorkerListenerScope.SHARED_GLOBAL, groupType = GroupType.SHARE) (2)
public class ReserverOrderEndpoint extends CommandEndpoint { (1)

    (3)
    @Override
    public void doProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException, NonRetryableExecutorException {
        log.info("Transaction Id: {},Idempotency Key: {}", consumerRecord.value().getTransactionId(), consumerRecord.value().getIdempotencyKey());
        consumerRecord.value().getCurrentDomainEntityStateForUpdate().ifPresent(domainEntity -> {
            String username = domainEntity.get("username").asText();
            log.info("username: {}", username);
            JsonNode productItems = domainEntity.get("product_items");
            for (JsonNode productItem : productItems) {
                String productId = productItem.asText();
                log.info("productId: {} ", productId);
            }
            try {
                //simulate some delay for make the request
                Thread.sleep(new Random().nextInt(1000, 3000));
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }
            // TODO: Reserve order with your logic here
        });
    }

    (4)
    @Override
    public void undoProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException {
        SagaPayload payload = consumerRecord.value();
        log.info("Transaction Id: {},Idempotency Key: {}", payload.getTransactionId(), payload.getIdempotencyKey());
        JsonNode lastDomainEntityState = payload.getDomainEntityState();
        String username = lastDomainEntityState.get("username").asText();
        log.info("username: {}", username);
        payload.getPrimaryExecutionException().ifPresent(primaryExecutionExceptionMetaData -> {
            log.info("primary error message : {}", primaryExecutionExceptionMetaData.getRealExceptionMessage());
        });
        // TODO: revert the reserve order with your logic here
            try {
                //simulate some delay for make the request
                Thread.sleep(new Random().nextInt(1000, 3000));
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }
        (5)
        payload.getHintStore().ifPresent(hintStore -> {
            hintStore.put("undone_reserve_order", "true");
        });
    }
}
1 预留订单属于状态变更操作,因此端点继承 CommandEndpoint 抽象类——要求同时提供正向业务 (doProcess) 与逆向补偿 (undoProcess) 实现。
2 使用 @SagaEndpoint 标记该类并配置其属性:
  • topicNameSuffix:该端点监听的主题后缀。框架添加前缀构建实际主题名——此处 order-service.reserve-order 变为 saga.order-service.reserve-order(仅在未包含 saga. 时添加)。这是 EventManager 分发 DO_ORDER_RESERVED 命令的目标主题。

  • listenerScope:用于消费该端点 do/undo 主题的容器作用域。此处使用 WorkerListenerScope.SHARED_GLOBAL,表示该端点共享全局监听器容器,而非创建专用容器。

  • groupType:该容器的 Kafka 消费协议。此处使用 GroupType.SHARE(Kafka 4 共享消费组协议 / KIP-932);使用 GroupType.CONSUMER 对应传统消费者组协议。

  • autoStart:此处未展示,默认值为 true;控制该端点的监听容器是否在应用启动时自动启动——参见 implementations:stacksaga-kafka-implementation/worker/worker-configuration.adoc#worker_auto_start。

3 重写 doProcess 方法实现正向订单预留逻辑。consumerRecord 携带了事务的全部信息——当前领域实体状态、事务 ID、幂等键以及 更多信息。
4 重写 undoProcess 方法实现释放订单预留的*补偿*逻辑。它能够获取领域实体最终状态以及引发回滚的原始异常,详见 更多信息。
5 undoProcess 方法还提供了 HintStore——用于补偿流程的键值暂存区。可使用它在各个补偿步骤之间传递标志或元数据(例如记录预留已成功回滚)。

有关更多细节,请参阅 工作节点端点说明。

配置属性 (Configuration Properties)

此处无需追加额外配置。因为 Worker 角色直接运行在驱动 Saga 的*同一个* order-service 应用程序内部,之前在 调整配置 中已经设置好了实例属性、事件存储与 Kafka 连接。

配置支付服务作为工作节点 (Configuring Payment-Service As a Worker)

正如我们对 user-service 所做的那样,我们将 payment-service 配置为工作节点。因为它是独立的服务(非编排器),所以需要显式添加 worker starter 依赖。 向 payment-service 添加以下依赖:

依赖管理 (Dependency Management)

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>
<dependency>
    <groupId>org.stacksaga</groupId>
    <artifactId>stacksaga-kafka-worker-spring-boot-starter</artifactId>
</dependency>

下一步是创建自定义 MakePaymentEndpoint 类来处理“执行支付”跨度。 由于执行支付是状态变更(命令)操作,因此它继承 CommandEndpoint 抽象类,提供正向操作 (doProcess) 与逆向补偿操作 (undoProcess)。 为了直观验证补偿流程,该端点故意通过抛出 NonRetryableExecutorException 随机模拟支付失败——这正是应当直接触发回滚而非重试的永久性业务故障(例如“余额不足”)。

自定义 MakePaymentEndpoint

@Slf4log
@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 {
        SagaPayload payload = consumerRecord.value();
        log.info("Transaction Id: {},Idempotency Key: {}", payload.getTransactionId(), payload.getIdempotencyKey());
        payload.getCurrentDomainEntityStateForUpdate().ifPresent(domainEntity -> {
            String username = domainEntity.get("username").asText();
            log.info("username: {}", username);
            double totalAmount = domainEntity.get("total_amount").asDouble();
            log.info("total amount: {}", totalAmount);
            // TODO: Make payment with your logic here
            try {
                //simulate some delay for make the request
                Thread.sleep(new Random().nextInt(1000, 3000));
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }
            //simulate a payment failure to demonstrate compensation.
            //run the request a few times until this branch is hit: the orchestrator will then
            //compensate the previously completed command span (UNDO_ORDER_RESERVED).
            if (new Random().nextBoolean()) {
                throw NonRetryableExecutorException
                        .buildWith(new RuntimeException("insufficient balance"))
                        .put("error_code", "MAKE_PAYMENT_FAILED")
                        .build();
            }
        });
    }

    @Override
    public void undoProcess(ConsumerRecord<String, SagaPayload> consumerRecord)
            throws JustRetryableExecutorException, RetryableExecutorException {
        SagaPayload payload = consumerRecord.value();
        log.info("Transaction Id: {},Idempotency Key: {}", payload.getTransactionId(), payload.getIdempotencyKey());
        JsonNode lastDomainEntityState = payload.getDomainEntityState();
        String username = lastDomainEntityState.get("username").asText();
        log.info("username: {}", username);
        double totalAmount = lastDomainEntityState.get("total_amount").asDouble();
        log.info("total amount: {}", totalAmount);

        payload.getPrimaryExecutionException().ifPresent(primaryExecutionExceptionMetaData -> {
            log.info("primary error message : {}", primaryExecutionExceptionMetaData.getRealExceptionMessage());
        });
        // TODO: revert the payment with your logic here
        payload.getHintStore().ifPresent(hintStore -> {
            hintStore.put("undone_make_payment", "true");
        });
    }
}
框架根据抛出的*异常类型*决定后续行为。NonRetryableExecutorException(此处使用)表示*永久性*业务故障并立即启动逆向补偿。对于*瞬态*故障(超时、网络抖动),应抛出 JustRetryableExecutorException / RetryableExecutorException,以便框架重试该跨度而非回滚整个 Saga。参见 异常处理参考。

有关更多详细信息,请参阅 工作节点端点说明。

配置属性 (Configuration Properties)

spring.application.name=payment-service
#kafka connection (1)
spring.kafka.bootstrap-servers=${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
#stacksaga-instance properties (2)
stacksaga.instance.cluster=local-cluster
stacksaga.instance.region=local-region
stacksaga.instance.zone=local-zone
1 工作节点连接至与编排器*相同*的 Kafka 集群——该 Broker 是其接收命令与回发应答的唯一通道。工作节点无需配置数据库或事件存储;它相对于 Saga 是无状态的。
2 声明 StackSaga 实例坐标。所有 Saga 参与方——无论是编排器还是工作节点——都必须共享相同的 cluster / region / zone,以便引擎能够正确关联它们。

运行应用并测试事务流程 (Running the Application and Testing the Transaction Flow)

在 Kafka Broker 和 MySQL 正常运行的前提下,启动所有三个服务——order-service、user-service 和 payment-service。在首次启动时,事件存储数据库 Schema 会被自动创建(通过 Liquibase)。所需的 Kafka 主题必须已存在——可以通过 Broker 的自动建主题行为 (auto.create.topics.enable=true) 创建,或者提前预先创建好。

向 order-service 的 /api/v1/order 端点发送 POST 请求触发 Saga 流程:

curl -X POST http://localhost:8080/api/v1/order \
  -H "Content-Type: application/json" \
  -d '{
        "username": "john.doe",
        "totalAmount": 149.99,
        "items": ["item-1", "item-2"]
      }'

由于 MakePaymentEndpoint 会随机模拟失败,请重复发起数次请求以观察*两种*执行结果。PlaceOrderHandler.onStateChanged(…​) 回调会打印每一次状态流转,您可以在 order-service 控制台中完整追踪 Saga 轨迹。

  • 正常流程 (Happy Path)(支付成功)—— Saga 顺利执行至终态:

2026-06-21T23:57:36.911+05:30  INFO 19868 --- [order-service] [nio-8080-exec-3] o.e.o.controller.PlaceOrderController    : Transaction Id : e0b34f8a-9702-4e23-87e8-8ae5b9654a61
2026-06-21T23:57:38.357+05:30  INFO 19868 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : PROCESSING
2026-06-21T23:57:41.177+05:30  INFO 19868 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : PROCESSING
2026-06-21T23:57:42.527+05:30  INFO 19868 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : PROCESS_COMPLETED
  • 补偿流程 (Compensation Path)(支付失败)—— 引擎自动逆向回滚此前已完成的命令跨度 (UNDO_ORDER_RESERVED),最终以 REVERT_COMPLETED 结束:

2026-06-21T23:54:19.094+05:30  INFO 17792 --- [order-service] [nio-8080-exec-1] o.e.o.controller.PlaceOrderController    : Transaction Id : 0864a84c-79a1-45a1-a688-662f431ca218
2026-06-21T23:54:26.248+05:30  INFO 17792 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : PROCESSING
2026-06-21T23:54:27.956+05:30  INFO 17792 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : PROCESSING
2026-06-21T23:54:29.980+05:30  INFO 17792 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : REVERTING
2026-06-21T23:54:32.546+05:30  INFO 17792 --- [order-service] [oundedElastic-3] o.e.o.handler.PlaceOrderHandler          : onStateChanged : REVERT_COMPLETED