StackSaga-Kafka 框架 (异步架构)

概述 (Overview)

StackSaga-Kafka 框架 (stacksaga-kafka-spring-boot-starter) 是 StackSaga 生态中的异步传输实现。 它利用 Kafka 消息骨干网取代同步的服务间调用,从而在微服务系统中为长事务 (LRT - Long-Running Transaction) 实现完全事件驱动的 Saga 执行。

与各个服务独立发布并响应领域事件的简单协同 (Choreography) 模式不同,StackSaga-Kafka 遵循*集中式编排 (Centralized Orchestration)* 模型:由单个编排器服务 (Orchestrator Service) 掌控完整的 Saga 生命周期,通过专用的 Kafka 主题 (Topic) 向目标 Worker 服务发出命令,并通过框架托管的基于领域的专用回调主题收集响应。 这种分离既保留了 Kafka 的高扩展性与服务解耦能力,又继承了 Saga 编排模式的确定性控制流与强审计追踪特性。

重要公告 — Kafka 共享消费组支持 (KIP-932)

StackSaga 现已在框架的编排器端和 Worker 端全链路端到端支持 Kafka 4 队列与共享组 (Kafka 4 Queues with Share Groups - KIP-932)。 共享消费组允许多个消费者实例以单条记录确认的方式协同从同一主题读取消息,因此 Saga 处理吞吐量不再受限于主题的分区数量 (Partition Count)。

可通过 groupType 属性在每个端点或事件管理器上独立启用 — 无需重新分区或重构主题结构。

该框架构建在 Spring Boot 之上,并与 StackSaga 核心引擎集成,提供事件溯源、状态管理、重试调度以及集群协同支持。 每次 Saga 执行都进行持久化 — SagaDomainEntity(Saga 的聚合根与载荷载体)的每一次状态转换都会通过 stacksaga-database-support 写入事件存储库,因此任何进行中的事务都可以在任何时间点进行恢复、重试或审查。

核心特性与优势:

  • 异步 Saga 编排 — 编排器通过 Kafka 向 Worker 服务派发命令,仅在收到 Worker 的应答后才推进 Saga 状态,等待期间不会占用或阻塞任何工作线程。

  • 专为支持 KIP-932 共享组的 Kafka 4 打造 — 编排器与 Worker 均既支持经典消费者组协议,也支持 Kafka 4 共享组协议 (KIP-932)。这是一种队列式消费模型,其并行度*不受分区数量限制*。 协议可通过 groupType 属性在每个端点和每个事件管理器上按需选择,无需重新分区即可将 Saga 处理吞吐量扩展至超越分区数。

  • 正向推进与逆向回滚恢复 — 原生内置补偿事务支持:当某个 Saga 步骤失败时,框架会以逆序调用 EventManager 上的 onNextRevert(),依次触发先前已成功完成的每个服务的撤销回滚动作。

  • 基于领域的专用回调拓扑 — 出向命令主题按 Worker 操作定义;入向应答则每个 Saga 领域使用一个专用回调主题。 拥有三个 Saga 领域的系统仅需使用三个回调主题,使 Kafka 分区管理清晰可控。

  • 通过事件溯源实现持久状态保障 — 完整的 SagaDomainEntity 载荷和状态(STARTED → IN_PROGRESS → COMPLETED | FAILED → COMPENSATING → COMPENSATED)在每次状态流转时均被持久化,支持任意时间点恢复与完整审计链路。

  • 基于令牌环分区的分布式重试 — 失败或停滞的事务由负责的编排器实例负责重试,该归属通过基于 Murmur3 的令牌环分区确定,并由 stacksaga-ring-coordinator 统一协同。

  • 全面支持响应式与命令式 — 框架在 Worker 端的 @SagaEndpoint 抽象同时支持响应式 (Mono/Flux) 与阻塞式处理器实现,开发者可自由选择编程模型而不影响全局 Saga 协调。

术语表 (Glossary)

以下技术术语贯穿本模块文档。在阅读技术细节之前,熟悉这些术语将大幅降低理解门槛:

术语 定义

长事务 (LRT - Long-Running Transaction)

跨越多个微服务且可能需要数秒、数分钟或更长时间才能完成的业务分布式事务。 LRT 需要持久化状态管理、分布式协同以及对局部故障的补偿回滚支持。

Saga 领域 (Saga Domain)

由其 DomainEntity 子类标识的特定 LRT 类型。 同一类型的所有 Saga 实例(例如所有 PlaceOrder 下单事务)都属于同一个 Saga 领域。 每个领域拥有专属的 EventManager、Kafka 回调主题以及一组跨度 (Spans)。

跨度 (Span)

Saga 中的单个原子执行步骤。 每个跨度通过专用的 Kafka 主题调用一个目标 Worker 服务。 一个 Saga 在正向主流程中由一个或多个顺序跨度组成,每个跨度可带有一个对应的可选补偿跨度。

