agents24 仓库 saga-orchestration 技能进阶指南:生产级 Saga 编排器实现、补偿事务链与可观测性治理

发布时间:2026/9/9 13:44:31
agents24 仓库 saga-orchestration 技能进阶指南:生产级 Saga 编排器实现、补偿事务链与可观测性治理 agents24 仓库 saga-orchestration 技能进阶指南生产级 Saga 编排器实现、补偿事务链与可观测性治理【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents导读本文基于本仓库 saga-orchestration 技能 的进阶参考文档 references/advanced-patterns.md 展开。当你在微服务架构中用异步消息取代两阶段提交2PC、编排跨库存、支付、物流与通知的分布式长事务时本文提供可直接继承复用的完整 Saga 编排器抽象基类、按步骤独立的超时看门狗、幂等且总是成功的补偿事务链以及配套的 Prometheus 指标、卡死检测与 DLQ 恢复手段。读完本文你将掌握如何为每个 Saga 类型定义有序步骤与反向补偿、如何实现 per-step 超时调度、如何规避补偿无限卡死以及如何把整套编排器接进生产可观测性体系。一、仓库上下文这份进阶文档在技能体系中的位置本仓库是一个面向 Claude Code、Codex、Cursor、OpenCode、GitHub Copilot 与 Google Antigravity 的多宿主 Agent 插件市场后端开发领域由 backend-development 插件 承载其中saga-orchestration技能采用渐进式披露progressive disclosure组织内容文件定位SKILL.md技能入口适用场景、输入输出、Dos / Donts、故障排查、指向进阶文档references/details.md核心概念编排/编排式对比、状态机、三套基础模板Order Fulfillment 编排式、Choreography 编排式、幂等步骤守卫references/advanced-patterns.md生产级实现完整抽象基类、按步超时、银行转账补偿链、监控与 DLQ ——即本文主体三份文档共同配合例如 SKILL.md 故障排查一节中关于超时早于慢但合法的步骤完成的结论指向 advanced-patterns.md 的TimeoutSagaOrchestrator与STEP_TIMEOUTSdetails.md 则说明全局超时不可取每步延迟特征不同。advanced-patterns.md 的定位是从核心技能抽取的复杂实现供深度参考适合大多数 Saga 不需要、但一旦需要就必须可靠的生产路径。二、完整 Saga 编排器抽象基类状态转移、补偿顺序与事件发布SagaOrchestrator是全部 Saga 类型的抽象基类其职责边界在文档类注释中写得很清楚按顺序通过异步命令消息执行步骤、任何失败时按逆序触发补偿、每次状态转移后持久化 Saga、完成或失败时发布领域事件。每个具体的 Saga 子类只需回答两个问题saga_type唯一标识如OrderFulfillment与define_steps()有序步骤清单。2.1 状态机与数据模型Saga 实例生命周期由SagaState枚举约束五个状态与 details.md 中的状态表一一对应状态语义STARTEDSaga 已启动第一步已派发PENDING等待参与服务的步骤回执COMPENSATING某步失败正在回滚已完成步骤COMPLETED所有正向步骤成功FAILED补偿全部完成后 Saga 以失败告终from abc import ABC, abstractmethod from dataclasses import dataclass, field from enum import Enum from typing import List, Dict, Any, Optional from datetime import datetime, timedelta import uuid class SagaState(Enum): STARTED started PENDING pending COMPENSATING compensating COMPLETED completed FAILED failed dataclass class SagaStep: name: str action: str compensation: str status: str pending result: Optional[Dict] None error: Optional[str] None executed_at: Optional[datetime] None compensated_at: Optional[datetime] None timeout_at: Optional[datetime] None dataclass class Saga: saga_id: str saga_type: str state: SagaState data: Dict[str, Any] steps: List[SagaStep] current_step: int 0 created_at: datetime field(default_factorydatetime.utcnow) updated_at: datetime field(default_factorydatetime.utcnow)SagaStep值得细读每个步骤同时承载正向action命令名与逆向compensation命令名status、result、error、executed_at、compensated_at、timeout_at构成完整审计字段方便状态转移日志记录对应 SKILL.md 最佳实践中在每次变更时记录saga_id、step_name、old_state → new_state。Saga.timeout_at的默认值由field(default_factory...)注入为超时模式预留了扩展点。2.2 基类实现启动、回执处理与补偿class SagaOrchestrator(ABC): Base class for all saga orchestrators. Responsibilities: - Execute steps in sequence via async command messages - Trigger compensation in reverse order on any failure - Persist saga state after every transition - Publish domain events on completion and failure def __init__(self, saga_store, event_publisher): self.saga_store saga_store self.event_publisher event_publisher abstractmethod def define_steps(self, data: Dict) - List[SagaStep]: Define the ordered saga steps for this workflow. pass property abstractmethod def saga_type(self) - str: Unique identifier for this saga type (e.g., OrderFulfillment). pass async def start(self, data: Dict) - Saga: Start a new saga instance. saga Saga( saga_idstr(uuid.uuid4()), saga_typeself.saga_type, stateSagaState.STARTED, datadata, stepsself.define_steps(data) ) await self.saga_store.save(saga) await self._execute_next_step(saga) return saga async def handle_step_completed(self, saga_id: str, step_name: str, result: Dict): Handle a successful step reply from a participant service. saga await self.saga_store.get(saga_id) for step in saga.steps: if step.name step_name: step.status completed step.result result step.executed_at datetime.utcnow() break saga.current_step 1 saga.updated_at datetime.utcnow() if saga.current_step len(saga.steps): saga.state SagaState.COMPLETED await self.saga_store.save(saga) await self._on_saga_completed(saga) else: saga.state SagaState.PENDING await self.saga_store.save(saga) await self._execute_next_step(saga) async def handle_step_failed(self, saga_id: str, step_name: str, error: str): Handle a step failure and begin compensation. saga await self.saga_store.get(saga_id) for step in saga.steps: if step.name step_name: step.status failed step.error error break saga.state SagaState.COMPENSATING saga.updated_at datetime.utcnow() await self.saga_store.save(saga) await self._compensate(saga) async def _execute_next_step(self, saga: Saga): Publish the command for the current step. if saga.current_step len(saga.steps): return step saga.steps[saga.current_step] step.status executing await self.saga_store.save(saga) await self.event_publisher.publish( step.action, { saga_id: saga.saga_id, step_name: step.name, **saga.data } ) async def _compensate(self, saga: Saga): Execute compensation steps in reverse order. for i in range(saga.current_step - 1, -1, -1): step saga.steps[i] if step.status completed: step.status compensating await self.saga_store.save(saga) await self.event_publisher.publish( step.compensation, { saga_id: saga.saga_id, step_name: step.name, original_result: step.result, **saga.data } ) async def handle_compensation_completed(self, saga_id: str, step_name: str): Mark a compensation step done and check if all are finished. saga await self.saga_store.get(saga_id) for step in saga.steps: if step.name step_name: step.status compensated step.compensated_at datetime.utcnow() break all_compensated all( s.status in (compensated, pending, failed) for s in saga.steps ) if all_compensated: saga.state SagaState.FAILED await self._on_saga_failed(saga) await self.saga_store.save(saga) async def _on_saga_completed(self, saga: Saga): await self.event_publisher.publish( f{self.saga_type}Completed, {saga_id: saga.saga_id, **saga.data} ) async def _on_saga_failed(self, saga: Saga): await self.event_publisher.publish( f{self.saga_type}Failed, {saga_id: saga.saga_id, error: Saga failed after compensation, **saga.data} )从源码结构可以梳理出的关键契约每次转移都落盘start先save再派发handle_step_completed/handle_step_failed/_execute_next_step/_compensate均在状态改变后立即save保证崩溃后可从持久化状态恢复。补偿只针对已完成步骤_compensate从current_step - 1逆序回退到0且仅当step.status completed才发布补偿命令未开始的步骤pending和已失败步骤failed都被跳过。这与 details.md 的补偿规则表完全一致步骤从未启动→无需补偿步骤完成→运行补偿步骤失败→标记失败即可。补偿消息携带original_result补偿命令将上一正向步骤返回的结果一并转发参与服务可以据此精确撤销例如release_reservation使用command[original_result][reservation_id]否则补偿方无从知道要回滚什么资源。幂等协调采用saga_idsaga_id作为 correlation id 贯穿每条命令与事件负载正是 SKILL.md 最佳实践中每个事件和日志都要流过 saga_id的落地方式。一个值得注意的边界handle_compensation_completed中all_compensated的判定允许pending/failed步骤存在——因为未执行或已失败的步骤不需要补偿只要所有已完成的步骤都变成compensatedSaga 即进入FAILED。补赏期间的任何回执成功/失败都会走这条收尾路径。2.3 从基类到具体 SagaOrder Fulfillmentdetails.md 给出了该基类的典型子类OrderFulfillmentSaga其四步正逆命令形成闭合回滚路径reserve↔release、payment↔refund、shipment↔cancel、confirmation↔cancellation-notice可作为撰写新 Saga 时define_steps的参照模板其InventoryService.handle_release_reservation以ReservationNotFoundError → pass方式演示补偿幂等、总是发布完成事件与本文档银行转账链的做法如出一辙。三、按步骤独立超时TimeoutSagaOrchestrator3.1 为什么不能只有一个全局超时参与服务 SLA 千差万别支付通常 30 秒内出结果而物流面单可能需 15 分钟。使用全局超时会带来两种系统性风险全局超时过短→ 合法但缓慢的步骤如高峰期的create_shipment被误杀触发多余的补偿造成业务抖动——这是 SKILL.md 故障排查明确列出的一种线上问题全局超时过长→ 无响应的参与方把 Saga 拖成僵尸PENDING卡死无法被发现。SKILL.md 的 Donts 中不要使用全局超时每个步骤有不同的延迟特征正是STEP_TIMEOUTS字典存在的理由。3.2 实现调度器 看门狗回调TimeoutSagaOrchestrator通过组合一个scheduler延迟任务调度器实现按步看门狗每个步骤进入executing时计算自己的timeout_at并注册一个唯一 job若截止时步骤仍是executing看门狗回调自动把它转为失败并触发补偿。class TimeoutSagaOrchestrator(SagaOrchestrator): Extends the base orchestrator with configurable per-step timeouts. # Override per saga subclass as needed STEP_TIMEOUTS: Dict[str, timedelta] { reserve_inventory: timedelta(minutes2), process_payment: timedelta(minutes1), create_shipment: timedelta(minutes15), send_confirmation: timedelta(minutes2), } def __init__(self, saga_store, event_publisher, scheduler): super().__init__(saga_store, event_publisher) self.scheduler scheduler async def _execute_next_step(self, saga: Saga): if saga.current_step len(saga.steps): return step saga.steps[saga.current_step] step.status executing step.timeout_at datetime.utcnow() self.STEP_TIMEOUTS.get( step.name, timedelta(minutes5) ) await self.saga_store.save(saga) # Schedule the timeout watchdog await self.scheduler.schedule( job_idfsaga_timeout_{saga.saga_id}_{step.name}, handlerself._check_timeout, payload{saga_id: saga.saga_id, step_name: step.name}, run_atstep.timeout_at ) await self.event_publisher.publish( step.action, {saga_id: saga.saga_id, step_name: step.name, **saga.data} ) async def _check_timeout(self, data: Dict): Called by the scheduler when a step deadline is reached. saga await self.saga_store.get(data[saga_id]) step next((s for s in saga.steps if s.name data[step_name]), None) if step and step.status executing: await self.handle_step_failed( data[saga_id], data[step_name], fStep {data[step_name]} timed out after {self.STEP_TIMEOUTS.get(data[step_name])} ) async def handle_step_completed(self, saga_id: str, step_name: str, result: Dict): Cancel the timeout job before processing the success reply. await self.scheduler.cancel(fsaga_timeout_{saga_id}_{step_name}) await super().handle_step_completed(saga_id, step_name, result)实现细节与调优要点job 命名须全局唯一且可取消saga_timeout_{saga_id}_{step_name}同时作为调度注册与取消的键正常回执到达时先cancel再走基类成功路径避免成功回执已处理、看门狗却误触发补偿的竞态。兜底默认值未在STEP_TIMEOUTS中登记的步骤使用timedelta(minutes5)兜底。从源码结构可推断向现有 Saga 追加新步骤而忘记登记超时时5 分钟默认值会兜住但应在新步骤 SLA 明确后补登记。看门狗是幂等门卫_check_timeout回调执行时若步骤已经completed例如消息延迟导致成功回执与超时竞态step.status executing判定为 False直接跳过——不会对已完成步骤二次补偿。超时错误信息应携带超时阈值错误串内插值STEP_TIMEOUTS.get(...)可读、可告警、可定位是哪一步因哪个阈值被判定超时。回调入参约定_check_timeout通过调度器的payload拿到{saga_id, step_name}这意味着你的调度器需要支持把结构化负载回传给 handler若使用现有 Celery / APScheduler / 自研延迟队列需做一层适配。在写集成测试时文档与 SKILL.md 一致建议故意在每个步骤索引处注入失败验证a失败步骤之前的已完成步骤都按逆序收到补偿命令b超时看门狗只对executing步骤生效c正常回执会先取消 job。四、补偿事务链以银行转账 Saga 为例补偿是最关键、最难的代码路径SKILL.md 最佳实践。本仓库的高级文档用跨服务银行转账给出完整范式。4.1 步骤定义每个正向操作都有一一对应的逆向操作class BankTransferSaga(SagaOrchestrator): Saga for transferring funds between accounts across services. property def saga_type(self) - str: return BankTransfer def define_steps(self, data: Dict) - List[SagaStep]: return [ SagaStep( namedebit_source, actionAccountService.DebitAccount, compensationAccountService.CreditAccount # reverse the debit ), SagaStep( namecreate_transfer_record, actionLedgerService.CreateTransfer, compensationLedgerService.VoidTransfer ), SagaStep( namecredit_destination, actionAccountService.CreditDestinationAccount, compensationAccountService.DebitAccount # reverse the credit ), SagaStep( namenotify_parties, actionNotificationService.SendTransferConfirmation, compensationNotificationService.SendTransferFailureNotice ), ]四个正向步骤分别跨越账户扣款、账本登记、入账、通知反向路径对称且可解释扣款↔冲正贷记、登记↔作废、入账↔冲正借记、确认通知↔失败通知。注意文档特意为通知也配了补偿发送失败通知体现凡有副作用、皆可补偿的原则。4.2 参与服务的正向与补偿处理幂等 总是发布class AccountService: async def handle_debit_account(self, command: Dict): idempotency_key fdebit-{command[saga_id]}-{command[account_id]} existing await self.ledger.find_by_key(idempotency_key) if existing: await self._publish_completed(command, {transaction_id: existing.id}) return try: txn await self.ledger.debit( account_idcommand[source_account_id], amountcommand[amount], idempotency_keyidempotency_key ) await self._publish_completed(command, {transaction_id: txn.id}) except InsufficientFundsError as e: await self._publish_failed(command, str(e)) async def handle_credit_account(self, command: Dict): Compensation: credit back a previously debited account. idempotency_key fcredit-comp-{command[saga_id]}-{command[account_id]} existing await self.ledger.find_by_key(idempotency_key) if not existing: await self.ledger.credit( account_idcommand[source_account_id], amountcommand[amount], idempotency_keyidempotency_key ) # Always publish — even if already credited await self.event_publisher.publish(SagaCompensationCompleted, { saga_id: command[saga_id], step_name: debit_source })从这段示例中可以提炼出补偿代码必须遵守的三条纪律补偿必须幂等idempotency_key形如credit-comp-{saga_id}-{account_id}与正向路径的debit-{saga_id}-{account_id}命名空间隔离find_by_key命中即说明此前已冲正直接跳过实际写操作。这覆盖了参与者服务重启后重放命令超时误判后重复补偿等场景。补偿必须总是成功即使底层资源已经在目标状态例如已被贷记也必须无条件发布SagaCompensationCompleted。这正是 SKILL.md 故障排查中Saga 卡死在 COMPENSATING问题的解药——补偿处理器一旦抛出未捕获异常且不发布完成事件编排器将永远等不到all_compensated。补偿使用确定性键而非随机数键由saga_id与业务维度account_id派生任何一次重放都得到同一个键可在账本表中直接命中已执行的补偿记录。五、生产监控与故障自愈Prometheus 指标、卡死告警与 DLQ 恢复5.1 Prometheus 指标与插桩编排器先定义四类指标计数器Counter衡量吞吐与成功/失败量直方图Histogram衡量按saga_type与outcome切分的耗时分布Gauge 反映当前卡死数量。随后用InstrumentedSagaOrchestrator包装基类在关键钩子上打点from prometheus_client import Counter, Histogram, Gauge import time saga_started_total Counter( saga_started_total, Total sagas started, [saga_type] ) saga_completed_total Counter( saga_completed_total, Total sagas completed successfully, [saga_type] ) saga_failed_total Counter( saga_failed_total, Total sagas that failed after compensation, [saga_type] ) saga_compensating_total Counter( saga_compensating_total, Total sagas that entered compensation, [saga_type] ) saga_duration_seconds Histogram( saga_duration_seconds, Saga execution duration, [saga_type, outcome], buckets[1, 5, 15, 30, 60, 300, 600] ) saga_stuck_gauge Gauge( saga_stuck_count, Sagas stuck in COMPENSATING or PENDING threshold, [saga_type, state] ) class InstrumentedSagaOrchestrator(SagaOrchestrator): Wraps base orchestrator with Prometheus instrumentation. async def start(self, data: Dict) - Saga: saga_started_total.labels(saga_typeself.saga_type).inc() saga await super().start(data) saga._start_time time.monotonic() return saga async def _on_saga_completed(self, saga: Saga): duration time.monotonic() - getattr(saga, _start_time, 0) saga_completed_total.labels(saga_typeself.saga_type).inc() saga_duration_seconds.labels( saga_typeself.saga_type, outcomecompleted ).observe(duration) await super()._on_saga_completed(saga) async def _on_saga_failed(self, saga: Saga): duration time.monotonic() - getattr(saga, _start_time, 0) saga_failed_total.labels(saga_typeself.saga_type).inc() saga_duration_seconds.labels( saga_typeself.saga_type, outcomefailed ).observe(duration) await super()._on_saga_failed(saga)设计上的几个关键点hook 复用而非逐点插桩文档只覆写了start与两个收尾事件_on_saga_completed/_on_saga_failed成功路径与失败路径共用同一套插桩点耗时用time.monotonic()记录于saga._start_timemonotonic 不受系统时钟跳变影响。耗时直方图分桶贴合业务buckets[1, 5, 15, 30, 60, 300, 600]秒跨度从 1 秒覆盖到 10 分钟与文档中订单流转支付 1 分钟内、物流 15 分钟级的时间尺度匹配落在最后桶外的长尾可作为优化信号。每个指标都带saga_typelabel同一集群内跑着OrderFulfillment、BankTransfer等多种 Saga 时可独立聚合定位到底是哪一类 Saga 在失败。注意saga_stuck_gauge需要另外的巡检任务定期扫描持久化的 Saga 表把超龄PENDING/COMPENSATING实例计数上报指标定义中已通过注释声明阈值口径stuck in COMPENSATING or PENDING threshold。5.2 卡死与健康度的 PromQL 告警文档给出两条可直接进 Alertmanager 的规则# Alert: saga stuck in compensation for 10 min increase(saga_compensating_total[10m]) - increase(saga_failed_total[10m]) 0 # Alert: saga completion rate drops below 95% ( rate(saga_completed_total[5m]) / (rate(saga_completed_total[5m]) rate(saga_failed_total[5m])) ) 0.95解读这两条查询的语义第一条用进入补偿的增量 − 结束于失败的增量刻画正卡在补偿流程中的存量只要 10 分钟内进入补偿的 Saga 多于宣告失败的就说明有 Saga 滞留COMPENSATING——通常是某个补偿处理器抛异常且未发布SagaCompensationCompleted对应 SKILL.md 故障排查首条。第二条监控完成率success rate5m窗口内的完成量与完成失败总量之比跌破 95% 即告警。分母刻意不含仍在运行的实例避免误报长流程。5.3 DLQ 恢复指数退避重放与毒信隔离补偿处理器未捕获异常时消息落入死信队列DLQ。生产环境需要恢复 worker 定期重放 DLQ 中的补偿消息文档给出的实现采用指数退避 上限兜底class SagaDLQRecovery: Replays failed compensation messages from the dead-letter queue. MAX_RETRIES 5 BASE_DELAY_SECONDS 10 async def process_dlq_message(self, message: Dict, attempt: int): delay self.BASE_DELAY_SECONDS * (2 ** attempt) if attempt self.MAX_RETRIES: await self._move_to_poison_queue(message) await self._alert_on_call(message) return await asyncio.sleep(delay) try: await self.event_publisher.publish(message[original_topic], message[payload]) except Exception as e: await self.process_dlq_message(message, attempt 1)设计要点退避序列为 10 / 20 / 40 / 80 / 160 秒BASE_DELAY_SECONDS * 2 ** attempt第 0 次尝试前睡 10 秒第 1 次前 20 秒……符合指数退避定义asyncio.sleep表明 recovery worker 运行在事件循环内不会阻塞其他重放任务。达到MAX_RETRIES5 次后不再盲目重试转入_move_to_poison_queue毒信队列隔离并触发_alert_on_call人工介入告警。这与 details.md 补偿规则表中补偿失败→退避重试→DLQ→人工介入告警的路径一致。重放对象是消息而非业务方法publish(message[original_topic], message[payload])说明 DLQ 里存的是原 topic 原负载重放即重新投递完整经过参与服务的幂等守卫安全无副作用。异常被吞并递归attempt 1形成有界重试若异常发生在publish之前的读操作需要把读取也纳入 try 或做整体事务封装这属于对生成代码的使用边界可在工程化落地时按需调整。六、扩展阅读进阶文档的See Also线索advanced-patterns.md 结尾的 See Also 把本技能与同插件内的相邻技能串成一条完整的数据流写作时值得一并参读核心概念与决策表saga-orchestration/SKILL.md含核心概念导读、Dos/Donts、故障排查与 references/details.md编排式 vs 编排式对比、状态机表、补偿规则表、Order Fulfillment / Choreography / 幂等守卫三套模板。阅读顺序建议SKILL.md → details.md → advanced-patterns.md。cqrs-implementation 技能SKILL.md每个 Saga 步骤完成后由投影更新读模型形成写侧走 Saga 编排、读侧走 CQRS 投影的完整形态。event-store-design 技能SKILL.md以事件存储充当 Saga 的持久化saga_store与event_publisher从而获得不可变审计日志与回溯重放能力——这是文档基类中saga_store/event_publisher两个依赖注入点的天然落点。若业务场景需要更高阶的工作流引擎Temporal、Conductor可进一步阅读同插件的 workflow-orchestration-patterns 技能其与手写 Saga 的取舍在 SKILL.md 的 Related Skills 一节有索引。整体上backend-development 插件 内的多个 Agent如 backend-architect可作为使用这些模式时的协作方。七、落地清单与自检把本文四段实现接入你的服务时建议按下列清单逐项核对每一条都能在仓库文档中找到对应出处继承SagaOrchestrator只实现saga_type与define_steps()正逆命令一一配对。每个参与服务的补偿处理器遵守幂等键 总是发布SagaCompensationCompleted宁可多余发布不可静默吞异常。用TimeoutSagaOrchestrator替换基类需要你的基础设施里有延迟调度器按步骤 SLA 登记STEP_TIMEOUTS并验证正常回执会cancel看门狗 job。接入InstrumentedSagaOrchestrator与四条 PromQL 告警让COMPENSATING卡死与完成率下滑在 10 分钟内被感知。为补偿消费者配 DLQ SagaDLQRecovery毒信人工介入从机制上杜绝无限补偿。参照 details.md 补上在每个步骤索引注入失败的集成测试用测试锁定逆序补偿与超时行为。以上全部代码实现均以 Markdown 形式沉淀在本仓库plugins/backend-development/skills/saga-orchestration/目录下可作为后续 Agent 会话按需加载的参考依据无需重复推导。【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考