Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

第 6 章 任务处理方法论:从短任务、长任务到可恢复的 Workflow 与 Agent 协作

任务处理解决的不是某个业务对象如何流转,而是一段工作如何被拆分、调度、执行、暂停、恢复和治理。本章以供应商数据同步为主案例,结合离线计算、流处理和 Agent 协作,建立同步请求、异步消息、Batch / 分布式 Job、Workflow、Stateful Agent Runtime 与 Multi-Agent 的统一选型框架。

前面的章节分别讨论了关键业务动作的一致性,以及业务对象跨阶段的生命周期。本章处理第三类问题:一个可执行任务如何可靠跑完,或者如何在没有一次性“跑完”的情况下持续保持正确。扫描超时订单、生成日报、同步供应商库存、重建索引、执行历史数据修复、训练模型、运行 Spark 作业,以及让代码 Agent 修改多个文件,都需要明确执行边界、进度记录和恢复机制。

从表面看,这些任务的技术栈差异很大:有的使用 HTTP,有的使用消息队列,有的运行在 Kubernetes Job 或数据计算集群上,有的由 Workflow 引擎驱动,有的还需要大模型选择工具和调整路径。但它们共享一个事实:任务一旦跨越请求、进程、时间、步骤、系统或角色边界,就不能再把调用栈当作唯一状态,也不能把“进程退出”当作完成证明。

本章不把某个框架当成答案,而是先回答五个问题:任务的输入和成功判据是什么?中断后从哪里继续?重复执行会不会产生额外副作用?谁拥有最终事实?自动路径什么时候必须暂停并交给人或另一种执行单元?这些问题回答清楚后,方案才有可能从同步请求逐步演进为异步任务、Batch、Workflow 或 Agent Runtime,而不是在故障发生后被动堆叠重试队列和运维脚本。

6.1 问题定义:任务为什么会失控

6.1.1 短任务与长任务的区别不在耗时

许多团队用“几秒以内是短任务、超过几分钟是长任务”来划分任务。这种经验可以作为容量讨论的线索,却不是架构边界。真正重要的是任务是否需要跨边界持续推进。

这里的边界包括:

  • 请求边界:调用方收到响应后,工作是否还要继续。
  • 进程边界:执行进程重启后,是否能从外部状态恢复,而不是依赖内存。
  • 时间边界:任务是否可能等待分钟、小时、天甚至更长时间。
  • 步骤边界:任务是否包含多个有不同失败语义的阶段。
  • 系统边界:任务是否调用外部 API、数据库、计算集群或人工服务。
  • 角色边界:任务是否需要审批、复核、转派或人工确认。

一个耗时两秒但会产生不可逆外部副作用的付款请求,仍然需要幂等和结果未知处理;一个运行一小时但完全可重算、无副作用的离线计算,也许只需要 Batch 和分片。相反,一个只运行十秒的代码 Agent,如果需要修改文件、执行测试、等待用户确认后继续,就已经具有长任务的状态与治理特征。

因此,本章使用如下定义:短任务是可以在一次执行上下文中完成,并且失败后可以安全地整体重试的工作;长任务是必须跨越至少一个执行边界,且需要把进度、等待、恢复或人工决策持久化的工作。这个定义比固定时长更能解释为什么同一个接口在业务规模增长后会突然需要 Run、Checkpoint 和操作台。

6.1.2 典型错误:把执行问题误认为组件问题

任务最初往往只是一个函数、一个 Cron 或一条消息消费逻辑。数据量小、依赖稳定时,这种实现能够工作;一旦执行时间变长、对象数量变多或外部依赖开始抖动,问题会以相似的方式出现:

  • 同步请求一直不返回,线程、连接和客户端重试被长期占用。
  • 一个 Consumer 同时完成拉取、转换、发布、通知和补偿,失败后无法判断应从哪里恢复。
  • 全量任务中途失败,只能从头执行,重复写入和重复调用不断累积。
  • 任务状态散落在日志、队列和开发者记忆中,值班人员无法判断它是否仍在推进。
  • 下游限流或主库繁忙时,重试反而放大流量,最后影响正常业务。Google SRE 将这种正反馈称为级联故障,并特别强调重试、队列和资源耗尽之间的相互放大关系。[2][3]

这些问题不能只靠“把消费者扩容”“换成更快的数据库”或“再加一个重试次数”解决。一个任务系统需要把意图、执行事实、执行结果、失败原因和下一步动作分开记录。消息队列可以负责把工作交给 Worker,却不能替代任务状态;日志可以解释发生了什么,却不能作为可靠的恢复点;模型输出可以帮助 Agent 做判断,却不能自动获得生产副作用的权限。

6.1.3 典型长任务场景

旧版本中被删减的场景仍然有助于建立方法的覆盖范围。它们可以归纳为五类:

场景任务为什么会变长首要治理问题
供应商同步、ETL、索引重建对象多、外部源不稳定、需要分片和质量校验Checkpoint、数据版本、发布门禁、对账
模型训练与模型交付训练耗时长,资源可能被抢占,结果还要评估和审批实验状态、产物版本、指标门槛、回滚
Spark 离线计算与 Flink 流处理前者由 Stage 和 Shuffle 推进,后者长期持有状态分区、状态快照、资源隔离、恢复
代码、调研和运维 Agent输入不确定,需要工具调用、计划调整和人工判断工具权限、上下文、预算、Checkpoint、审查
审批、发布、迁移和平台运维有等待、角色边界和不可逆动作状态机、定时器、审批、暂停和人工接管

供应商同步是本章的主案例。模型训练和 Spark 任务说明 Batch 的边界不只存在于业务中台;Flink 的有状态处理则说明“长”也可以意味着作业长期运行,而不是一次执行很久。Apache Spark 以可容错的 RDD 和血缘重算来处理节点故障,Flink 则通过状态快照和 Checkpoint 恢复有状态算子;两者都说明恢复能力必须建立在明确的中间事实之上,而不是依赖某个 Worker 恰好没有重启。[9][10]

