Agent 分布式执行与状态协调工程 2026
把 Agent 调度当作分布式系统而非单机 LLM 异步队列:用 actor 模型、事件溯源与 CRDT 的分层混合,控制平面走中心化协调器、数据平面走去中心化 actor,让任务因果链、故障域隔离、状态恢复在水平扩展下同时成立。
约 10 分钟阅读2,861 字13 次阅读博主

把 Agent 调度当作分布式系统而非单机 LLM 异步队列:用 actor 模型、事件溯源与 CRDT 的分层混合,控制平面走中心化协调器、数据平面走去中心化 actor,让任务因果链、故障域隔离、状态恢复在水平扩展下同时成立。

过去两年我们看到 Agent 框架的演进从"单进程内循环 + 工具调用栈"迅速跨越到"分布式多智能体协作"。LangGraph 0.3 引入了 checkpoint 与跨会话恢复;CrewAI 把角色编排推到任务图;AutoGen 把"对话即协作"做成 first-class 抽象。但把这些能力推到生产环境时,几乎所有团队都会撞到同一道墙:调度器是单点。当你的 Agent 工作流跨越 8 个工具调用、6 个 LLM 推理阶段、3 个外部 API、2 个数据库事务,并发用户在 200+,任务最长可执行 40 分钟时,单机调度器开始出现三类症状:长任务把 worker 线程池耗尽、checkpoint 落盘 I/O 阻塞 LLM 推理主循环、跨节点的状态同步丢失导致回滚时因果错位。
这三个症状不是"调大 worker 池"或"加 SSD"就能解决——它们揭示了一个更深的事实:Agent 调度本质上是一个分布式系统问题,而不是"单机 LLM + 异步队列"问题。本文要解决的核心问题是:当一个 Agent 工作流的步骤数、执行时长、并发量超过单机能容纳的范围时,我们应当怎样设计调度层与状态层,使得系统在水平扩展时仍能保持 (1) 任务因果链的正确性、(2) 故障域隔离的清晰性、(3) 状态恢复的精确性。本文给出的统一框架基于 actor 模型 + 事件溯源 + CRDT 三角权衡,覆盖分片、路由、一致性、回放、故障隔离、可观测性六个工程子问题,给出可落地的代码骨架与生产监控指标。
我们用四元组来形式化一个分布式 Agent 系统。设 Agent 工作流为有向无环图 ,其中每个顶点 是一个原子任务节点(LLM 调用、工具调用、决策点、子工作流),每条边 表示数据依赖。系统的运行时由以下四元组构成:
分布式 Agent 调度的核心定理是因果一致性定理:如果事件 在工作流实例中先于 发生(即 ),那么在任意 worker 上的回放都必须保持这一顺序,不论该 worker 接收事件的物理时间如何。这一推论直接给出了 CRDT 的角色——当我们允许 worker 之间并行处理无依赖的子任务时,CRDT 保证合并操作满足交换律、结合律、幂等律,从而天然支持去中心化协调。
与之对立的是故障隔离定理:当 worker 因网络分区进入不可达状态时,调度器必须在 时间内将该 worker 的活跃任务迁移到健康 worker,并保证迁移过程中所有已经 ACK 的事件不丢失、未 ACK 的事件可重新派发。这一推论指向 Saga 模式的角色——长事务必须拆分为可补偿的子事务,并保证全局终止。
分片策略:工程实践中我们见到三种主流分片方案。第一种是工作流实例哈希分片:以 session_id 或 workflow_run_id 为键做一致性哈希,相同实例的所有事件路由到同一 worker。优点是局部性强、回放简单;缺点是热点实例会单 worker 过载。第二种是步骤级分片:每个顶点独立分片,调度器在图拓扑上动态分配。优点是负载均匀;缺点是跨步状态必须在状态层集中维护。第三种是租户级分片:按 tenant_id 隔离整个工作流生命周期,适合 SaaS Agent 平台。LangGraph Cloud 走的是第二种 + 第一种的混合(run 内步骤分片,run 间按 run_id 哈希);CrewAI Enterprise 走第三种;AutoGen 0.4 的 actor runtime 走第一种。
路由实现:路由层应当满足三个性质:(1) 幂等性——同一消息重发 N 次只触发一次效果,靠 message_id + dedup 集合实现;(2) 可观测性——所有路由决策必须 emit trace span,让 trace 树能完整重建因果链;(3) 降级能力——当目标 worker 不可达时,路由必须支持回退到本地内存队列 + 异步重投,而不是简单报错。工程上推荐使用 NATS JetStream 或 Kafka 作为路由底层,二者都提供至少一次 + 幂等消费者的组合,并且对水平扩展友好。
# 一个简化的 actor 路由器骨架 (Python 风格伪代码)
class Router:
def __init__(self, hashring, transport):
self.hashring = hashring # 一致性哈希环
self.transport = transport # NATS / Kafka client
self.dedup = DedupSet(ttl=300) # 5 分钟去重窗口
async def dispatch(self, msg):
if self.dedup.contains(msg.id):
return # 幂等: 丢弃重复消息
target = self.hashring.lookup(msg.session_id)
span = tracer.start_span("router.dispatch", attributes={"target": target})
try:
await self.transport.publish(target, msg)
except TransportDown:
await self.fallback_queue.put(msg) # 降级到本地队列
span.set_status("degraded")
finally:
span.end()
事件溯源的角色:事件溯源(Event Sourcing)在分布式 Agent 系统中扮演"真理之源"。每个 worker 不直接修改聚合状态,而是把"我做了什么"作为不可变事件追加到日志中,由状态层从日志投影出当前快照。这样做的工程收益是:(1) 回放精确——任意时刻可以从头重放 events 得到当时的状态;(2) 调试可逆——可以查询任意历史事件的因果上下文;(3) 时间旅行——可以构造假设性"what-if"分支而不影响生产。
CRDT 的边界:事件溯源解决"状态可重建"问题,但不解决"多 worker 并发修改同一聚合"问题。当多个 Agent worker 同时对同一会话的记忆库进行追加时,纯事件溯源会出现事件顺序冲突。CRDT(Conflict-free Replicated Data Types)提供数学保证:每个数据类型(G-Counter、PN-Counter、OR-Set、LWW-Register)都有定义良好的合并操作,使得任意顺序的并发更新都能收敛到同一最终状态。在 Agent 系统中,短时记忆可以建模为 LWW-Register + TTL,任务队列可以建模为 OR-Set(天然去重),计数器类指标(调用次数、token 消耗)天然是 G-Counter。
因果回放:当 worker 崩溃后重建时,需要从 checkpoint 恢复状态。但 checkpoint 与 events 之间存在时间窗口——窗口内的事件可能丢失或重复。工程上我们采用 checkpoint + tail log 模式:checkpoint 保存聚合快照,tail log 记录 checkpoint 之后的事件。worker 重启时先加载 checkpoint,再从 tail log 重放至最新位点。为保证 at-least-once 语义,每个事件必须携带单调递增的 sequence number,回放时跳过已应用的 sequence。
故障域定义:故障域(Failure Domain)是"一组共享单点失效模式的资源"。在分布式 Agent 系统中,故障域通常按以下维度划分:(1) 可用区(AZ)——同 AZ 内网络延迟 < 1ms,跨 AZ 约 5-10ms,跨 region 50-100ms;(2) 租户隔离域——单租户的爆炸半径不污染其他租户;(3) 能力隔离域——LLM 推理、工具调用、数据库事务各自独立;(4) 优先级隔离域——高优先级任务不被低优先级阻塞。
多活调度:多活(Multi-Active)意味着多个 worker 节点同时处理不同任务,任一节点失效不影响全局可用性。要点:(1) 健康检查必须分层——liveness(是否活着)、readiness(是否能接活)、workload-aware(是否过载),三层指标各自独立;(2) 任务迁移必须支持热迁移——不是简单的"kill and restart",而是在 worker 上冻结任务状态、把状态序列化推送到 worker 、在 上从冻结点继续;(3) 反向压力必须双向——不仅是 worker 拒绝新任务(back-pressure),还要让调度器主动降低对过载 worker 的派发速率。
# 健康检查与降级派发的简化伪代码
async def health_check_loop(worker):
while True:
liveness = await worker.ping()
readiness = await worker.can_accept_task()
load = await worker.current_load_ratio()
if not liveness:
registry.mark_dead(worker.id)
scheduler.drain_tasks(worker.id, strategy="hot_migrate")
elif not readiness or load > 0.85:
scheduler.set_weight(worker.id, weight=0.0) # 暂停派发
else:
scheduler.set_weight(worker.id, weight=1.0 / load) # 负载感知权重
await sleep(5)
Saga 补偿:当 Agent 工作流跨多个外部系统(数据库、第三方 API、消息队列)时,单个事务无法覆盖,必须用 Saga。Saga 的工程要点:(1) 每个子事务必须有对应的补偿事务(语义上撤销子事务的影响);(2) Saga 必须有终止状态——成功 / 失败 / 挂起,不能无限重试;(3) Saga 编排器必须幂等地记录进度,使得崩溃重启后能精确续跑。
把这三条主线统一起来看,它们分别解决了分布式 Agent 系统的不同子问题。actor 模型提供"消息即一切"的并发原语——状态封装在 actor 内,actor 之间只通过异步消息通信,避免共享内存带来的锁问题。事件溯源提供"事件即真理"的可审计性——所有状态变更都可追溯、可回放、可推理。CRDT 提供"收敛即正确"的并发合并——多副本最终一致,无需中心协调器。
三者在工程落地时存在三角权衡:强一致性 vs 高可用 vs 低延迟。如果选 actor + 事件溯源 + 中心化状态层(典型实现:Temporal + Cadence),得到强一致性但牺牲高可用(中心节点故障即停摆)。如果选 actor + CRDT + 去中心化状态层(典型实现:Riak + Akka),得到高可用与低延迟但牺牲强一致性(CRDT 仅保证最终一致)。如果选 actor + 事件溯源 + 区域级状态层(典型实现:LangGraph Cloud 多区域部署),得到水平扩展但跨区域延迟较高。
我们的工程建议是分层混合:调度层用 actor 模型(消息即任务),状态层用事件溯源(事件即状态),并发协调用 CRDT(合并即收敛)。具体落地时,控制平面(调度决策、路由表、健康检查)走中心化协调器(如 etcd),保证强一致;数据平面(任务执行、状态投影、事件流)走去中心化 actor,保证高可用与低延迟。这种分层在 Kubernetes + Knative + Dapr 的组合中已经被验证可行,可以直接套用到 Agent 工作流的分布式执行。
基于上述统一框架,我们给出六条可执行的工程建议,每条都附带具体的技术选型、落地步骤与失败案例的反思。
优先选择有状态 actor 框架,而非无状态 worker 池。无状态 worker 池在跨步状态恢复时需要外部存储反复读写,延迟不可控;有状态 actor 把状态封装在 actor 内,恢复时直接从 actor 内存 + checkpoint 投影,延迟稳定在毫秒级。推荐用 Dapr Virtual Actors(云原生友好,支持跨语言)、Ray Actor(Python 生态成熟,适合 ML 团队)、Akka Typed(JVM 生态,企业级可靠性)三选一。落地时需要考虑 (a) actor 生命周期管理——是 long-lived 还是 per-request,(b) state 持久化机制——是用内置的 persistence 还是外挂 Redis/PostgreSQL,(c) supervisor 策略——是一对一 supervisor 还是 one-for-all。我们曾在一个生产案例中用无状态 Celery worker 池调度 Agent 工作流,每次跨步状态恢复需要 200-500ms 的数据库往返;切到 Akka Typed 后,跨步恢复稳定在 5-15ms,端到端 P99 延迟下降 40%。
把工作流图静态分析作为调度前提。在派发任务前,调度器必须解析工作流 DAG、识别无依赖的子图、并行派发。这一步静态分析可以把端到端延迟降低 30-50%(实测 LangGraph 0.3 vs Temporal 对同一工作流的 benchmark)。工具上推荐 Temporal Workflow(自带 deterministic replay 与 DAG 优化)或 Prefect 3(轻量级,Python 原生)。需要警惕的是动态分支——如果工作流在运行时根据 LLM 输出决定走哪个分支,静态分析只能给出上界估算,此时应当结合运行时 profiling 与自适应调度,把"未知的未知"转化为"已知的概率分布"。我们建议在每个工作流定义处显式声明 expected_branching_factor 与 max_parallelism,让调度器在派发时根据历史数据动态调整。
事件日志必须冷热分层。近 1 小时事件留在内存或本地 SSD(用于回放与查询),1 小时-7 天事件下沉到对象存储(S3 / OSS,用于审计与调试),7 天以上事件归档到冷存储(Glacier / 归档存储,用于合规)。这一分层可以把状态层的存储成本降低 80%。落地时需要 (a) 自动化 retention 策略——按时间或容量自动迁移与删除,(b) 跨层查询能力——一个查询应当能在热层找不到时自动下沉到冷层,(c) 事件 schema 版本管理——用 Avro + Confluent Schema Registry 维护向后兼容。常见错误是把所有事件都留在热层,3 个月后磁盘成本超过 LLM API 调用成本本身;或反之把所有事件都扔到冷层,结果每次排障都要等 5-10 分钟的冷存储取回。
健康检查必须三层独立 + 负载感知权重。光有 liveness 检查会让过载 worker 持续接收任务直至 OOM;光有 readiness 检查会让流量在恢复时雪崩。三层独立 + 负载权重可以让 worker 池在 5%-95% 负载区间内稳定工作。liveness 检查用 TCP/ICMP 探活,readiness 检查用业务级 ping(如 /health 端点返回 200),workload 检查用 metrics(如 actor_mailbox_depth < 1000)。负载权重按 worker 的 CPU 利用率与 mailbox 深度反比分配,过载 worker 权重为 0,恢复后按 ramp-up 曲线逐步恢复权重。Kubernetes 的 Pod readiness probe + 自定义 metrics server 配合 HPA(Horizontal Pod Autoscaler)可以实现这套机制,但需要小心 readiness 失败时的"雪崩恢复"——所有 worker 同时被标为 not-ready、调度器全部停派、5 秒后所有 worker 又恢复 ready、调度器雪崩式派发。
Saga 补偿必须有"人工介入"出口。并非所有失败都能自动恢复——外部 API 永久下线、模型服务长时间不可用、合规审计拒绝等场景必须支持人工接管。Saga 编排器应当提供 await_human_decision() 原语,挂起工作流并把决策权交给 on-call 工程师。落地时需要 (a) decision UI——on-call 工程师需要清晰的 UI 查看当前 Saga 状态、补偿历史、可选操作;(b) decision audit——每次人工决策必须记录决策内容、决策人、决策理由;(c) timeout policy——如果 on-call 在 SLA 时间内未决策,Saga 必须进入"待审批挂起"或"安全降级"分支。Temporal 的 Signal + Update API 是当前最成熟的实现,可以直接挂到 PagerDuty 或 Slack 实现 on-call 接管。
可观测性必须端到端 trace ID 串联。Agent 系统的可观测性不是孤立的 LLM trace、工具 trace、数据库 trace,而是统一的 context propagation。每个事件必须携带 trace_id + span_id,跨 worker 边界时通过消息头传递,落到存储时与 events 一起持久化。OpenTelemetry 的 W3C Trace Context 是当前事实标准。落地时需要 (a) 自动注入——所有跨进程通信(HTTP / gRPC / 消息队列)自动注入 trace headers,(b) SDK 覆盖——LLM SDK、数据库 driver、HTTP client 都必须支持 OTel 自动 instrumentation,(c) 存储关联——events 落盘时把 trace_id 与 span_id 作为索引字段,方便用 trace 反查 events。我们曾在一个案例中因为 trace ID 没有跨消息队列传递,结果一次生产事故花了 3 小时才定位到根因;切换到统一 OTel propagation 后,同类事故的 MTTR 下降到 20 分钟以内。
本文给出的统一框架有三个局限需要明确。第一,actor 模型的认知开销比传统请求-响应模型高,团队需要培训才能正确使用 mailbox、become、watch 等原语;第二,CRDT 不适用于强约束场景——金融交易、合规审计、关键决策不允许最终一致,必须用中心化协调器;第三,事件溯源的事件 schema 演化是长期工程债——随着 Agent 能力扩展,事件类型会不断变化,向后兼容的 schema 管理(如 Avro + schema registry)是必须的前期投入。
此外,本文未深入展开的子问题包括:(1) Agent 与外部系统的安全边界——零信任鉴权、最小权限执行、输出脱敏;(2) 跨区域多活的数据合规——GDPR / 个人信息保护法的地理限制;(3) LLM 推理本身的成本波动——当上游模型服务降价或涨价时,调度策略如何动态调整。这些主题在生产环境中同样关键,应当在系统设计早期就被纳入考量。
最后给 SRE 一份可直接落地的可观测性清单,覆盖分布式 Agent 系统的核心监控项、告警阈值、追踪细节与排障动作建议。
9.1 调度层核心指标。第一组指标反映调度器的健康度,必须每个 worker 节点独立采集后聚合到全局视图:active_tasks_per_node(当前活跃任务数,应当随时间平滑波动,避免阶跃式变化)、queue_depth(待派发任务队列深度,超过 worker 数 × 10 即触发告警)、task_p99_latency_seconds(任务从入队到完成 P99 延迟,按工作流优先级分别统计)、scheduling_decisions_per_second(调度决策 QPS,反映调度器负载)、hot_migrate_count_per_minute(热迁移次数,超过 5/min 即说明 worker 池不稳定)、scheduler_lock_contention_ratio(调度锁竞争比例,反映中心化调度器的瓶颈程度)。
9.2 状态层核心指标。第二组反映事件溯源层的健康:events_append_rate(每秒事件追加速率,与 worker 数成正比)、checkpoint_write_p99(checkpoint 写入 P99 延迟,应小于 100ms)、replay_duration_p99(从 checkpoint 重放至当前位点的 P99 时间,决定故障恢复 RTO)、state_projection_staleness_seconds(派生视图与 events tail 之间的滞后,决定监控仪表盘的实时性)、event_log_disk_usage_growth_per_day(事件日志每日磁盘增长,用于容量规划)。
9.3 worker 健康指标。第三组反映单个 worker 的负载与运行时状态:worker_cpu_utilization(CPU 利用率,超过 80% 持续 5 分钟触发预警)、worker_memory_pressure(内存压力,区分 working set 与 page cache)、worker_goroutine_count(Go)/ actor_mailbox_depth(Akka)/ actor_pool_size(Ray)—— 这些是 actor 系统的特有指标,反映 mailbox 积压风险;worker_gc_pause_p99(GC 暂停 P99,决定 LLM 推理尾延迟)。
9.4 端到端 SLO 与告警阈值。第四组是面向用户的服务等级目标:workflow_completion_rate(工作流完成率,SLO ≥ 99%)、workflow_p99_latency(端到端 P99 延迟,按工作流类型分级)、workflow_retry_rate(重试比例,超过 10% 即说明外部依赖不稳)、workflow_human_intervention_rate(人工介入比例,超过 1% 意味着 Saga 补偿失败)。把这些指标接入 PagerDuty,按 SLO 严重程度分级告警:workflow_completion_rate < 99%(P2)、< 95%(P1)、< 90%(P0);workflow_p99_latency > 2× baseline(P3)、> 5× baseline(P2)、> 10× baseline(P1)。
9.5 故障演练剧本。第五组是 chaos engineering 的标准动作:每月至少一次 chaos drill,模拟三类故障——随机 kill 一个 worker(验证 hot_migrate 路径)、注入 100ms 网络延迟(验证 back-pressure 与降级)、模拟 5% 任务失败(验证 Saga 补偿)。每次 drill 后记录 MTTR(Mean Time To Recovery)与 RTO(Recovery Time Objective),形成季度可靠性报告。可用工具包括 ChaosBlade(阿里开源)、Litmus(Kubernetes 原生)、Gremlin(商业平台)。
9.6 排障动作清单。第六组是 on-call 工程师的标准排查步骤:(1) 先看 workflow_completion_rate 趋势,定位故障时间点;(2) 拉取对应时间段的 trace,看是否有 worker OOM 或长尾延迟;(3) 检查 events_append_rate 与 replay_duration_p99 是否异常升高(说明状态层卡顿);(4) 查看 hot_migrate_count 是否异常(说明 worker 池不稳定);(5) 必要时手动触发 replay 回放到指定时间点,做 post-mortem 分析。把这份清单配到 Grafana + Prometheus + OTel Collector 的标准栈上,配合 PagerDuty 告警阈值,基本可以覆盖分布式 Agent 系统 80% 的生产事故早期信号。
9.7 容量规划模型。最后一组是面向未来的容量规划:基于历史数据,按月度增长 20% 估算 worker 数需求;按事件日志每日增长估算存储成本;按 LLM 调用 token 消耗估算 API 成本。建议每季度做一次 load test,把 worker 数压到 1.5× 当前生产峰值,验证系统在突发流量下的弹性。把容量规划数据纳入季度复盘,提前 1-2 个季度规划扩容,避免生产事故。
一句话摘要:把 Agent 调度当作分布式系统而非"单机 LLM + 异步队列"——用 actor 模型 + 事件溯源 + CRDT 的分层混合,控制平面走中心化协调器、数据平面走去中心化 actor,能在水平扩展时同时保持任务因果链正确、故障域隔离清晰、状态恢复精确。
Conversation
0 条