Saga 执行协调器 (SEC - Saga Execution Coordinator)

负责推动 Saga 向前演进的核心引擎组件:评估 EventManager、向 Worker 分发命令、接收应答、持久化状态转换以及在需要时触发补偿。 StackSagaKafkaTemplate 是开发者面向 SEC 的核心交互入口。

领域实体 (Domain Entity)

单个 Saga 实例的聚合根。 在整个 Saga 生命周期中承载累积的完整业务载荷与当前执行状态。 每次状态转换都会作为快照持久化到事件存储库中。

事件管理器 (EventManager)

编排器端用于定义 Saga 领域路由逻辑的核心组件。 决定在每个跨度成功完成后接下来触发哪个主题 (onNext()),并为补偿序列提供路由钩子 (onNextRevert())。

执行器 / 端点 (Executor / Endpoint)

执行特定 Saga 步骤业务逻辑的 Worker 端处理器。 实现为 QueryEndpoint(只读,无补偿)或 CommandEndpoint(修改状态,具备补偿动作)。

补偿 (Compensation)

当 Saga 发生故障时,撤销先前已成功完成步骤的逆向流程。 由 EventManager 中的 onNextRevert() 与 CommandEndpoint 实现中的 undoProcess() 按逆序驱动执行。

事件存储库 (Event Store)

存储所有 Saga 状态转换与领域实体快照的持久化存储层。 由 stacksaga-database-support 提供底层支撑。 支持按时间点恢复、重试与完整的审计跟踪。

补偿提示存储库 (Revert Hint Store)

在补偿序列中向前传递元数据的键值存储库。 在某次 undoProcess() 调用期间写入的值可供后续补偿步骤读取。

主题键 (Topic Key)

分配给每个 Saga 主题的唯一 float 浮点数值。 在 Kafka 消息头中使用,取代原始长字符串主题名,保持网络传输消息紧凑并解耦主题重命名。

环形协调器 (Ring Coordinator)

管理用于重试归属判定的令牌环分区的独立微服务 (stacksaga-ring-coordinator)。 在编排器实例之间分配 Murmur3 令牌子区间,使每个停滞事务由且仅由一个实例重试,无需分布式锁。

为什么选择 StackSaga-Kafka 而非原生 Kafka?

在 Saga 中仅使用原生 Kafka 进行服务间通信并不能自动获得 Saga 编排能力。 原生 Kafka 拓扑要求每个参与微服务都必须感知全局事务上下文:当前处于第几步、哪些步骤已成功、失败时回滚哪些操作以及如何跟踪全局进度。 这些编排逻辑不可避免地泄漏到各个业务服务中,导致系统脆弱且难以维护迭代。

以下对比清晰展示了 StackSaga-Kafka 所填补的架构鸿沟:

集中式事务状态管理

关注点 详细对比

原生 Kafka

Kafka 仅提供持久的消息传递,没有业务事务概念。缺乏内置机制来知晓“第 2 步已完成(共 5 步)”或“事务当前处于补偿回滚状态”。

StackSaga-Kafka

编排器通过严密的状态机跟踪全局生命周期:`STARTED → IN_PROGRESS → COMPLETED

FAILED → COMPENSATING → COMPENSATED`。SagaDomainEntity 是全局单一事实来源:承载事务累积载荷与当前执行游标。每一步完成或失败时,状态变更原子化持久化到事件存储库。

核心优势

弹性与容错保障

关注点 详细对比

原生 Kafka

Kafka 提供消息持久性和至少一次投递保证,但当消息被消费后下游业务处理失败时,Kafka 不知道该如何应对。重试策略、超时和回滚决策必须完全自行编写。

StackSaga-Kafka

框架实现了可配置的重试窗口、边界保护阈值以及自动补偿触发机制。如果 Worker 服务返回失败,Saga 引擎立即启动 EventManager 的逆向遍历,按逆序依次调用每个已完成前向步骤的 onNextRevert()。针对临时瞬态故障(超时、基础设施波动),重试子系统利用 Murmur3 令牌环调度器从最后一个未确认步骤重新调用 Saga。

核心优势

即使在局部微服务故障下,也能跨服务保障最终数据一致性,无需在各业务服务中侵入式嵌入重试或回滚逻辑。

极大简化业务开发

关注点 详细对比

原生 Kafka

开发者必须在参与业务事务的每个服务中手动实现状态机、幂等检查、步骤路由和补偿协调。

StackSaga-Kafka

编排器端抽象(StackSagaKafkaTemplate、EventManager、SagaDomainEntity)与 Worker 端抽象(@SagaEndpoint)为每个角色提供了清晰严谨的契约。开发者只需声明步骤*做什么*以及*按何种顺序*执行;路由、持久化、重试和补偿调度全由框架自动打理。

核心优势

消除海量模板样板代码,将分布式事务控制收敛至极小的可测试层。

运营可观测性

关注点 详细对比

原生 Kafka