Agent 场景又增加了一层复杂度。代码 Agent 可能要理解需求、扫描仓库、规划修改、调用编辑工具、运行测试、读取失败输出,再决定返工;调研 Agent 可能并行检索资料、去重、交叉验证并由 Reviewer 收敛结论;运维 Agent 可能需要查询日志、诊断指标、生成变更计划,等待人工批准后再执行。它们的共同点不是“用了大模型”,而是执行路径包含不确定性,同时仍然需要传统任务系统的状态、权限、审计和恢复骨架。

6.1.4 本章的范围与非目标

本章聚焦以下问题:

  1. 如何识别任务的执行跨度、数据规模、流程确定性、认知复杂度和副作用风险。
  2. 如何选择同步请求、异步消息、Batch、Workflow、单 Agent、Stateful Agent Runtime 或 Multi-Agent。
  3. 如何设计 Task、Run、Step、Shard、Checkpoint、Lease、状态历史和结果账本。
  4. 如何处理重复投递、结果未知、局部失败、限流、背压、乱序、发布异常和人工接管。

本章不展开某个具体 Workflow 或 Agent 框架的 API,也不把模型提示词写作当作任务治理。第 4 章深入 Saga、Outbox、补偿和最终一致性;第 5 章深入业务对象的长生命周期状态机;本章只在这些主题交界处保留任务执行所必需的契约,并把 Agent 当成一种可能的执行单元,而不是一种可以替代业务事实的万能架构。

6.2 约束与指标:先明确什么算“跑得对”

6.2.1 用复杂度画像代替技术名词

任务选型前,建议先为任务建立一张复杂度画像。至少记录以下九个维度:

维度评审问题设计含义
执行跨度能否在一次请求和一个进程内闭环?决定是否需要异步交接与持久状态
对象规模处理一条记录、一个批次还是数十亿事件?决定分页、分片和资源调度
流程确定性路径是否固定,还是根据中间结果选择下一步?决定 Batch / Workflow 与 Agent 的边界
状态需求是否需要等待、暂停、恢复和历史查询?决定 Run、Step、Checkpoint 与状态历史
依赖复杂度是否调用多个系统、外部机构或人工角色?决定超时、结果未知、补偿与审批
副作用风险重复执行会不会重复扣费、发布或修改生产数据?决定幂等键、版本条件和人工门禁
认知复杂度是否需要阅读非结构化输入、归纳和动态规划?决定是否引入单 Agent 或 Agent Runtime
协作复杂度一个执行者能否承担规划、执行和审查?决定是否拆为多个角色化 Agent
治理复杂度是否需要配额、审计、回滚、SLO 和合规留痕?决定控制面、操作台和权限模型

这九个维度不需要一开始都量化成精确分数,但必须写出证据和假设。例如“每天同步五千万对象”是规模假设;“供应商每分钟最多 600 次调用”是外部约束;“价格波动超过上一版本 30% 必须人工确认”是业务门禁。容量数字如果只是方案假设,要明确标记,运行后的实际吞吐、P99、队列长度和最老任务年龄再用于校正。[1]

6.2.2 成功不能只有一个布尔值

长任务至少要区分四种成功:

  • 受理成功:任务定义、输入范围和权限检查通过,Run 已创建。
  • 执行成功:Worker 已完成某个 Step 或 Shard,并把结果与状态一起落地。
  • 业务成功:输出满足业务质量门槛,正式版本已经发布或业务对象已收敛。
  • 治理成功:所有失败、跳过、人工豁免和外部差异都有记录,后续责任边界清楚。

例如供应商同步接口返回 HTTP 202,只代表受理成功,不能让调用方把它展示为“库存已更新”;一个 Shard 的 Fetch 成功,也不意味着数据已经通过校验和发布;一次 Run 标记 SUCCEEDED,还应能解释隔离了多少坏数据、是否存在容忍范围内的差异,以及这些差异由谁负责处理。

指标应围绕这些层次建立:

指标类别示例需要回答的问题
完成性输入数、成功数、失败数、跳过数是否所有对象都有结论
正确性质量错误率、版本冲突数、对账差异输出是否满足不变量
时效性Run 时长、最老 Shard 年龄、等待时间是否在截止时间和 SLO 内推进
恢复性Checkpoint 命中率、重试成功率、恢复耗时故障后是否能从安全位置继续
资源性并发、CPU、内存、连接池、外部配额任务是否挤压在线流量
治理性人工任务逾期、回滚次数、审计缺失系统是否可控、可解释、可追责

6.2.3 Task、Run、Step、Shard 与 Checkpoint

无论最终使用脚本、消息队列、Kubernetes Job 还是 Workflow 引擎,都建议先持久化执行事实。一个通用模型包含五层:

实体含义必备字段示例
Task可复用的任务定义任务类型、输入契约、规则版本、资源策略
Run某次实际执行实例Run ID、输入快照、发起人、截止时间、配置版本
StepRun 内具有独立责任的阶段阶段状态、输入输出引用、重试预算、门禁
Shard可独立调度和重试的工作切片输入范围、租约、幂等范围、Checkpoint
Checkpoint已经落地且可验证的恢复位置游标、版本、输出摘要、提交时间、校验和

Task 表达“要做什么”,Run 表达“这次究竟处理了什么”。Run 应保存输入范围、规则与配置版本、发起来源、开始时间、截止时间和最终摘要。不能用一个全局任务状态覆盖历史运行记录,否则无法区分今天正在执行的任务和昨天遗留的失败。

Step 的边界应按责任和失败语义划分,而不是按代码函数随意切分。供应商同步可以把 Fetch、Normalize、Validate、Publish 和 Reconcile 分开;Fetch 失败可以重试外部请求,Validate 失败可以隔离坏对象,Publish 失败则需要保护正式版本。阶段边界越清晰,恢复时越不需要重复已经验证过的工作。

Shard 是扩展性与故障隔离的核心。常见切法包括供应商 / 商家维度、ID 范围、时间窗口、分页游标和哈希桶。分片太大,失败恢复成本高;分片太小,调度开销、状态数量和数据库写入压力又会过大。一个好分片应能估算工作量、单独重试、追溯输入范围,并且不会因并行而破坏同一业务对象的顺序约束。

