我把结算对账流程改成 Temporal 有状态工作流后丢单问题彻底消失了说出来有点丢人我们跑了快两年的结算对账隔三差五就“丢单”。不是数据库丢了也不是消息队列没发到而是流程跑到一半卡住了重启后没人记得它跑到哪一步。凌晨两点被叫醒查账成了家常便饭。直到我把这套流程整体迁到 Temporal 上才发现原来对账这件事天生就不该交给无状态的消息消费者来做。老方案Cron Kafka 无状态 Worker先交代一下背景。我们做的是聚合支付对账每天需要把几十家渠道前一天的交易流水跟我们自己平台订单逐笔核对。核对成功后生成差异单再推给财务系统。老架构很简单00:30 cron 触发 - Kafka topic reconcile - 5 个无状态 consumer 并行处理每个 consumer 拿到消息后开始跑从渠道接口拉流水本地按订单号分组调用内部订单服务查状态写对账结果表发事件通知财务看起来每个步骤都很清晰但问题藏在“无状态”三个字里。Consumer 一重启当前内存里的上下文全没了。如果第三步调用订单服务超时消息重试 3 次后 ack 失败就会进死信队列。人工去死信队列里捞数据再手动补账补错一笔就是一笔坏账。更麻烦的是重复消费。Kafka 重试时 consumer 可能已经执行到第四步幂等没做好结果表就写了两遍。财务看到“对账成功金额”和“订单金额”对不上群里直接 全员。为什么消息队列不适合做“长事务”我后来复盘发现根因不是 Kafka 不好而是我们把它用错了场景。消息队列的设计是“尽快处理、尽快确认”天然鼓励短、快的任务。而对账流程是典型长事务跨多个外部服务有明确的步骤依赖失败后需要精确回到断点继续而不是从头再来。无状态 consumer 没有“回到断点”这个能力。你能做的只有两种选择重试整条消息副作用不可控ack 后进死信流程断掉需要人工缝合。无论哪种都无法保证“恰好一次”的执行语义。方案用 Temporal 把工作流状态持久化我们最终选型了 Temporal。它不是又一个消息队列而是一个有状态的工作流执行引擎。核心思想是把业务流程写成 Workflow把真实 side effect 写成 ActivityWorkflow 的状态会被持久化到数据库Activity 失败可以按策略重试整个流程可以挂起、恢复、重放。对账流程改写成 Temporal 后架构变成Scheduler 触发 Workflow - Temporal Server 调度 - Worker 执行 Activity关键变化是流程本身成为了一个长生命周期的对象。如果 worker 重启Temporal Server 会根据持久化的历史把 workflow 重新构建到断点。Activity 失败会按配置重试Activity 成功后不会再重复执行因为历史已经记录。重复记账、丢单、死信队列的问题从根上消失了。代码落地Workflow Activity Saga下面是核心代码片段已经脱敏。语言用的是 Java SDK我们的技术栈是 Spring Boot。1. 定义 Workflow 接口WorkflowInterfacepublicinterfaceReconcileWorkflow{WorkflowMethodReconcileResultreconcile(StringchannelCode,LocalDatebillDate);SignalMethodvoidapproveManualDiff(StringdiffId);QueryMethodReconcileStatusgetStatus();}reconcile是主入口approveManualDiff是信号方法用于人工确认差异后继续流程getStatus可以在运行期查询当前进度。2. Workflow 实现publicclassReconcileWorkflowImplimplementsReconcileWorkflow{privatefinalActivityOptionsoptionsActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofMinutes(5)).setRetryOptions(RetryOptions.newBuilder().setInitialInterval(Duration.ofSeconds(2)).setMaximumInterval(Duration.ofMinutes(1)).setMaximumAttempts(5).build()).build();privatefinalReconcileActivitiesactivitiesWorkflow.newActivityStub(ReconcileActivities.class,options);privateReconcileStatusstatusnewReconcileStatus();OverridepublicReconcileResultreconcile(StringchannelCode,LocalDatebillDate){status.setState(PULLING);// 1. 拉取渠道流水ListChannelBillbillsactivities.fetchChannelBills(channelCode,billDate);status.setBillsFetched(bills.size());// 2. 查询本地订单status.setState(MATCHING);ListDiffRecorddiffsactivities.matchBillsWithOrders(channelCode,bills);status.setDiffsFound(diffs.size());// 3. 对账记账status.setState(BOOKING);activities.writeReconcileResult(channelCode,billDate,diffs);// 4. 通知财务系统status.setState(NOTIFYING);activities.notifyFinance(channelCode,billDate);status.setState(DONE);returnnewReconcileResult(true,diffs.size());}OverridepublicvoidapproveManualDiff(StringdiffId){activities.confirmManualDiff(diffId);}OverridepublicReconcileStatusgetStatus(){returnstatus;}}注意Workflow 里不能有阻塞 IO、随机数、当前时间等不可重放操作。所有真实动作都委托给 Activity。3. Activity 实现补偿用 Saga对账流程里第四步“通知财务”如果失败不应该回滚前面的记账但需要在状态上标记为“已记账未通知”。更复杂的场景比如记账后需要撤销我们用 Saga 补偿。ActivityInterfacepublicinterfaceReconcileActivities{ListChannelBillfetchChannelBills(StringchannelCode,LocalDatebillDate);ListDiffRecordmatchBillsWithOrders(StringchannelCode,ListChannelBillbills);voidwriteReconcileResult(StringchannelCode,LocalDatebillDate,ListDiffRecorddiffs);voidnotifyFinance(StringchannelCode,LocalDatebillDate);voidconfirmManualDiff(StringdiffId);CompensatingvoidrollbackReconcileResult(StringchannelCode,LocalDatebillDate);}Saga 模式在 Temporal 里很简单用Saga类把补偿 Activity 注册进去主流程异常时统一触发。SagasaganewSaga(newSaga.Options.Builder().build());try{activities.writeReconcileResult(channelCode,billDate,diffs);saga.addCompensation(activities::rollbackReconcileResult,channelCode,billDate);activities.notifyFinance(channelCode,billDate);}catch(ActivityFailureExceptione){saga.compensate();throwe;}迁移过程没有一刀切上线不是直接全量替换。我们分了三步影子跑新 Temporal 流程和老 Kafka 流程并行执行只对比结果不写真实账。小渠道灰度挑两家日单量少的渠道先切到 Temporal观察一周。全量切换验证无误后所有渠道和日期全部迁移。灰度期间发现一个很重要的问题老的渠道接口偶尔会 5xx 几分钟。以前这种情况下 Kafka consumer 会快速耗尽重试次数进死信Temporal 则把 Activity 挂起按指数退避重试等服务恢复后自动继续。这让我第一次意识到“重试”和“恢复”其实是两回事。踩过的几个坑1. Activity 必须幂等Temporal 能保证 Activity 历史只记录一次但 Activity 内部如果调用外部 HTTP 接口网络抖动时仍可能实际执行多次。我们对每个 Activity 都加了幂等键比如channelCode billDate stepName外部服务根据这个键去重。2. Workflow 不能随意调用 Thread.sleepWorkflow 里用Workflow.sleep(Duration.ofHours(1))而不是Thread.sleep因为前者会被 Temporal 记录为事件真正执行时不会占用 worker 线程后者会卡住 worker 线程池。3. 版本兼容性很重要第一次改 Workflow 代码后正在运行的老版本实例直接报错。后来每次修改都加Workflow.getVersion()分支保证老实例按旧逻辑跑完新实例走新逻辑。intversionWorkflow.getVersion(use-new-matching,Workflow.DEFAULT_VERSION,1);if(version1){diffsactivities.matchBillsWithOrdersV2(channelCode,bills);}else{diffsactivities.matchBillsWithOrders(channelCode,bills);}写在最后迁移 Temporal 后我们对账流程的丢单率从之前的每月若干笔降到了零。更重要的是on-call 的凌晨电话几乎没了流程卡住就卡住它会自动重试我们只需要在白天看看 Temporal Web UI 里的异常列表。不是说消息队列不能做对账而是当业务流程本身有状态、有步骤、需要断点续传时把工作流状态交给专门的引擎来管比自己用数据库消息队列硬拼可靠得多。如果你的业务里也有“跑一半卡住的批处理”不妨试试 Temporal 这种思路。代码改动量不算小但换来的睡眠时间和财务老师的信任我觉得值。