快速入门示例:StackSaga:Synchronous:MySQL
本快速入门示例将引导您在 StackSaga 同步(阻塞式)编排引擎上构建一个完整的*下单 (order-placement)* Saga,并使用 MySQL 作为事件存储 (Event Store)。 我们将从最小可用的正常流程 (Happy Path) 实现入手,随后逐步叠加补偿 (Compensation)、重试 (Retry) 以及可观测性 (Observability) 能力。
在同步模型中,单个*编排器服务 (Orchestrator Service)* 在进程内运行整个 Saga 流程——每个执行器依次内联调用其目标服务(通过 HTTP、gRPC 或本地调用)并阻塞等待返回。这最大限度减少了系统活动组件。如果您需要基于 Kafka 的完全事件驱动、非阻塞式执行,请参阅 异步 (Kafka) 快速入门示例。
概述 (Overview)
我们编排的业务场景是一个简化的*下单 (place-order)* 操作,建模为由三个原子步骤(执行器 (Executors))组成的单个长事务 (Long-Running Transaction, LRT):
| # | 步骤 (执行器) | 类型 | 补偿 (doRevert) |
|---|---|---|---|
1 |
验证用户 ( |
查询 (Query / 只读) |
— (查询操作从不补偿) |
2 |
预留订单 ( |
命令 (Command / 状态变更) |
释放已预留的订单 |
3 |
执行支付 ( |
命令 (Command / 状态变更) |
— (最后一步;其后无后续步骤触发回滚) |
如果任何步骤失败,StackSaga 会自动按倒序*补偿 (Compensate)* 此前已成功执行的*命令*步骤——例如在支付失败时释放已预留的订单——从而确保系统最终始终处于一致状态。
|
架构提醒:连续型长事务 vs. 可暂停长事务
本快速示例演示了 连续型长事务 (Continuous Long-Running Transaction, Continuous LRT),其中跨微服务的执行按顺序连续推进且不中断挂起。 如果您的实际业务场景涉及异步回调、外部支付 Webhook 或人工介入 (Human-in-the-Loop, HITL) 审核(例如等待餐厅接单确认或经理审批),StackSaga 还原生支持 可暂停长事务 (Pausable Long-Running Transaction, Pausable LRT)。
在这些场景中,执行器可调用 |
我们将分三个渐进阶段进行构建:
-
阶段一 —— 最小实现 (Stage-1 · Minimal Implementation):创建领域实体、三个执行器、事件处理器以及触发 Happy-Path 与补偿的 REST 控制器。
-
阶段二 —— 添加重试能力 (Stage-2 · Adding Retry Capability):引入环形协调器 (Ring Coordinator) 应用与重试子系统,使流程能够应对瞬态故障。
-
阶段三 —— 可观测性 (Stage-3 · Observability):接入 StackSaga 链路追踪窗口 (Trace Window),实时可视化与监控事务。
阶段一 [最小实现] (Stage-1 [Minimal Implementation])
阶段一构建一个最小但完整的 Saga 流程——领域实体、三个执行器(包含补偿逻辑)、事件处理器以及 REST 控制器——全部以 MySQL 作为事件存储。我们将按以下顺序实现:
前置准备 (Prerequisites)
在开始之前,请确保已安装并运行以下组件:
-
Java 21+
-
Spring Boot 4.x
-
MySQL 8+ —— 由编排器用作事件存储 (Event Store)。
-
Maven
启动本地数据库最便捷的方式是 Docker——例如 mysql:8 容器。与异步 (Kafka) 示例不同,同步引擎不需要任何消息代理。
|
起步 - 初始项目搭建 (Getting Started - Initial project setup)
在本示例中,我们使用 Spring MVC 作为 Web 框架,并使用 MySQL 作为事件存储的主数据库。
因此,首先创建一个新的 Spring Boot 应用程序,并在 pom.xml 中添加以下依赖项与插件配置:
| 建议使用 StackSaga Initializer 为您的项目生成 StackSaga 依赖配置代码,以确保版本兼容性与初始配置正确无误。 |
<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.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency> (2)
<groupId>org.stacksaga</groupId>
<artifactId>stacksaga-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 Spring Boot Starter 依赖项,包含了 StackSaga 的核心功能,例如编排引擎、事件处理和事务管理。它为在应用中实现 StackSaga 模式提供了必要的组件。 |
| 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) 是事务流程中流转的核心实体,它代表在整个事务期间被操作和处理的主数据结构。了解更多
@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 the Executors)
下一步是为长事务 (LRT) 的每个跨度(原子步骤)创建执行器。执行器负责执行事务流程中各个步骤的具体业务逻辑。了解更多
在此快速示例中,我们将创建 3 个执行器:分别用于*验证用户 (Validating User)、*预留订单 (Reserving Order) 以及*执行支付 (Make Payment)*。(本例仅用于演示,实际生产应用中业务逻辑与流转步骤会更为复杂)。
创建 ValidateUserExecutor
ValidateUserExecutor 负责在事务流程中验证用户信息。它检查用户是否有效,并相应地更新领域实体。
(1)
@SagaExecutor(executeFor = "user-service", value = "ValidateUserExecutor")
public class ValidateUserExecutor implements QueryExecutor<OrderDomainEntity> { (2)
@NonNull
@Override
(3)
public ProcessStepManager<OrderDomainEntity> doProcess(
OrderDomainEntity currentDomainEntityState,
ProcessStepManagerUtil<OrderDomainEntity> stepManager,
String idempotencyKey
) throws RetryableExecutorException, NonRetryableExecutorException {
//call the user-service via http or rpc to validate the user information.
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
(4)
return stepManager.next(ReserveOrderExecutor.class, "VALIDATED_USER");
}
}
| 1 | @SagaExecutor 注解将此类标记为 StackSaga 执行器,指定其负责的服务名称 (user-service) 以及执行器名称 (ValidateUserExecutor)。 |
| 2 | ValidateUserExecutor 类实现了 QueryExecutor 接口,该接口专用于执行*只读*操作或数据校验、且不改变系统外部持久状态的执行器。了解更多 |
| 3 | doProcess 方法是执行执行器主干业务逻辑的统一入口。它接收领域实体的当前状态、用于管理流程的步骤管理器工具类 (ProcessStepManagerUtil),以及用于确保操作在重试时安全防重的幂等键 (Idempotency Key)。在此处可实现验证用户信息的逻辑,例如调用外部服务或查询数据库。 |
| 4 | 调用 stepManager.next 方法声明事务流程应推进至下一个步骤,在本例中为 ReserveOrderExecutor。动作事件名称 ("VALIDATED_USER") 作为字符串直接传递并持久化至事件存储中。TIP: 有关动作事件命名规范(动词过去分词优先)以及维护实现 SagaExecutionEventName 的集中枚举这一推荐实践,请参阅 执行器架构与最佳实践。 |
第一个执行器既可以是 QueryExecutor 也可以是 CommandExecutor。框架对事务流程中第一个执行器的类型没有限制。
|
创建 ReserveOrderExecutor
ReserveOrderExecutor 负责在事务流程中预留订单。它检查订单是否可预留并执行相应的库存预扣操作,同时更新领域实体。
@SagaExecutor(executeFor = "order-service", value = "ReserveOrderExecutor") (1)
public class ReserveOrderExecutor implements CommandExecutor<OrderDomainEntity> { (2)
@NonNull
@Override
(3)
public ProcessStepManager<OrderDomainEntity> doProcess(
OrderDomainEntity currentDomainEntityState,
ProcessStepManagerUtil<OrderDomainEntity> stepManager,
String idempotencyKey
) throws RetryableExecutorException, NonRetryableExecutorException {
//call the internal service and reserve the order
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
(4)
return stepManager.next(MakePaymentExecutor.class, "RESERVED_ORDER");
}
@NonNull
@Override
(5)
public SagaExecutionEventName doRevert(
NonRetryableExecutorException primaryExecutionException,
OrderDomainEntity finalDomainEntityState,
RevertHintStore revertHintStore,
String idempotencyKey
) throws RetryableExecutorException {
//call the internal service to revert the reserve order action
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
(6)
return SagaExecutionEventName.of("REVERTED_ORDER");
}
}
| 1 | @SagaExecutor 注解将此类标记为 StackSaga 执行器,指定其执行的服务名称 (order-service) 与执行器名称 (ReserveOrderExecutor)。 |
| 2 | ReserveOrderExecutor 类实现了 CommandExecutor 接口,用于对领域实体执行*状态变更*操作的执行器。了解更多 |
| 3 | doProcess 方法是执行器主干业务逻辑的入口。它接收领域实体当前状态、步骤管理器工具类以及幂等键。在此处实现预留订单的逻辑,例如调用外部服务或执行数据库预扣操作。在此场景下,订单服务即为编排器应用自身,因此可直接进行内部调用而无需产生额外网络开销。 |
| 4 | 调用 stepManager.next 方法指示事务流程应执行下一步骤(即 MakePaymentExecutor)。动作事件名称 ("RESERVED_ORDER") 作为字符串直接传入。TIP: 在该连续流中, stepManager.next(…) 会立即推进至下一步。如果您的业务逻辑需要在此暂停、等待外部 Webhook、异步回调或人工审批后再进入支付环节,您可以返回 stepManager.pause(…)。详情请参阅 可暂停长事务 (Pausable LRT)。 |
| 5 | 当后续步骤发生失败导致事务流程需要回滚时,框架会自动调用 doRevert 方法执行补偿。它接收导致回滚的异常、领域实体的最终状态、用于暂存逆向补偿上下文的提示存储 (RevertHintStore),以及幂等键。在此处实现撤销订单预留的补偿逻辑,例如调用服务或数据库操作解除库存冻结。在此场景下,订单服务为编排器应用自身的内联服务,可直接处理而无需网络调用。 |
| 6 | SagaExecutionEventName.of("REVERTED_ORDER") 提供表示订单预留已成功回滚的事件名称。若维护了实现 SagaExecutionEventName 的集中枚举(例如 OrderExecutionEvent),可直接返回对应的枚举值。了解更多事件最佳实践。 |
创建 MakePaymentExecutor
MakePaymentExecutor 负责在事务流程中处理支付。它校验支付资格并执行扣款操作,随后更新领域实体。
@SagaExecutor(executeFor = "payment-service", value = "MakePaymentExecutor") (1)
public class MakePaymentExecutor implements CommandExecutor<OrderDomainEntity> { (2)
@NonNull
@Override
(3)
public ProcessStepManager<OrderDomainEntity> doProcess(
OrderDomainEntity currentDomainEntityState,
ProcessStepManagerUtil<OrderDomainEntity> stepManager,
String idempotencyKey
) throws RetryableExecutorException, NonRetryableExecutorException {
//call the payment-service via http or rpc or with any protocol to make the payment.
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
(4)
//simulate payment failure randomly to test the compensation mechanism
if (new Random().nextBoolean()) {
throw NonRetryableExecutorException
.buildWith(new RuntimeException("insufficient balance"), SagaExecutionEventName.of("FAILED_PAYMENT"))
.build();
} else {
return stepManager.complete("MADE_PAYMENT");
}
}
@NonNull
@Override
(5)
public SagaExecutionEventName doRevert(
NonRetryableExecutorException primaryExecutionException,
OrderDomainEntity finalDomainEntityState,
RevertHintStore revertHintStore,
String idempotencyKey
) throws RetryableExecutorException {
throw new UnsupportedOperationException("this will not be executed until there is another executor after this executor in the flow.");
}
}
| 1 | @SagaExecutor 注解将此类声明为 StackSaga 执行器,指定执行服务 (payment-service) 与执行器名称 (MakePaymentExecutor)。 |
||
| 2 | MakePaymentExecutor 类实现了 CommandExecutor 接口,用于执行*状态变更*的执行器。了解更多 |
||
| 3 | doProcess 方法是执行器主要逻辑的入口。它接收当前领域实体状态、步骤管理器以及幂等键。在此处实现扣款逻辑(例如调用第三方支付网关)。本示例中通过随机抛出异常模拟支付失败,以演示 StackSaga 引擎的故障补偿机制。若支付失败,抛出携带失败事件名称的 NonRetryableExecutorException。若支付成功,调用 stepManager.complete 标记事务流已顺利完成,并记录支付成功事件。若支付失败,StackSaga 引擎将按相反顺序自动触发先前已完成执行器的 doRevert 方法执行补偿,首先调用 ReserveOrderExecutor 撤销订单预留。
|
||
| 4 | 在 stepManager.complete("MADE_PAYMENT") 中,表示支付成功的动作事件名称作为字符串直接传入。若支付失败,NonRetryableExecutorException.buildWith(…, SagaExecutionEventName.of("FAILED_PAYMENT")) 发出不可重试的永久性失败信号并携带事件名,触发先前步骤的逆向补偿。 |
||
| 5 | 当事务因后续步骤失败需要回滚时调用 doRevert 方法。由于 MakePaymentExecutor 是该事务流程中的最后一步,其后没有其他执行器,因此当本步骤自身失败时只需回滚之前的预留订单步骤,无需自身执行 doRevert。因此在 doRevert 中直接抛出 UnsupportedOperationException。若未来在支付步骤之后追加了发货等其他执行器,则必须在此实现支付的退款补偿逻辑。在真实的支付业务中,务必根据实际情况实现逆向退款或撤销逻辑。 |
创建事件处理器 (Creating the Handler)
事件处理器 (Handler) 负责在事务执行过程中接收 StackSaga 引擎分发的事务生命周期事件。了解更多
@Slf4j
@Component(1)
public class PlaceOrderHandler implements TransactionEventListener<OrderDomainEntity> {(2)
@Override(3)
public void onStateChanged(TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> transactionState) {
log.info("onStateChanged : {}", transactionState.getCurrentStatus());
}
}
| 1 | 为类添加 @Component 注解,声明为 Spring Bean,以便 Spring 容器自动扫描并注册。 |
| 2 | PlaceOrderHandler 类实现了 TransactionEventListener 接口,用于监听事务执行流中由 StackSaga 引擎发出的状态变更事件。了解更多NOTE: TransactionEventListener 提供了阻塞与非阻塞版本。完整细节请参阅 TransactionEventListener 与 ReactiveTransactionEventListener。 |
| 3 | 重写 onStateChanged 方法处理事务状态变更。它接收事务的当前状态作为参数,其中包含当前事务状态、执行历史、当前领域实体快照以及各阶段时间戳。本示例中仅在状态发生变化时打印日志。您可以根据实际业务场景扩展此方法,例如向用户推送通知或触发下游系统同步。 |
创建控制器以访问 SagaTemplate 触发事务 (Creating the Controller to trigger the transaction accessing SagaTemplate)
在此我们创建一个简单的 REST 控制器,通过暴露端点注入并调用 SagaTemplate(与 StackSaga 引擎交互的核心 API)来触发 Saga 并按需查询其执行状态。了解更多
这是一个普通的 Spring REST 控制器,包含两个端点:
-
POST /api/v1/order—— 提供初始领域实体状态并指定起始执行器,从而启动 Saga 流程。 -
GET /api/v1/order/status?order-id={id}—— 从事件存储中按需拉取指定事务当前的状态快照。
@Slf4j
@RestController
@RequestMapping("/api/v1/order")
@RequiredArgsConstructor
public class PlaceOrderController {
private final SagaTemplate<OrderDomainEntity> stacksagaTemplate; (1)
@PostMapping
public String placeOrder(@RequestBody PlaceOrderRequest placeOrderRequest) {
final String transactionId = this
.stacksagaTemplate
(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(ValidateUserExecutor.class)
(5)
.fireAndForget()
(6)
.execute();
(7)
return "Order placed successfully with transaction id: " + transactionId;
}
@GetMapping("/status")
public TransactionState<?, ?> orderStatus(@RequestParam("order-id") String orderId) {
return this
.stacksagaTemplate
(8)
.getCurrentState(orderId)
.fetch();
}
(9)
@Data
public static class PlaceOrderRequest {
private String username;
private double totalAmount;
private String[] items;
}
}
| 1 | 自动装配 SagaTemplate<DE>,这是与 StackSaga 引擎交互的核心 API。它提供了流畅的链式调用 API 来定义和执行事务流。DE 为领域实体类型——此处为 OrderDomainEntity。NOTE: SagaTemplate 提供了阻塞式(execute(), fetch())与响应式(executeAsync(), fetchAsync())两套方法。在响应式 (WebFlux) 应用中,请使用响应式方法。了解更多 |
| 2 | 调用 init 方法,使用领域实体的初始状态初始化事务流程。它接收一个返回初始领域实体状态的 Supplier 函数,在本示例中根据请求体数据进行构建。 |
| 3 | peek 方法是一个可选步骤,允许您在 Saga 流程正式启动前对领域实体执行本地准备工作。同时也可以在此访问由引擎生成的全局唯一事务 ID,便于日志记录与链路追踪。 |
| 4 | 调用 startWith 方法指定事务流程中首个执行的执行器,在本示例中为 ValidateUserExecutor。 |
| 5 | 调用 fireAndForget 方法指示事务以“即发即弃”模式异步执行,即调用方无需阻塞等待事务全流程结束即可先行返回,继续处理其他请求。这在长事务场景下尤为推荐,能够避免占用和阻塞 Web 容器的工作线程。您可以根据应用类型(响应式或 Servlet)通过实现 ReactiveTransactionEventListener 或 TransactionEventListener 持续监听事务状态变更事件。 |
| 6 | 调用 execute 方法以既定配置启动事务流程。它返回由引擎为该事务生成的唯一事务 ID,可用于后续的追踪与状态查询。 |
| 7 | 端点返回响应,包含下单成功的提示信息与该事务的唯一 ID。 |
| 8 | orderStatus 端点按需查询事务的当前状态。调用 .getCurrentState(orderId).fetch() 直接从事件存储中检索指定事务 ID 的 TransactionState 状态快照,涵盖完成状态、每一步的执行历史以及领域实体最新数据。了解更多 |
| 9 | 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.domain-entity-scan=org.example.orderservice.domain (2)
#stacksaga-datasource properties (3)
stacksaga.datasource.r2dbc.url=r2dbc:mysql://localhost:3306
stacksaga.datasource.r2dbc.database=order_service_event_store
stacksaga.datasource.r2dbc.username=${MYSQL_USER:root}
stacksaga.datasource.r2dbc.password=${MYSQL_PASSWORD:password}
| 1 | stacksaga.instance 属性用于配置 StackSaga 引擎的实例拓扑信息,例如集群 (cluster)、区域 (region) 与可用区 (zone)。这些属性用于在分布式环境中精准标识事务,以便进行故障重试与宕机恢复。请注意,如果您指定了除 default 之外的自定义 region,则必须声明自定义 SagaRegionResolver Spring Bean。在此示例中配置为 local-cluster、local-region 和 local-zone 仅供演示。 |
| 2 | stacksaga.domain-entity-scan 属性用于指定领域实体类所在的包路径。StackSaga 引擎将扫描该包并注册其中的领域实体类供事务流程使用。如果领域实体分布在多个包中,可以通过逗号分隔列出。 |
| 3 | stacksaga.datasource.r2dbc.* 属性配置响应式 R2DBC 连接,供 StackSaga 引擎向 MySQL 事件存储中持久化事务检查点 (Checkpoint) 与执行尝试历史记录。指定连接 URL、数据库名称、用户名和密码。建议将事件存储配置为独立于应用业务库的专属 Schema,以避免数据库连接池与事务锁竞争。 |
运行应用并测试事务流程 (Running the Application and Testing the Transaction Flow)
在 MySQL 实例运行的前提下启动 order-service。在首次启动时,框架会自动初始化事件存储数据库 Schema(通过 Liquibase)。
向 /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"]
}'
由于 MakePaymentExecutor 会随机模拟失败,请多次重复发送请求以观察*两种*执行结果。PlaceOrderHandler.onStateChanged(…) 回调会实时打印每一次状态流转,您可以在控制台中观察 Saga 的完整轨迹。
-
正常流程 (Happy Path)(支付成功)—— Saga 顺利执行至终态:
2026-09-17T23:49:26.609+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [nio-8080-exec-3] c.d.s.s.controller.PlaceOrderController : Transaction Id : orde-01a0b098-4e11-797f-bb12-c1b5c0af5df2
2026-09-17T23:49:29.316+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-4] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : PROCESSING
2026-09-17T23:49:31.142+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-4] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : PROCESSING
2026-09-17T23:49:34.083+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-4] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : PROCESS_COMPLETED
-
补偿流程 (Compensation Path)(支付失败)—— 引擎自动逆向回滚先前已预留的订单,最终以
REVERT_COMPLETED结束:
2026-09-17T23:46:02.870+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [nio-8080-exec-4] c.d.s.s.controller.PlaceOrderController : Transaction Id : orde-01a0b095-322f-7cda-884f-cf03ae37c864
2026-09-17T23:46:04.950+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-2] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : PROCESSING
2026-09-17T23:46:07.685+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-2] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : PROCESSING
2026-09-17T23:46:10.549+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-2] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : REVERTING
2026-09-17T23:46:13.212+05:30 INFO 10644 --- [stacksaga-framework-mysql-demo-servlet] [oundedElastic-2] c.d.s.saga.handler.PlaceOrderHandler : onStateChanged : REVERT_COMPLETED
按需查询事务状态 (Querying Transaction State On-Demand)
除了 PlaceOrderHandler.onStateChanged(…) 在执行期间以推模式实时接收事件通知外,您还可以在任何时刻通过 /api/v1/order/status 端点主动拉取查询事务的当前状态与完整执行历史:
该端点调用 sagaTemplate.getCurrentState(orderId).fetch() 直接从 MySQL 事件存储中检索 TransactionState 快照:
curl -X GET "http://localhost:8080/api/v1/order/status?order-id=orde-01a0b098-4e11-797f-bb12-c1b5c0af5df2"
-
已成功完成的事务 (Happy Path) —— 返回状态为
PROCESS_COMPLETED,且executionHistory中包含全部正向步骤:
{
"active": true,
"completed": false,
"current_domain_entity_state": {
"initialized_version": {
"major": 1,
"minor": 0,
"patch": 0
},
"payment_reference_id": null,
"product_items": [
"item-1",
"item-2"
],
"total_amount": 149.99,
"transaction_id": "orde-01a0b098-4e11-797f-bb12-c1b5c0af5df2",
"transaction_token": 8120724162383601234,
"user_validated": false,
"username": "john.doe"
},
"current_status": "PROCESS_COMPLETED",
"execution_history": {
"0": {
"current_status": "SUCCESS",
"event_name": "VALIDATED_USER",
"executed_datetime": "2026-09-17T18:19:29.300755Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.ValidateUserExecutor",
"executor_type": "QUERY_EXECUTOR"
},
"1": {
"current_status": "SUCCESS",
"event_name": "RESERVED_ORDER",
"executed_datetime": "2026-09-17T18:19:31.125944Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.ReserveOrderExecutor",
"executor_type": "COMMAND_EXECUTOR"
},
"2": {
"current_status": "SUCCESS",
"event_name": "MADE_PAYMENT",
"executed_datetime": "2026-09-17T18:19:34.070392Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.MakePaymentExecutor",
"executor_type": "COMMAND_EXECUTOR"
}
},
"freezing_reason": null,
"non_retryable_executor_exception": null,
"revert_hint_store": null,
"started_date_time": "2026-09-17T18:19:26.609259Z",
"transaction_id": "orde-01a0b098-4e11-797f-bb12-c1b5c0af5df2"
}
-
已补偿完成的事务 (Compensation Path) —— 返回状态为
REVERT_COMPLETED,既包含正向执行记录,也包含类型为DO_REVERT_MAIN的逆向补偿步骤:
curl -X GET "http://localhost:8080/api/v1/order/status?order-id=orde-01a0b095-322f-7cda-884f-cf03ae37c864"
{
"active": true,
"completed": false,
"current_domain_entity_state": {
"initialized_version": {
"major": 1,
"minor": 0,
"patch": 0
},
"payment_reference_id": null,
"product_items": [
"item-1",
"item-2"
],
"total_amount": 149.99,
"transaction_id": "orde-01a0b095-322f-7cda-884f-cf03ae37c864",
"transaction_token": 6469422457067738268,
"user_validated": false,
"username": "john.doe"
},
"current_status": "REVERT_COMPLETED",
"execution_history": {
"0": {
"current_status": "SUCCESS",
"event_name": "VALIDATED_USER",
"executed_datetime": "2026-09-17T18:16:04.903157Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.ValidateUserExecutor",
"executor_type": "QUERY_EXECUTOR"
},
"1": {
"current_status": "SUCCESS",
"event_name": "RESERVED_ORDER",
"executed_datetime": "2026-09-17T18:16:07.634302Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.ReserveOrderExecutor",
"executor_type": "COMMAND_EXECUTOR"
},
"2": {
"current_status": "FAILED",
"event_name": "FAILED_PAYMENT",
"executed_datetime": "2026-09-17T18:16:10.52115Z",
"execution_mode": "DO_PROCESS",
"executor": "com.demo.stage1.saga.executor.MakePaymentExecutor",
"executor_type": "COMMAND_EXECUTOR"
},
"3": {
"current_status": "SUCCESS",
"event_name": "REVERTED_ORDER",
"executed_datetime": "2026-09-17T18:16:13.189065Z",
"execution_mode": "DO_REVERT_MAIN",
"executor": "com.demo.stage1.saga.executor.ReserveOrderExecutor",
"executor_type": "COMMAND_EXECUTOR"
}
},
"freezingReason": null,
"non_retryable_executor_exception": {
"data": {
"@message": "insufficient balance",
"@exception_name": "java.lang.RuntimeException",
"@event_name": "FAILED_PAYMENT",
"@process_executor_name": "MakePaymentExecutor"
},
"event_name": "FAILED_PAYMENT",
"real_exception_message": "insufficient balance",
"executor_name": "MakePaymentExecutor",
"localized_message": "insufficient balance",
"message": "insufficient balance",
"real_exception_name": "java.lang.RuntimeException",
"suppressed": []
},
"revert_hint_store": {},
"started_date_time": "2026-09-17T18:16:02.873825Z",
"transaction_id": "orde-01a0b095-322f-7cda-884f-cf03ae37c864"
}
与在事务执行期间接收内存中推送事件的实时 onStateChanged 事件监听回调不同,getCurrentState(orderId).fetch() 是一种基于拉模式的主动查询,会按需从事件存储中检索持久化快照。
|
阶段二 [添加重试能力] (Stage-2 [Adding Retry Capability])
阶段一涵盖了正常流程与故障补偿。但在现实分布式生产环境中,您还会经常遇到值得通过重试而非立即回滚来解决的*瞬态故障 (Transient Failures)*——例如偶发网络抖动、下游服务短暂不可用或数据库死锁。在这一阶段,我们引入 StackSaga 重试子系统,该系统能够根据灵活配置的重试策略自动重新执行失败步骤,从而赋予事务弹韧性。
重试由专属子系统负责统一协调。启用该功能主要包含以下两个步骤:
-
创建 环形协调器应用 (Ring Coordinator Application):使用
stacksaga-ring-coordinator-spring-boot-starter构建一个与编排器应用并行的独立应用,负责在集群中各可用编排器实例之间动态协调与分配令牌区间 (Token Range)。 -
向编排器应用中添加
stacksaga-ring-coordinator-connector依赖,以建立与环形协调器的集成。该连接器负责接收由环形协调器从节点 (Slave) 分配给本实例的专属令牌区间。
创建环形协调器应用(主从一体模式)(Creating the Ring Coordinator Application (Master & Slave Hybrid))
创建一个新的 Spring Boot 应用,并在 pom.xml 中引入以下依赖以启用环形协调器功能。无需引入 Web 或其他非必要依赖,当然您也可以根据项目需要自由扩展其他依赖。
<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-ring-coordinator-spring-boot-starter</artifactId>
</dependency>
| 建议使用 StackSaga Initializer 获取相关依赖配置,以确保版本与依赖组合的一致性。 |
添加依赖后,在 application.properties 或 application.yml 中配置以下属性以启动环形协调器服务:
spring.application.name=order-service-ring-coordinator
#rsocket-server for master (1)
spring.rsocket.server.address=localhost
spring.rsocket.server.port=4455
#master and slave both (2)
stacksaga.coordinator.target-services=order-service
stacksaga.coordinator.instance-type=master,slave
#stacksaga-instance properties (3)
stacksaga.instance.cluster=local-cluster
stacksaga.instance.region=local-region
stacksaga.instance.zone=local-zone
#Master connection details for slave to connect to master (4)
stacksaga.coordinator.slave.target-master.port=4455
stacksaga.coordinator.slave.target-master.host=localhost
| 1 | spring.rsocket.server 属性用于配置环形协调器应用的 RSocket 服务端。您需要提供 RSocket 服务的监听地址与端口,供 Slave 实例连接至 Master 实例。 |
| 2 | stacksaga.coordinator.target-services 属性用于指定环形协调器所管理的目标微服务。在此示例中指定 order-service 为目标服务。stacksaga.coordinator.instance-type 用于指定当前环形协调器实例的角色类型,本例中针对本地单机测试将 master,slave 声明为同一实例。在生产环境中,通常部署一个独立的 Master 节点与多个 Slave 节点以获得更好的伸缩性与容灾高可用。 |
| 3 | stacksaga.instance 属性用于配置环形协调器应用的实例元数据。cluster 与 region 必须与编排器应用中配置的值完全一致,以确保二者处于相同的集群和区域以开展正常协作。有关区域解析规则,请参阅 SagaRegionResolver。 |
| 4 | 由于当前实例同时作为 Slave 运行,它需要连接至 Master 节点以接收令牌区间分配信息,因此通过 stacksaga.coordinator.slave.target-master 属性提供 Master 节点的连接配置。 |
在现有编排器中添加 stacksaga-ring-coordinator-connector
进入现有编排器应用的 pom.xml,添加以下依赖以启用与环形协调器的连接支持:
<dependency>
<groupId>org.stacksaga</groupId>
<artifactId>stacksaga-ring-coordinator-connector</artifactId>
</dependency>
随后在 application.properties 或 application.yml 文件中配置环形协调器连接属性:
#retry properties
stacksaga.sql.transaction.recovery.default.retry.delay=1m (1)
stacksaga.coordinator.connector.enabled=true (2)
stacksaga.coordinator.connector.master.requester.host=localhost (3)
stacksaga.coordinator.connector.master.requester.port=4455
# Optional: Recovery-Node batch streaming and concurrency tuning
#stacksaga.transaction.recovery-node.batch-size=100
#stacksaga.transaction.recovery-node.concurrency=8
| 1 | stacksaga.sql.transaction.recovery.default.retry.delay 属性设置停滞事务在被暴露给重试系统之前需等待的时间跨度。在此配置下,停滞事务将在 1 分钟后进入可重试候选集。您还可以针对特定领域配置独立的覆盖时长(例如 stacksaga.sql.transaction.recovery.order-domain.retry.delay=30s)、宕机恢复延迟以及事务最大生存周期。完整细节请参阅 SQL 数据库配置属性。 |
| 2 | stacksaga.coordinator.connector.enabled 属性用于在编排器应用中激活环形协调器连接器。将其设置为 true 会将标准编排器实例升级为能够接收环形协调器分配的令牌区间并执行重试任务的 重试节点 (Retry Node)。同时激活 stacksaga.transaction.recovery-node.* 流水线。 |
| 3 | 连接 Master 实例的主机名与端口。建立连接后,Master 会指派一个可用的 Slave 节点(当前单机测试下仅有 1 个),并在内部与该 Slave 节点建立通道并订阅分配的令牌区间信息。 |
模拟重试场景 (Simulating retry scenario)
接下来我们改造 ReserveOrderExecutor,模拟一个可通过重试恢复的瞬态故障(例如外部服务临时不可用)。我们将在第一次执行尝试时主动抛出 RetryableExecutorException,而在重试系统发起的第二次尝试时顺利执行,以直观验证重试机制的运作效果。
@SagaExecutor(executeFor = "order-service", value = "ReserveOrderExecutor")
public class ReserveOrderExecutor implements CommandExecutor<OrderDomainEntity> {
(1)
private final AtomicReference<Map<String, Integer>> counter = new AtomicReference<>(new HashMap<>());
@NonNull
@Override
public ProcessStepManager<OrderDomainEntity> doProcess(
OrderDomainEntity currentDomainEntityState,
ProcessStepManagerUtil<OrderDomainEntity> stepManager,
String idempotencyKey
) throws RetryableExecutorException, NonRetryableExecutorException {
(2)
if (counter.get().containsKey(currentDomainEntityState.getTransactionId())) {
counter.get().put(currentDomainEntityState.getTransactionId(), counter.get().get(currentDomainEntityState.getTransactionId()) + 1);
} else {
counter.get().put(currentDomainEntityState.getTransactionId(), 1);
throw RetryableExecutorException.of("Simulate a retryable exception for transactionId: " + currentDomainEntityState.getTransactionId());
}
//call the internal service and reserve the order
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
return stepManager.next(MakePaymentExecutor.class, "RESERVED_ORDER");
}
@NonNull
@Override
public SagaExecutionEventName doRevert(
NonRetryableExecutorException primaryExecutionException,
OrderDomainEntity finalDomainEntityState,
RevertHintStore revertHintStore,
String idempotencyKey
) throws RetryableExecutorException {
//call the internal service to revert the reserve order action
try {
//simulate some delay for make the request
Thread.sleep(new Random().nextInt(1000, 3000));
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
return SagaExecutionEventName.of("REVERTED_ORDER");
}
}
| 1 | 我们使用 AtomicReference 维护一个计数器 Map,记录每个事务 ID 的执行尝试次数。这仅用于测试模拟以确定何时抛出 RetryableExecutorException。在实际生产场景中,您通常会根据捕获的具体异常类型(如外部远程调用超时、临时连接拒绝)或业务状态来决定是否抛出 RetryableExecutorException。 |
| 2 | 在 doProcess 方法中检查当前事务 ID 的执行计数。若为第 1 次尝试(计数为 1),主动抛出 RetryableExecutorException 以模拟可重试的瞬态故障。若为第 2 次尝试(计数为 2),则继续正常业务流程并调用 stepManager.next 推进至下一步。通过这种方式,我们清晰地验证了重试机制:第一次尝试失败触发重试等待,第二次重试尝试成功完成。 |
阶段三 [通过链路追踪窗口实现可观测性] (Stage-3 [Observability via trace window])
在前面的阶段中,我们成功实现了具备故障补偿与重试机制的事务流。现在让我们为其赋予可观测性。StackSaga 链路追踪窗口 (Trace Window) 是一个安全的可视化 Web 门户,能够实时监控与图形化展现您的事务流程——它将每个事务绘制为其流转步骤的有向图,直观展示各步骤的执行状态以及执行期间发出的事件,帮助您快速定位瓶颈并全方位监控 Saga 的健康状况。在本阶段中,我们将其接入正在运行的编排器。
添加链路追踪窗口连接器 (Adding Trace Window Connector)
要启用与链路追踪窗口的集成,请在编排器应用的 pom.xml 中引入以下依赖:
<dependency>
<groupId>org.stacksaga</groupId>
<artifactId>stacksaga-trace-window-connector-servlet</artifactId>
</dependency>
随后在 application.properties 或 application.yml 中配置以下属性以开放链路追踪窗口 API:
#Trace-Window
stacksaga.trace-window.secure-api=false (1)
stacksaga.trace-window.enable-api=true (2)
| 1 | stacksaga.trace-window.secure-api 属性用于配置链路追踪窗口 API 是否需要鉴权与认证。设置为 false 表示链路追踪 API 无需认证即可直接访问,适合本地开发与调试。然而在生产环境中,强烈建议将其设置为 true 并配置访问令牌。了解更多 |
| 2 | stacksaga.trace-window.enable-api 属性用于在编排器应用中启用链路追踪窗口 API。该 API 默认处于禁用状态。 |
重启应用服务器后,打开浏览器访问 Trace-Window 页面,输入本地运行的服务地址 http://localhost:8080,并填入之前调用 /api/v1/order 接口所获取的事务 ID,即可实时图形化呈现与追踪该事务的完整执行轨迹。