Checkpoint 必须在副作用与状态事实都落地后推进。正确顺序通常是:读取输入 → 执行幂等写入或生成暂存结果 → 保存输出摘要和版本 → 在同一事务或可验证协议中推进 Checkpoint。只把页码打印到日志,或在落库前更新游标,都不能安全恢复。

6.2.4 状态迁移与版本契约

Run 可以经历 PENDING、PLANNED、RUNNING、PAUSED、SUCCEEDED、FAILED、CANCELLED、MANUAL_REVIEW;Shard 可以经历 PENDING、RUNNING、RETRYABLE、UNKNOWN、SUCCEEDED、FAILED 和 CANCELLED。状态不是为了增加字段,而是为了避免把“已入队”误判成“已完成”,也避免把“一个分片失败”误判成“整次运行已经无价值”。

每次迁移都应携带原因、发生时间、操作者或执行器标识、关联错误、版本号和审计引用。状态历史最好追加保存,当前状态服务于调度,历史服务于审计、复盘和容量优化。任务定义、分片规则、质量阈值和发布策略发生变化时,新的 Run 必须固化配置版本;历史 Run 按创建时的版本解释,不能在重试时悄悄读取一份已经变化的全局配置。

6.3 参考架构:调度、执行、状态与观测如何协作

6.3.1 控制面与数据面分离

一个可演进的任务系统不必一开始就有独立平台,但职责需要清晰。可以把它分成控制面和数据面:

控制面:Trigger / Scheduler / Policy / Operator Console
  -> 创建 Run、生成分片、检查配额、暂停或取消
  -> 读取状态、批准发布、触发受控重试

数据面:Queue / Dispatcher / Worker Pool / External Systems
  -> 领取 Shard、调用依赖、生成结果、提交 Checkpoint
  -> 执行幂等写入、返回错误和心跳

事实层:State Store / Object Ledger / Input Snapshot / Audit Trail
  -> 保存 Run、Step、Shard、版本、错误、外部请求号和人工动作

观测层:Metrics / Logs / Traces / Alerts / Reconciliation
  -> 解释推进速度、积压、失败原因、版本差异和恢复时间

调度器只负责创建运行实例、规划分片和发出可执行意图,不应在内存里长时间循环处理全部数据。Worker 只负责领取可执行分片、在租约内完成工作、持久化结果并续租或释放,不应自行猜测全局进度。状态存储保存权威的运行事实;消息队列负责分发,不应成为唯一状态来源。Airflow 的 DAG 与 Task 模型也强调,DAG 负责依赖、重试和时间控制,Task 才是实际执行单元;这正是控制面与数据面分工的一个具体例子。[8] DolphinScheduler 的官方项目资料同样把 DAG 编排、Worker 隔离、暂停恢复和批量回填作为调度平台能力,说明调度平台负责执行治理,业务任务仍要自己定义输入、结果和幂等边界。[18]

6.3.2 可靠交接:意图、消息与状态同时可追踪

如果请求需要异步交接,建议先在本地事务中写入 Run 意图,再写入 Outbox 事件,由 Relay 投递到队列。这样可以避免“数据库已创建任务但消息没发出”或“消息先发出但任务记录不存在”的断裂。Transactional Outbox 只解决本地业务事实与事件之间的双写窗口,Relay 仍可能在发送成功、更新投递状态前崩溃,因此消费者必须按事件 ID 或业务幂等键吸收重复。[5][6] 如果一个 Step 还要协调多个业务服务,则应把每个服务的本地事实和补偿边界显式记录,不能把 Outbox 误当成跨服务全局事务;Saga 的价值正在于把长事务拆为本地事务,并为已完成步骤定义业务补偿。[4]

消息契约至少应有 event_id、task_id、run_id、step_id、shard_id、业务主键、幂等键、输入版本、发生时间和过期时间。对于大对象,消息只保存对象存储引用与校验和;不能把可能无限增长的输入和上下文直接塞进队列。消费者确认成功的条件也要写清楚:是完成一次外部调用,还是完成本地事务、写入结果和推进 Checkpoint?如果只是调用成功但状态未落地,重新投递仍是正常路径。

6.3.3 Lease、心跳与重新领取

Worker 领取 Shard 时写入 owner_id、lease_until、attempt 和 heartbeat_at。在租约内,Worker 可以续租;进程崩溃或网络隔离后,租约到期,调度器将 Shard 标记为可恢复并重新投递。租约不是锁住整个任务的长期数据库事务,而是一种“当前执行者暂时负责推进”的可过期声明。

重新领取意味着执行至少可能一次,所以必须配合幂等和版本校验。一个 Worker 可能已经完成外部副作用,却在写成功状态前失联;另一个 Worker 随后重新执行。如果下游有稳定幂等键,可以查询并吸收这次重复;如果没有查询能力,就必须把结果标为 UNKNOWN,进入对账或人工队列。把“租约过期”直接等价为“前一次一定没有执行”是危险的。

6.3.4 查询、管理和数据接口

任务系统通常需要三类接口:提交接口、查询接口和管理接口。提交接口返回 run_id、受理状态和查询地址;查询接口返回 Run、Step、Shard 的进度、最近错误、下一次动作和是否需要人工处理;管理接口支持暂停、取消、恢复、范围重试、跳过和回滚,但每个动作都要受权限、状态机和审计约束。

对外暴露的状态不能只返回 success / failed。至少要区分 ACCEPTED、RUNNING、PARTIAL_FAILURE、SUCCEEDED、FAILED、PAUSED、CANCELLED 和 MANUAL_REVIEW。调用方得到的是受理结果,还是业务完成结果,必须在接口命名、字段和文档中明确,否则任务系统会把后端的异步不确定性转化成前端的错误承诺。

6.4 方案谱系与选型框架

6.4.1 七种主流执行方案

从简单到复杂,可以把常见方案归纳为七类。它们不是互斥的产品,而是对应不同主导矛盾的执行形态。

