架构演进 - StackSaga-Kafka 框架 (异步)

以下各节跨三个部署阶段逐步介绍 StackSaga-Kafka 架构。 每个阶段均构建于前一阶段之上,渐进式注入生产就绪能力。 这种分阶段方式使团队能够从最小化初始拓扑平滑起步,并随着业务规模与高可用需求的增长进行系统加固。

  1. 阶段 1:基础设置 — 基于 Kafka 的编排器与 Worker 通信。

  2. 阶段 2:重试就绪设置 — 通过环形协调器实现分布式重试。

  3. 阶段 3:监控可观测设置 — 通过 Trace Window 实现 Saga 级别全链路可观测性。

服务角色分类:编排器、Worker 与标准通用服务

在深入探讨架构之前,首先必须理清 StackSaga 生态系统中的不同服务类型。 理解服务的分类*方式*,将有助于后续章节的内容(分阶段架构、依赖关系表、请求处理流程)清晰映射到您自身的微服务系统上。

编排器 (Orchestrator) 与 Worker 是 StackSaga 生态系统中的角色标签,而非独立的新增部署单元或专用微服务。 系统中的每个微服务(order-service、payment-service、user-service 等)本质上首先都是一个*标准通用微服务 (Standard Utility Service)*:它一如既往地暴露自身的 REST/gRPC 端点、执行自身的业务逻辑,并照常为其现有的消费者提供服务。

向这些既有服务引入 StackSaga 依赖并不会取代或限制其原有功能,而是在其上叠加了额外的 Saga 角色能力。 例如,如果 order-service 引入了 stacksaga-kafka-orchestrator,它依然保留全部原有能力(暴露自身端点、调用其他内部服务等),同时*额外*获得了驱动特定 Saga 的编排器角色。 在 StackSaga 生态语境下,该服务被称为*编排器服务 (Orchestrator Service)* — 但该身份仅针对该特定的 Saga 而言。 Worker 服务同理:payment-service 继续充当正常的业务服务,同时通过 @SagaEndpoint 接收 Saga 命令,从而获得了 Worker 服务 (Worker Service) 的标签。

这是 StackSaga 的最大架构优势之一:它可以零侵入、零破坏地引入到任何现有微服务中。 无需从业务领域中强行抽离逻辑、无需部署全新的专用编排服务,也无需重构现有端点和业务逻辑。 您只需在既有微服务中引入对应的 StackSaga 依赖,它便可在现有职责*之外*自然承担编排器或 Worker 角色。 架构演进是平滑增量的,而非破坏性重构。

引入这些角色标签纯粹是为了让 Saga 拓扑结构更加易于理解与推演。 系统中的每个服务均属于以下三类之一(或在兼具通用职责时兼任其一):

角色 识别特征

编排器服务 (Orchestrator service)

被赋予端到端驱动 Saga 业务事务全生命周期额外角色的任何既有服务。通过 stacksaga-kafka-orchestrator 依赖(以及用于事件存储持久化的 stacksaga-database-support)识别。在 Saga 内部,它拥有 StackSagaKafkaTemplate、EventManager 以及全部 onNext() / onNextRevert() 步骤控制逻辑 — 在 Saga 之外,原服务功能保持完全不变。

工作服务 (Worker service)

被赋予参与 Saga 并与编排器交换命令/应答消息额外角色的任何既有服务。通过 stacksaga-kafka-worker 依赖以及注册的 @SagaEndpoint 处理器识别。在处理 Saga 命令之外,它继续作为标准通用服务正常运行。

标准通用服务 (Standard utility service)

每个服务的基础底色。*未引入*任何 StackSaga 依赖且从未被编排器调用的服务纯粹保留为标准通用服务 — 其行为不受 StackSaga-Kafka 任何影响。

单个微服务通常同时兼具标准通用服务与 StackSaga 角色(例如,payment-service 保留其常规 API,同时通过 @SagaEndpoint 监听并处理 Saga 命令)。 编排器/Worker 标签仅描述服务*在给定 Saga 中的参与职责*,而非独立的物理部署单元。 每个 Saga 业务事务中仅由一个服务担任编排器。

有关每个角色对应的具体依赖项,请参阅下文的阶段 1 依赖表。

阶段 1 — 基础设置 (Stage 1: Basic Setup)

该阶段建立了最小可行拓扑:一个能够启动和驱动 Saga 的编排器服务,以及一个或多个能够通过 Kafka 接收命令并返回结果的 Worker 服务。 此阶段尚未配置重试协调与外部监控。

StackSaga-Kafka 架构:阶段 1 — 基础设置

依赖项 (Dependencies)

服务角色 依赖工件 核心用途

编排器 (Orchestrator)

stacksaga-kafka-orchestrator

提供 Saga 引擎、StackSagaKafkaTemplate、EventManager 以及用于出向命令主题的 Kafka 生产者。

编排器 (Orchestrator)

stacksaga-database-support

提供事件存储库适配器,用于持久化 SagaDomainEntity 状态和执行历史。

工作节点 (Worker)

stacksaga-kafka-worker