仅提供消费者延迟 (Consumer Lag) 和 Broker 级别的基础监控指标,无法感知“订单 X 停滞在第 3 步”或“订单 Y 的补偿退款失败”。

StackSaga-Kafka

stacksaga-trace-window-connector 暴露 API 供 StackSaga Trace Window 界面使用,提供按事务维度的步骤级追踪、执行时间线、故障断点、重试次数和补偿状态详情。

核心优势

运维与研发团队可以在业务事务维度精确诊断停滞或失败的 Saga,而不仅停留在 Kafka 消费者级别。

架构核心组件 (Components)

StackSaga-Kafka 框架组件总览图

StackSaga-Kafka 框架划分为两个职责清晰的工件模块。 它们之间完全通过 Kafka 进行通信:编排器向 Worker 指定的主题发送出向命令消息,Worker 将应答消息回传到框架托管的基于领域的专用回调主题。

stacksaga-kafka-orchestrator(编排器)

stacksaga-kafka-orchestrator 是添加到*编排器服务 (Orchestrator Service)* 的核心运行时依赖项 — 负责发起并推动 Saga 全生命周期的单一服务。

它提供以下核心抽象:

  • StackSagaKafkaTemplate — 启动新 Saga 执行或恢复已恢复事务的核心入口。接收已初始化的 SagaDomainEntity 并移交至 Saga 引擎。

  • EventManager — 定义 Saga 的有序步骤与补偿步骤。引擎在前向演进时顺序调用 onNext(),在补偿回滚时以逆序调用 onNextRevert()。

  • SagaDomainEntity — Saga 实例的聚合根。承载累积的业务载荷和当前执行状态,在每次状态流转时序列化持久化到事件存储库中。

要将您的微服务配置为编排器服务,首先在项目中添加 stacksaga-kafka-orchestrator-starter 依赖:

<dependencyManagement>
    <dependencies>
        <dependency> <!--用于 stacksaga 依赖版本统一定义-->
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>

<dependencies>
    <dependency>
        <groupId>org.stacksaga</groupId>
        <artifactId>stacksaga-kafka-orchestrator-starter</artifactId>
    </dependency>
</dependencies>
该框架面向 Kafka 4 与 Spring Boot 4,并直接与 Spring for Apache Kafka 的监听器容器深度集成。 除经典消费者组协议外,全面支持用于队列消费模型的 Kafka 4 共享组协议 (KIP-932),可通过 groupType 属性在每个事件管理器与每个端点上自由选择。 业务代码支持响应式 (Mono/Flux) 与命令式阻塞风格编写 — 两者均为一等公民并获完全支持。
如果您的编排器服务同时需要充当由另一个编排器拥有的其他 Saga 领域的 Worker,无需单独引入 Worker 依赖 (stacksaga-kafka-worker-starter),因为编排器依赖已经传递依赖了 Worker 模块。

此外,编排器服务还必须包含:

  • stacksaga-database-support — 提供事件存储库集成(按配置支持 MySQL、PostgreSQL、Oracle、Cassandra 或 ScyllaDB),用于持久化 SagaDomainEntity 快照和状态转换。

  • stacksaga-ring-coordinator-connector (可选,分布式重试所需) — 将编排器实例连接至 Ring Coordinator 服务,注册为重试节点并接收分配的 Murmur3 令牌子区间,使重试子系统能够重新调用属于该实例的失败事务。

stacksaga-kafka-worker(工作节点)

stacksaga-kafka-worker 是添加到参与编排器发起的 Saga 的每个 Worker 服务 (工作服务)(亦称目标服务或执行服务)的依赖项。

Worker 依赖项提供:

  • Kafka 监听器基础设施,订阅框架主题配置为该服务分配的主题。

  • @SagaEndpoint — 业务开发者在此抽象中实现该服务所负责的每个 Saga 步骤的正向处理逻辑与补偿回滚逻辑。

  • 自动响应分发 — 在端点处理器完成(成功或失败返回)后,Worker 会自动将处理结果通过编排器的领域专用回调主题回传给框架。业务开发者无需手动编写这一应答流。

Worker 服务不与事件存储库或环形协调器直接交互。 它们在 Saga 视角下是无状态的:接收命令、执行相关业务操作并返回结果。

单个 Worker 服务可以同时处理来自多个编排器领域(即多种不同 Saga 类型)的步骤,框架会自动将消息路由到正确的 @SagaEndpoint 实现。

在 Worker 项目中添加以下依赖:

<dependencyManagement>
    <dependencies>
        <dependency> <!--用于 stacksaga 依赖版本统一定义-->
            <groupId>org.stacksaga</groupId>
            <artifactId>stacksaga-bom</artifactId>
            <version>1.0.0-SNAPSHOT</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>
<dependencies>
    <dependency>
        <groupId>org.stacksaga</groupId>
        <artifactId>stacksaga-kafka-worker-starter</artifactId>
    </dependency>
</dependencies>