方案解决的主导矛盾获得的能力主动牺牲的能力升级信号
同步短任务请求内快速闭环调用简单、结果直接长等待和局部恢复请求超时、连接长期占用
异步消息 / Consumer主链路不应阻塞解耦、削峰、快速接入多阶段治理和细粒度审计Consumer 出现多个阶段和补偿
Batch / 分布式 Job大量对象如何推进分片、并行、Checkpoint动态规划与人工流程需要门禁、审批或发布版本
阶段化 Workflow多阶段如何治理等待、重试、审计、人工接管部分自由度和平台成本流程定义频繁变化或需动态推理
单 Agent不确定输入如何理解语义理解、工具选择、动态计划长任务可恢复性和确定性需要跨会话恢复或高风险副作用
Stateful Agent Runtime智能任务如何持续推进状态、Checkpoint、预算、审查实现与观测成本一个 Agent 无法兼顾规划、执行、审查
Multi-Agent多职责如何分工并行探索、专业角色、交叉审查通信、协调、Token 和终止成本角色边界明确且协作收益可验证

选择顺序不应从“要不要用 Kafka、Airflow 或 Agent”开始,而应从复杂度来源开始:先看是否跨边界,再看是否要批量恢复,再看是否需要阶段治理,最后才看是否需要动态推理和角色协作。Anthropic 对 Workflow 与 Agent 的区分也采用了类似原则:预定义代码路径适合可预测流程,模型动态决定过程和工具使用时才进入 Agent 语境;即使使用 Agent,也应从最简单、可组合的结构开始。[14]

6.4.2 同步、异步与 Batch 的边界

同步请求适合校验、查询、单次规则计算和小范围写入。它的价值是及时给出确定结果,而不是把所有工作都放在请求线程里。一旦执行时长不可预测、需要大量 I/O 或调用方只需要“已受理”,就应建立异步交接边界。

异步消息适合一条消息可以独立完成的动作,例如发送通知、刷新单个商品索引或补写非关键投影。消息体需要携带事件 ID、幂等键、业务主键和输入版本,调用方得到的是 accepted,不是已完成。Consumer 如果开始拉取多个数据源、维护多个阶段、等待人工、发布正式版本,就已经超出单消息任务的边界。

Batch / 分布式 Job 适合全量同步、ETL、对账、索引重建、历史修复和训练任务。它的核心不是“后台慢慢跑”,而是把“处理全部对象”拆成可观察、可重试、可恢复的分片推进。Kubernetes Job 将一次性任务、并行完成计数和工作队列模式区分开,并允许挂起后恢复,这说明调度层的并行与业务层的对象状态仍然是两个需要分别建模的问题。[11]

6.4.3 Workflow、Agent Runtime 与 Multi-Agent 的边界

Workflow 适合阶段依赖、外部等待、审批门禁、发布保护和人工转派。它把执行路径显式写出来,让运行时负责持久化等待、超时、重试、回放和状态查询。Temporal 将 Workflow Execution 描述为可持久、可恢复的执行单位,并通过事件历史与重放恢复进度;这类能力正是长任务需要的执行骨架,但业务事实仍应由领域服务拥有。[7] 中文云工作流资料也把顺序、选择、并行、状态跟踪、错误重试和长时间执行列为工作流编排的核心能力;这些能力可以由平台提供,但不应替代业务状态和发布门禁。[17][19]

单 Agent 适合高认知密度、低到中等执行跨度的问题,例如一次性分析、轻量诊断、资料归纳和小范围代码解释。它不应直接获得无限工具权限,也不能因为输出了“完成”就被视为完成了副作用。Stateful Agent Runtime 则把 Goal、Plan、Step、Tool Call、Observation、Review、Budget 和 Checkpoint 变成正式事实,让 Agent 可以暂停、恢复、返工和等待人工。

Multi-Agent 不是“多几个模型一起聊天”,而是把规划、执行、审查、协调或领域专家角色拆开。只有当一个执行单元同时承担这些职责会产生明显冲突,并且角色分工可以用指标证明收益时,才值得引入它。AutoGen 的原始论文把多 Agent 对话描述为可组合的应用基础设施,同时也意味着需要显式定义角色、消息和交互行为,而不是把自然语言对话当作隐式协议。[15]

6.4.4 选型 ADR:记录获得与牺牲

以“供应商同步是否采用 Workflow + Batch”为例,ADR 不应只写“推荐使用某某框架”,而应记录:

ADR 项目内容
背景与问题供应商分页、限流、质量异常和发布风险使单个 Consumer 无法安全恢复
候选方案串行脚本、MQ Consumer、Batch、Workflow + Batch、Agent 主导流程
决策驱动因素局部恢复、版本发布、人工门禁、外部配额、可审计性
最终决策Workflow 管理阶段,Batch 管理 Shard,Agent 只用于低风险异常归因
获得的能力可观察状态、分片恢复、质量门禁、版本切换、人工接管
主动牺牲额外状态存储、调度运维成本、最终发布延迟和模型调用成本
已接受风险部分差异需要人工处理,Workflow 和 Worker 之间仍可能重复投递
验证指标最老任务年龄、Shard 恢复时长、发布差异率、人工队列逾期量
重新评估条件数据规模、供应商协议、质量规则或人工处理量发生明显变化

ADR 的价值在于把“为什么现在选择它”与“未来什么条件下应重新选择”分开。它也迫使团队承认获得的能力与牺牲的能力必须用同一组维度比较,不能只说“可靠性更高”而不说明成本、延迟、复杂度和人工负担。[16]

6.4.5 一条渐进式升级路径

实践中可以沿着下面的路径渐进升级:

同步短任务
  -> 异步单元任务
  -> 分片并行 Batch
  -> 阶段化 Workflow
  -> 在局部步骤引入单 Agent
  -> Stateful Agent Runtime
  -> 有明确分工和终止条件的 Multi-Agent

这不是规定所有系统必须走完的路线。能用同步短任务解决,就不要先上 Workflow;能用 Batch 解决,就不要默认上 Agent Runtime;能用单 Agent 解决,就不要为了“更像 AI 系统”直接拆成 Multi-Agent。阿里云的工程实践也把单体函数、事件触发和 Workflow 编排区分为不同适用区间:流程长、需要状态持久化和自定义重试时,才有理由承担 Workflow 的编排成本。[20] 真正的升级信号应来自新出现的事实:请求被阻塞、局部失败无法恢复、阶段需要门禁、输入需要动态理解,或者一个执行单元已经不能兼顾规划、执行和审查。

6.5 完整案例:供应商数据同步如何设计