提供 Kafka 消费者监听基础设施和 @SagaEndpoint 抽象。

请求处理与执行流程 (Execution Flow)

  1. 入向 HTTP 请求(例如 POST /order)到达编排器服务的 OrderController。

  2. 控制器实例化 PlaceOrderDomainEntity(订单提交 Saga 的 SagaDomainEntity 子类),填充其初始载荷,并调用 StackSagaKafkaTemplate.init(…​).startWith(..).execute();。 从此时起,Saga 引擎接管全权控制。

  3. 引擎评估 PlaceOrderEventManager 以确定首个步骤,并传入当前领域实体调用 onNext()。

  4. onNext() 的实现将命令消息投递到为目标服务分配的出向 Kafka 主题(例如支付服务的命令主题)。 该消息携带 Saga 步骤标识符与已序列化的领域实体载荷。

  5. 引入了 stacksaga-kafka-worker 依赖的 Worker 服务通过注册的 @SagaEndpoint 处理器消费其分配主题中的消息。 它执行业务操作(例如扣除支付款项),并将应答消息发送回框架托管的*基于领域的专用回调主题*。 对于 PlaceOrderDomainEntity 领域,这是专属于该领域的单一回调主题。

  6. 编排器的内部 Kafka 消费者从领域回调主题中读取应答消息。 如果步骤成功,引擎更新领域实体状态,将新快照持久化到事件存储库中,并通过再次调用 onNext() 推进至下一步骤。

  7. 如果 Worker 应答失败,引擎将 Saga 状态流转为 FAILED,持久化失败状态,并通过逆序对每个先前已完成的步骤调用 onNextRevert() 启动补偿回滚。 每个 onNextRevert() 向对应的 Worker 主题发送一条补偿命令。 Worker 通过其专用的 @SagaEndpoint 补偿处理器处理补偿命令,并通过相同的领域回调主题回传应答。

  8. 每次状态流转(步骤完成、失败、补偿启动、补偿完成)均通过 stacksaga-database-support 写入事件存储库。

上图为简洁起见展示了单个跨度 (MakePaymentSpan)。 在实际业务中,一个 Saga 可以跨越不同 Worker 服务的任意数量跨度,每个服务拥有自己的命令主题。 所有应答(无论由哪个 Worker 发送)均回流至单个基于领域的专用回调主题。 针对 GroupType.CONSUMER,回调主题消费者配置为 auto.offset.reset=earliest(GroupType.SHARE 共享消费组不使用该属性 — 参见 编排器监听器模型)。Saga 步骤的至少一次处理由框架自身的重试系统保障,而非依赖 Kafka 的消息重投递 — 参见下文的 阶段 2。 框架为每条消息提供一个基于跨度的幂等键,以便在消费者再平衡或偏移量提交延迟产生重复投递时,业务能通过幂等处理识别并安全防御,而不是二次执行 Saga 步骤。框架不会自动静默丢弃重复消息,利用幂等键保持业务幂等性是开发者的责任。
编排器上的回调主题与 Worker 上的命令主题均通过 Spring for Apache Kafka 监听器容器进行消费。每个容器均可采用经典的 Kafka 消费者组 (consumer-group) 协议或 Kafka 4 共享组 (share-group) 协议 (KIP-932) — 这是一种并行度不受分区数限制的队列式模型,可通过 groupType 属性在每个事件管理器和每个端点上自由选择。参见 编排器监听器模型 与 Worker 监听器模型。

双重 LRT 模型:连续直通处理 (Continuous STP) 与可暂停人机协同 (Pausable HITL)

StackSaga-Kafka 编排器引擎原生支持两种执行动态:

  • 连续型长事务 (Continuous LRT / 直通式处理 - STP):在完全自动化、不间断的处理流水线中,编排器向 Worker 主题派发命令消息并顺序消费其应答,中途无暂停。整个跨服务事务以 Kafka 事件流速度从发起高速运行至最终完成或补偿。

  • 可暂停长事务 (Pausable LRT / 业务等待状态与人机协同 - HITL):当事务必须暂停等待外部异步事件(例如等待第三方支付网关 Webhook、仓储物理发货确认、或主管/KYC 人工审核)时,事务进入预期的 PAUSED 等待状态。编排器在事件存储库中记录等待状态,并释放 Worker 线程与 Kafka 消费者资源。一旦回调或审核通过事件到达,执行将安全地从暂停里程碑精准恢复。

阶段 2 — 重试就绪设置 (Stage 2: Retry-Ready Setup)

阶段 2 引入了分布式重试能力。 若无此阶段,因瞬态基础设施故障(例如 Worker 服务临时不可用、Kafka Broker 分区短暂停摆或网络超时)而停滞的 Saga 将无限期保持未完成状态。 阶段 2 使系统针对此类故障具备自我愈合 (Self-Healing) 能力。

StackSaga-Kafka 架构:阶段 2 — 重试就绪设置

此阶段引入的新组件

服务角色 新增组件 核心用途

环形协调器服务 (Ring Coordinator Service)

stacksaga-ring-coordinator