假设平台每天从供应商拉取酒店、房型、价格和库存。供应商接口按城市分页,存在限流、游标过期和字段缺失;同步结果必须经过质量校验后才能覆盖线上可售数据。这个案例同时包含大批量推进、外部不确定性、版本发布和人工治理,适合演示完整任务模型。

6.5.1 Plan:固化输入、规则和分片

调度器首先创建一条 Run,固化供应商、同步范围、输入快照时间、规则版本、发布目标、截止时间和操作者。随后根据城市、数据量和供应商配额生成 Shard。每个 Shard 记录起始游标、预计页数、优先级、最大重试次数和资源池。

Plan 阶段要验证三个不变量:分片不重叠,所有目标范围都被覆盖;规则版本和发布目标存在,不能中途从全局配置读取;输入快照和截止时间可解释,不能把不同时间采集的数据假装成同一批次。如果配置缺失、范围异常或上一次 Run 仍占用同一资源,应在真正拉取之前拒绝本次 Run。

6.5.2 Fetch:拉取与暂存分离

Worker 领取 Shard 后,以供应商配额为上限拉取分页数据。单次请求应有超时、请求号和稳定幂等键;如果游标失效,先判断供应商是否支持按对象版本或时间窗口重建游标,不能无条件从第一页开始制造重复流量。原始响应至少保留来源请求号、拉取时间、页码或游标、规则版本和快照引用,必要时保存压缩快照与校验和。

转换结果先写入暂存区,并用 run_id + shard_id + source_object_id 作为业务幂等键。暂存的目的不是多存一份数据,而是把“供应商返回了什么”“平台如何转换”“正式发布了什么”分开。重复拉取、分片重试和 Worker 重启都不能制造额外记录,也不能覆盖其他 Run 已确认的结果。

6.5.3 Normalize、Validate:对象级与批次级门禁

对象级校验检查必填字段、日期范围、币种、引用关系、数值上下限和来源版本。错误对象进入失败池,携带输入引用、错误分类、规则版本和推荐动作;对象级错误不一定阻塞全部城市。批次级校验则检查缺失比例、数据量突变、价格波动、库存变化、关键字段覆盖率和与上一有效版本的差异。

质量门禁应明确三种结果:通过,可以进入发布;可容忍,记录差异并在授权范围内继续;阻断,保留暂存版本并进入 MANUAL_REVIEW。门禁阈值不应被写成一个无法解释的常数。例如“价格下降 50%”可能对促销商品合理,对酒店库存则可能代表供应商接口异常;阈值必须带有业务范围、数据来源、版本和豁免理由。

阶段成功判据可恢复点失败策略
Plan范围覆盖且分片不重叠分片计划已持久化拒绝运行并修正配置
Fetch原始数据可追溯、游标已确认最后一页已暂存限流退避;永久错误转失败池
Normalize字段映射完成且版本记录完整对象级暂存结果隔离坏数据并记录规则版本
Validate对象和批次质量门槛通过校验报告已保存超阈值阻断发布、转人工
Publish目标版本原子切换完成发布版本号回退有效版本或冻结范围
Reconcile输入、暂存、发布和下游可见数可核对差异报告已保存补数、重放或再次核验

6.5.4 Publish:版本切换而不是逐行污染线上数据

发布是数据任务最危险的一步。推荐将暂存结果写入候选版本,例如 supplier-A:v19,通过校验后再把城市或供应商的“当前有效版本”从 v18 原子切换到 v19。如果只能按城市发布,也要为每个城市保留上一有效版本、切换时间和回滚点。不要把暂存数据逐行更新到线上主表,否则中途失败会让用户同时看到一半新数据和一半旧数据。

发布前应执行权限和资源检查:本次 Run 是否有权影响可售数据,发布范围是否与 Plan 一致,是否有其他版本正在切换,是否超过单次影响半径。高风险发布可以按城市、租户或供应商分批灰度,并在每批之间检查错误率、查询延迟和业务指标。发布后再触发索引、缓存和通知等派生任务,避免尚未通过质量门禁的数据扩散到更多系统。

6.5.5 Reconcile:用对账结束 Run

Run 结束前需要比较输入对象数、成功转换数、隔离错误数、暂存数、发布数和下游可见数。数量相等也不代表正确,还要比较版本、业务键、缺失范围和关键字段摘要。某些差异可以在预先定义的容忍范围内结束,某些差异必须把 Run 标为 MANUAL_REVIEW;只有所有必需阶段收敛,并且差异已被解释或转入明确的后续任务,Run 才能进入 SUCCEEDED。

对账发现问题时,应创建新的、可审计的修复 Run,而不是让运营人员直接修改结果表。修复 Run 需要引用原始 Run、差异报告、修复策略和审批记录。这样可以区分“原始供应商数据有问题”“转换规则有问题”“发布过程有问题”和“下游索引没有追上”,也能避免一次人工修改破坏后续重放。

6.6 Agent 任务的状态化与协作

6.6.1 Agent 增加的是不确定性,不是免除治理

Workflow 的执行路径通常由代码预先定义;Agent 则可以根据输入和中间观察结果选择下一步。ReAct 论文将推理与行动交错,让模型在行动结果的基础上更新计划,这解释了 Agent 为什么适合检索、诊断和开放式探索;但它也意味着每一次工具调用都可能改变环境,必须被当作正式执行事实记录。[13]

可以把 Agent 看成一个具有动态规划能力的执行单元:

Goal
  -> Plan
  -> Select Tool / Ask Human
  -> Execute
  -> Observe Ground Truth
  -> Update State
  -> Review / Continue / Stop

其中 Plan 不是承诺,Tool Result 才是环境事实;模型的自然语言总结不是状态机,Checkpoint 也不能用“对话上下文里大概记得什么”替代。Agent 每次需要继续工作时,都应从持久化的 Goal、Step、工具结果、审查结论和预算中恢复。Anthropic 的工程实践也强调,Agent 在每一步都需要从环境取得真实反馈,并在检查点或阻塞点暂停等待人类判断。[14]

6.6.2 Stateful Agent Runtime 的最小模型

一个可生产化的 Agent Runtime 至少需要以下实体:

实体作用不应承担的责任
Goal用户目标、范围、成功判据不直接作为执行状态
Plan当前计划和候选步骤不保证每一步一定执行
Step一个可验证的执行单元不隐藏多个不可分的副作用
Tool Call一次工具请求、输入和结果不默认代表业务成功
Observation环境返回的事实、日志或测试结果不把模型猜测当事实
Review自动或人工审查结论不绕过权限和策略
Checkpoint可恢复的状态快照不只保存对话文本
Budget时间、Token、调用次数和副作用额度不用无限预算掩盖失控

每个 Step 至少有 PENDING、RUNNING、WAITING、SUCCEEDED、RETRYABLE、BLOCKED、REJECTED 和 CANCELLED 状态。工具调用需要保存工具版本、参数摘要、请求 ID、权限上下文、返回摘要和是否产生副作用。对代码 Agent 来说,修改文件、执行命令和提交变更的风险等级不同;对运维 Agent 来说,查询监控、生成变更计划和执行生产变更也必须有不同权限。

6.6.3 代码、调研与运维 Agent 的边界

代码 Agent 的安全链路可以是:读取需求 → 生成计划 → 扫描仓库 → 修改受限文件 → 运行测试 → 生成差异摘要 → 等待人工确认 → 提交或回滚。计划阶段可以由模型主导,修改和测试需要确定性工具,提交和发布则应由权限策略与人工门禁控制。测试失败不是“模型再想一会儿”就算恢复,而是一个带输出、环境和重试次数的失败 Step。

调研 Agent 更适合并行检索、提取证据、去重、冲突检测和摘要,但每个结论都要保留来源、访问时间和证据片段的引用关系。模型可以提出“需要继续搜索”的判断,不能把没有来源的段落自动升级为事实。对于高影响结论,可用第二个 Reviewer 检查来源覆盖和推理跳跃,并把返工次数纳入预算。

运维 Agent 必须把查询与变更分离。查询日志、指标和配置属于低风险观察动作;生成迁移计划、扩容或重启属于中风险动作;删除数据、切换流量或修改生产配置属于高风险动作。高风险动作应该要求明确授权、显示影响范围、支持 dry-run,并在执行前保存回滚点。Agent 的价值在于缩短诊断和准备时间,不等于应该把最终控制权交给模型。

6.6.4 Tool Contract、权限与成本预算

工具契约应包含名称、版本、输入 Schema、输出 Schema、超时、可重试错误、不可重试错误、幂等键、权限范围、副作用等级和审计字段。工具返回值最好分为 success、retryable_error、permanent_error、unknown,而不是只返回一段自然语言。这样 Runtime 才能决定重试、返工、等待人工还是结束。

预算至少分为四类:最大执行时间、最大模型调用次数、最大工具调用次数和最大副作用额度。预算耗尽后,任务进入 BLOCKED 或 MANUAL_REVIEW,不能通过自动循环悄悄重置。上下文也应分层保存:用户目标和权限是稳定上下文,步骤状态和工具结果是执行事实,临时推理是可丢弃上下文。把全部对话拼成一个越来越长的 Prompt,会使恢复、成本和隐私边界都变得不可控。

6.6.5 Multi-Agent 的角色、共享状态与终止条件

复杂调研或研发任务可以采用最小三角结构:Planner 负责拆解,Worker 负责执行,Reviewer 负责检查。只有在角色责任、输入输出和失败路径清晰后,才考虑增加 Router、Supervisor 或领域 Specialist。角色拆分应减少职责冲突,而不是把同一问题让多个 Agent 重复回答。

共享状态不应只是一个公共聊天房间。更稳妥的做法是维护任务账本,每个 Agent 以结构化消息写入:完成了什么、使用了哪些来源或工具、产出了什么、仍缺什么证据、推荐下一步是什么。Planner 可以读取账本生成下一轮计划,Reviewer 可以针对某个 Step 返回 ACCEPT、REWORK 或 REJECT。消息必须有版本和幂等键,避免重投后重复创建任务。

终止条件至少包括:目标判据满足、所有必需 Step 已收敛、预算耗尽、连续返工超过上限、遇到无法自动判断的冲突、人工拒绝或截止时间到达。没有终止条件的 Multi-Agent 很容易进入循环讨论;没有最终裁决者的协作系统也无法解释为什么把互相矛盾的结果交付给用户。Agent 之间的分工获得了并行和专业化能力,但牺牲了通信成本、可复现性、延迟和调试简单性,这些都应进入 ADR 和运行指标。

6.7 故障与治理:默认假设会重复、会超时、会中断

6.7.1 幂等、重复和结果未知

任何 Shard、Tool Call 或消息都可能被重复执行:消息至少一次投递、租约超时、进程在提交确认前重启,都可能造成重复。数据库写操作可以用唯一约束、条件更新或 UPSERT;外部接口应携带稳定请求号并支持按请求号查询。幂等键必须表达业务意图,而不是简单使用随机消息 ID。例如 run_id + shard_id + source_object_id 能识别同一次同步对象处理,order_id + operation_type 能识别同一次业务动作。

最危险的是“结果未知”。一次调用超时只说明调用方没有在截止时间内收到响应,并不说明下游没有执行。正确顺序是:把步骤写成 UNKNOWN,记录外部请求号;优先查询同一幂等键;确认已完成则补写执行事实;确认未执行才允许使用同一键重试;仍无法确认则进入延迟查询、对账或人工队列。直接生成一个新请求号重发,可能造成重复扣费、重复发布或两份资源占用。

6.7.2 重试、超时与失败池

重试必须分类。网络抖动、短暂限流和临时不可用通常可以指数退避并加入随机抖动;字段非法、权限拒绝、规则缺失和质量门禁失败属于永久或业务错误,应该隔离而不是重试;超过重试预算的任务应停止自动尝试。Google SRE 建议限制每次请求和客户端的重试预算,并避免多层同时重试,因为多层重试会产生组合式流量放大。[2][3]

超时至少分为单次调用、Step / Shard 执行和整次 Run 三层。单次超时保护连接与线程;Shard 超时防止死循环或异常大数据;Run 截止时间避免旧任务和下一轮长期重叠。超时触发后,先判断是否存在未知副作用,再决定取消、重试、查询或对账。失败池或 DLQ 不应是垃圾桶,每条记录至少保存原始输入引用、错误分类、尝试次数、最近错误、首次失败时间、下次动作和权限范围;运维人员可以修复输入后重放,但所有动作都要留下审计。

6.7.3 并发、限流与背压

任务系统常见事故不是完全不跑,而是跑得太快。并发度应按被保护资源设置:供应商配额、租户范围、数据库连接池、计算集群、模型服务和发布能力都可能需要独立限制。实时库存修复不应和月度索引重建共享没有隔离的资源池;紧急补数可以有受控加急通道,但不能绕过配额、权限和审计。

背压意味着下游处理不过来时,上游主动减速。需要观察队列积压、Shard 等待时间、接口 429 比例、数据库连接池、Worker 利用率、发布延迟、重试率和最老任务年龄。超过阈值后,可以降低拉取并发、延长退避时间、暂停新 Run、只保留高优先级任务或返回“稍后查询”。AWS 和 Google 的可靠性资料都把超时、退避、抖动、限流和负载削减视为同一个过载控制问题,而不是互不相关的参数调节。[2][3]

6.7.4 暂停、取消、恢复与人工接管

暂停表示保留 Run、Shard 和 Checkpoint,稍后继续;取消表示不再领取新工作,并让已运行的 Shard 在安全边界退出。Worker 应在页、批或对象边界检查控制标记,在完成当前原子单元并写入结果后退出。强行杀进程可以是故障处置手段,但不应被当作正常取消协议,因为它会把状态留在未知位置。

恢复程序应先回收过期租约,再扫描 RETRYABLE、长期 RUNNING 和 UNKNOWN Shard。对未知外部副作用先查询或对账,确认安全后再投递;对长期无进度的执行器,要结合心跳、资源指标和最后 Checkpoint 判断是卡死、慢任务还是外部依赖阻塞。人工接管不能通过直接改库绕过状态机,操作台应展示进度、最近错误、外部请求号、影响范围和建议动作,并要求操作者填写理由、证据和后续验证。

6.7.5 顺序、版本与数据一致性

任务并行化后,并不是所有对象都需要全局顺序,但同一酒店、商品或账务账户的连续更新常常需要按版本收敛。不能依赖“队列看起来有序”:扩容、重放、延迟投递和人工补发都会破坏全局接收顺序。更可靠的做法是让对象携带来源版本、更新时间或单调序号,写入时只接受不旧于当前版本的数据。

示例中的条件更新只表达版本保护思想:

UPDATE staged_inventory
SET quantity = :quantity,
    source_version = :source_version,
    updated_at = :updated_at
WHERE supplier_id = :supplier_id
  AND source_object_id = :source_object_id
  AND source_version <= :source_version;

实际比较依据必须符合来源系统语义。如果上游版本不可信,不能伪造一个“看似有序”的时间戳;应保留原始快照、运行批次和冲突记录,把无法自动判定的新旧关系转入人工核验。对 Agent 任务也一样:旧的计划或工具结果不能因为晚到就覆盖新的人工拒绝或策略版本。

6.7.6 观测、对账与发布保护

可观测性要让值班人员在几分钟内回答:Run 是否在推进,停在哪个 Step,哪些 Shard 受影响,下游是否被压垮,是否需要暂停或介入。日志、指标和追踪应携带 task_id、run_id、step_id、shard_id、幂等键、输入版本和外部请求号。OpenTelemetry 提供跨服务关联 Trace、Metric 和 Log 的统一语义,但业务字段仍需要由任务系统主动写入;没有这些字段,技术链路可见也不等于任务事实可见。[12]

层级必备事实主要告警
Run开始、截止、输入版本、总进度、最终摘要长期未结束、整体失败率超阈值
Step耗时、输入输出数、异常率、等待原因阶段耗时突增、门禁失败
Shard租约、重试数、处理速率、Checkpoint、最近错误租约长期占用、重复失败、无进度
下游队列积压、限流、连接池、发布延迟背压触发、配额耗尽、发布阻塞
Agent工具调用、Token、预算、返工次数、人工等待循环调用、权限拒绝、成本突增

对账负责发现遗漏、重复和错位,在线任务负责快速推进,二者不能互相替代。供应商同步至少核对输入、暂存、发布和下游可见版本;索引任务核对主库版本与索引版本;Agent 调研核对结论与来源引用;模型训练核对数据集、代码、参数、评估结果和产物签名。发现差异后,创建新的修复 Run 或人工任务,不要直接删除错误历史。

发布保护应有最小质量门槛、影响半径和停止条件。质量阈值突破、下游错误率持续上升、回滚次数超过上限或人工确认未完成时,系统应暂停后续发布并报警。版本化发布的回滚对象是“当前有效版本或发布开关”,不是试图删除所有已经传播的消息;如果外部副作用已经发生,回滚应被建模为新的补偿任务,并记录影响范围、进度和验证结果。

6.7.7 任务系统的测试与演练

长任务不能只测试主路径。至少要演练:Worker 在外部调用成功后崩溃、消息重复投递、租约过期、Checkpoint 写入失败、游标失效、下游限流、数据库主从切换、旧版本晚到、发布门禁阻断、人工审批超时、Agent 工具返回歧义结果,以及恢复程序重复运行。

测试结果应验证不变量,而不是只验证最终 HTTP 状态。例如同一幂等键不能产生两次副作用;旧版本不能覆盖新版本;暂停后不能新增未授权工作;取消后已运行的原子单元可以安全结束;人工拒绝后 Agent 不能自动重新提交同一高风险动作;对账差异必须能指向具体 Run、Step、Shard 和责任人。对于不可确定的模型行为,应使用固定任务集、工具桩、轨迹评估和人工抽检建立回归集,不要把一次“看起来合理”的回答当成可靠性证明。[14]

6.8 方法论总结:按问题升级,而不是按潮流堆组件

任务处理可以压缩成一条判断链:先判断是否跨执行边界,再判断是否需要大量对象推进,再判断是否需要阶段治理,最后判断是否需要动态推理与角色协作。同步调用、异步消息、Batch、Workflow、Agent Runtime 和 Multi-Agent 不是互相替代的技术潮流,而是不同复杂度区间中的执行形态。