管理令牌环的独立微服务。跟踪可用的编排器实例,在其间分配 Murmur3 令牌子区间,并在实例加入或离开集群时处理区间再平衡。详见 基于重试协调器的事务重试架构 (Transaction Retry Architecture With Retry Coordinator)。

编排器 (Orchestrator)

stacksaga-ring-coordinator-connector

将编排器实例连接至环形协调器(通过 RSocket request-stream),接收并持有其分配的令牌子区间,并驱动本地重试调度器扫描事件存储库中 Murmur3 哈希值落入本实例持有区间的停滞事务。

重试调度机制 (Retry Mechanism)

通过向编排器服务添加 stacksaga-ring-coordinator-connector,该实例被提升为 重试节点 (Retry Node)。 环形协调器为每个已注册的编排器实例分配 Murmur3 令牌环的一个连续子区间。 每个实例上的重试调度器定期扫描事件存储库中处于非终态(如 IN_PROGRESS、COMPENSATING)且其事务 ID 哈希落入本地持有区间的事务。 一旦发现此类事务,调度器将其重新提交给 Saga 引擎,从最后一个未完成的步骤恢复重新执行。

这种分区机制确保在多实例编排器部署中,每个停滞事务由且仅由一个实例进行重试 — 杜绝重试任务重复执行,且无需分布式锁。 当实例重启或新实例加入集群时,环形协调器会自动再平衡令牌范围,并通过现有的 RSocket 流下发新的分区分配。

环形协调器是一个独立的微服务 (stacksaga-ring-coordinator-spring-boot-starter),需单独部署。 它是一个极轻量级的协同服务,不参与具体业务逻辑或 Kafka 消息流传输。
建议深入阅读 基于重试协调器的事务重试架构 (Transaction Retry Architecture With Retry Coordinator) 以深入了解重试机制及环形协调器的工作原理。

阶段 3 — 监控与全链路可观测性设置 (Stage 3: Monitoring & Observability)

阶段 3 引入了可观测性层,使 StackSaga Trace Window 能够从编排器服务中实时查询实时与历史 Saga 执行数据。

StackSaga-Kafka 架构:阶段 3 — 监控与重试就绪设置

此阶段引入的新组件

服务角色 新增组件 核心用途

编排器 (Orchestrator)

stacksaga-trace-window-connector

暴露一组内部 API(供 StackSaga Trace Window UI 消费),直接从事件存储库中调取按事务维度的执行链路、步骤级时间线、故障详情、重试历史与补偿状态。

呈现的可观测能力

引入链路追踪窗口连接器后,StackSaga Trace Window 提供:

  • 按 Saga 维度的执行拓扑图,显示每个步骤、其执行时间戳、耗时、状态及任何错误载荷。

  • 补偿追踪链路,显示触发了哪些 onNextRevert() 调用、执行顺序以及是否回滚成功。

  • 重试审计日志,显示进行了多少次重试尝试、哪个重试节点处理了每次尝试以及最终结果。

  • 进行中 Saga 的实时事务健康状态。

该层不改变 Kafka 拓扑或重试行为 — 它是被动的可观测性连接器,仅从既有事件存储库读取数据并通过 API 暴露给 Trace Window 界面。

完成阶段 3 后,完整的 StackSaga-Kafka 部署即可全面生产就绪:支持异步 Saga 执行、带令牌环分区的分布式重试以及基于 Trace Window 的全景可观测性。

Saga 事务状态机 (Saga Transaction State Machine)

每个 Saga 实例都会流转经历一组严密定义的状态。 理解这些状态及其流转触发条件,对于在 Trace Window 中解读事务历史以及在端点和 EventManager 中实现正确的异常处理至关重要:

状态 类型 触发条件

STARTED

中间态

成功调用 StackSagaKafkaTemplate.execute() 或 executeAsync(),且首个跨度已加入调度派发队列。

IN_PROGRESS

中间态

收到首个 Worker 应答,引擎推进至下一个跨度。

COMPLETED

终态 — 成功

最后一个跨度成功完成后,从 onNext() 返回 actionUtil.complete()。

FAILED

中间态

从 Worker 的 doProcess() 抛出 NonRetryableExecutorException,或从 EventManager 的 onNext() 抛出任何异常(或通过 actionUtil.error(…​) 返回)。 立即触发逆向补偿序列。

COMPENSATING

中间态

紧随 FAILED 状态。引擎开始以逆序调用 onNextRevert() 并向各 Worker 派发补偿命令。

COMPENSATED

终态 — 成功

所有补偿跨度均成功执行完毕。

补偿失败 (Compensation Failed)

终态 — 故障

从 EventManager 的 onNextRevert() 抛出未捕获异常,或某个补偿步骤耗尽了配置的最大重试上限。 有关重试上限和死信配置,请参阅 stacksaga-database-support。

任何处于非终态(STARTED、IN_PROGRESS、FAILED、COMPENSATING)的 Saga,如果发生停滞,都有资格被环形协调器重试调度器恢复重试。 仅有三种终态 — COMPLETED、COMPENSATED 以及补偿失败 (Compensation Failed) — 被视为最终永久敲定,不再触发任何重试。