本章的核心判断如下:

  1. 先定义输入、输出、成功判据、截止时间和资源边界,再讨论框架。
  2. 只要中断后必须恢复,就需要持久化 Run、分片和可信 Checkpoint。
  3. 只要存在副作用,就把重复、延迟、乱序和结果未知当作默认路径。
  4. 任务级成功不等于对象级成功;批量任务必须能解释局部失败和跳过。
  5. 队列负责分发,状态存储负责事实,日志负责解释;三者不能互相冒充。
  6. Batch 解决大规模推进,Workflow 解决阶段依赖与治理,二者经常组合。
  7. Agent 可以处理不确定输入,但不能替代状态机、权限、审计、对账和人工门禁。
  8. Multi-Agent 只有在角色分工能带来可验证收益时才成立,并且必须有共享状态和终止条件。
  9. 重试、限流、背压、失败池、对账和演练共同构成可靠性,而不是事后补丁。

在进入实现前,可以用下面的评审清单快速检查一个方案:

评审问题若回答不清,通常意味着
这次任务的业务输入、输出和权威事实是什么?任务可能把日志、队列或缓存误当成事实
失败后从哪个最小安全单元继续?任务可能只能全量重跑
同一动作重复两次会怎样?幂等、版本或结果查询能力不足
哪些错误可以重试,哪些必须隔离或人工处理?重试策略会放大故障或吞掉数据错误
谁能暂停、取消、批准发布和人工放行?控制面与权限边界不完整
如何证明已经完成,如何发现遗漏?缺少对象级结果和对账闭环
为什么需要 Agent,为什么还不需要 Multi-Agent?方案可能由技术潮流而非复杂度驱动

成熟的任务系统不是一组“后台脚本”,也不是让模型自由循环的对话。它是一套能够跨故障、跨时间、跨人员交接持续收敛的执行机制:有明确的事实、有安全的恢复点、有可见的失败、有受控的副作用,也有在自动化无法安全判断时及时停下来的能力。

6.9 参考资料

正文中的编号用于就近说明来源支撑的事实或方法;其中“本章推导”表示本书结合供应商同步、数据处理和 Agent 协作场景做出的工程归纳,不等同于来源原文的直接结论。链接优先指向作者、出版方或长期维护的官方页面,访问与加入日期为 2026-09-21。

[1] Kleppmann M. Designing Data-Intensive Applications[M]. O’Reilly Media, 2017. 用于数据系统的可靠性、可扩展性、可维护性和数据流处理权衡。

[2] Beyer B, Jones C, Petoff J, Murphy N R, eds. Site Reliability Engineering[M]. O’Reilly Media, 2016. 用于服务目标、过载、重试预算和生产治理。

[3] Google SRE. Addressing Cascading Failures[EB/OL]. Google, 2016. 用于解释重试、资源耗尽和级联故障的正反馈关系。

[4] Richardson C. Saga Pattern[EB/OL]. Microservices Patterns. 用于区分跨服务流程协调、补偿和全局原子提交。

[5] Richardson C. Transactional Outbox[EB/OL]. Microservices Patterns. 用于说明本地事实与事件可靠交接的双写边界。

[6] Richardson C. Idempotent Consumer[EB/OL]. Microservices Patterns. 用于说明至少一次投递下的重复消费处理。

[7] Temporal. Workflow Execution overview[EB/OL]. Temporal Documentation, 2026. 用于说明持久执行、事件历史、重放、暂停和恢复。

[8] Apache Airflow. Core Concepts[EB/OL]. Apache Software Foundation, 2026. 用于说明 DAG、Task、依赖、重试和任务执行架构。

[9] Apache Spark. RDD Programming Guide[EB/OL]. Apache Software Foundation, 2026. 用于说明 RDD、分区、血缘和节点故障恢复。

[10] Apache Flink. 有状态流处理[EB/OL]. Apache Software Foundation, 2026. 用于说明状态、Checkpoint、Savepoint 和有状态流处理容错。

[11] Kubernetes. Job 中文文档[EB/OL]. Cloud Native Computing Foundation, 2026. 中文官方资料,用于说明一次性任务、并行 Job、工作队列和挂起恢复。

[12] OpenTelemetry. Documentation[EB/OL]. OpenTelemetry, 2026. 用于说明 Trace、Metric、Log 的关联观测基础。

[13] Yao S, Zhao J, Yu D, et al. ReAct: Synergizing Reasoning and Acting in Language Models[C/OL]. ICLR, 2023. 用于说明推理与工具行动交错的 Agent 执行模式。

[14] Anthropic. Building effective agents[EB/OL]. Anthropic, 2024. 英文工程资料,用于区分 Workflow 与 Agent、工具契约、环境反馈和人类检查点。

[15] Wu Q, Bansal G, Zhang J, et al. AutoGen: Enabling Next-Gen LLM Applications via Multi-Agent Conversation[C/OL]. arXiv, 2023. 用于说明 Multi-Agent 的角色化对话、工具和人工输入组合。

[16] 周志明. 《凤凰架构:构建可靠的大型分布式系统》[EB/OL]. 2026. 中文开源技术资料,用于补充分布式系统的可靠性、状态和架构演进视角。

[17] 阿里云. 协调多个分布式任务执行——Serverless 工作流[EB/OL]. 阿里云帮助中心. 中文官方资料,用于说明顺序、选择、并行、状态跟踪、重试和长时间流程编排。

[18] Apache DolphinScheduler. Apache DolphinScheduler[EB/OL]. Apache Software Foundation, 2026. 官方项目资料,用于补充 DAG 工作流、暂停恢复、工作组隔离和批量回填等调度能力。

[19] 阿里云. 云工作流:功能特性[EB/OL]. 阿里云帮助中心, 2024. 中文官方资料,用于说明状态节点、选择、并行、循环和标准模式的长流程持久化。

[20] 阿里云开发者社区. Serverless 工作流适用场景及最佳实践[EB/OL]. 阿里云开发者社区, 2020. 中文工程实践资料,用于补充单体函数、事件触发和 Workflow 编排之间的适用边界;具体产品行为以官方文档为准。