☰
构建可靠AI Agent:以分布式状态机思维设计长期运行任务
2026/9/26 7:38:53 网站建设 项目流程

1. 内容整体设计与思路拆解

1.1 为什么“Agent”和“长期运行的分布式状态机”会扯上关系

先说个我踩过的坑。前阵子接了一个客服类 Agent 项目,需求很朴素:用户半夜问一个问题,Agent 按流程查库存、查订单、再算运费,最后给个答复。刚开始我们做得也“很朴素”,直接拿 LangChain 串一个循环,LLM 调用一次工具、拿到结果、再调一次,完事。上线头三天看着挺正常,直到有用户反馈:同一笔订单,客服机器人上午说能退,下午又说不能退;还有一个更离谱,Agent 执行到第三步时你刷新了一下页面,整个对话就“失忆”了,用户得重新说一遍问题。

问题的根子在哪?在于我们默认 Agent 是一次性的、无状态的计算过程:你给个 prompt,它调几个工具,吐一段回复,任务结束。但在真实业务里,Agent 的执行时间可能跨分钟、跨小时,甚至跨天。一个完整任务往往要经过多轮工具调用、多个步骤、多次失败重试,期间还会有新事件进来(用户补充信息、审批人点了同意、上游系统返回结果)。这个过程中的中间状态,比如“当前执行到哪一步了”“哪些工具已经调用过”“已经收集到了哪些事实”“下一步该做什么”,如果不被显式地保存和管理,Agent 就谈不上可靠。

把 Agent 理解成分布式状态机,我第一次看到这个说法是在一次内部架构评审上。一位搞过多年订单系统的老哥说:“你们这个 Agent,其实就是个状态机:初始状态、等待工具结果、判断分支、终止。问题是它跑在分布式环境里,进程随时可能被杀、机器随时可能宕、接口随时可能超时,所以状态不能只存在内存里。”这个类比一下就击中了我。传统支付、订单、工作流系统能撑住高并发和故障,靠的就是状态机建模加持久化:状态都落库,事件都可追溯,重试都幂等。Agent 本质上也是在处理一系列有前置依赖的步骤,区别只是步骤的“决策者”从固定代码变成了 LLM,而不确定性更高、状态更加动态。既然决策者变了,基础设施就更不能糊弄。

1.2 长期运行给基础设施带来了什么新问题

传统接口是“请求-响应”模型,几秒内返回,挂了就重试,重试再挂就报错,超时熔断。Agent 完全不一样,它是长跑型任务:一次任务可能由数十个步骤组成,中间还有多轮 LLM 推理,每次推理几百毫秒到几秒不等,再加上工具调用的网络 IO,整体耗时可观。长运行带来的第一个问题就是进程不再可靠——部署、扩容、崩溃都会中断执行。

第二个问题是部分完成不再是“不可接受”的。传统事务追求要么全成功要么全失败,但 Agent 任务往往已经执行了一部分外部操作(发了邮件、创建了工单、扣了款),这时候你不可能回滚整个世界。所以系统必须能回答:已经发生了哪些副作用?哪些步骤已确认完成?哪些还在重试?哪些已经放弃?

第三个问题是LLM 的不确定性被放大。同样的输入,模型可能这次选择工具 A,下次选择工具 B;过程中出现 JSON 解析失败、幻觉参数、循环调用(比如反复搜索同一个关键词)的概率并不低。你没法用“写死代码”来兜底,只能靠编排层、状态持久化、重试策略、超时熔断一并解决。

我把这些痛点总结成一句话:Agent 可靠性的本质,是把“不可预测的模型决策”放进“可预测的执行框架”里。状态机就是那个框架:它在每个步骤之间显式定义状态,所有转移都是可被记录、可恢复、可审计的;LLM 只负责在某个状态下决定“下一步做什么”,而执行引擎负责保证这些步骤最终被可靠地跑完。

1.3 本文会拆解什么内容

结合前阵子搭建 Agent 生产环境的经验,这篇会按下面几条线展开:

  • 状态机的核心定义:有哪些状态、状态怎么迁移、状态存到哪里,这是所有可靠性的基石。
  • 基础设施五条关键链路:队列、持久化、幂等、可观测、人工介入,每一条对长期运行任务都缺一不可。
  • 实操搭建:给出一个可以直接落地的 Agent 执行骨架,包含事件日志、检查点存储、崩溃恢复和幂等控制的关键代码。
  • 排障实录:把我实际遇到过的典型故障整理成速查表,包括 Agent 意外终止时你第一时间应该查什么。

2. 核心细节解析与实操要点

2.1 先把状态机本身定义清楚

如果你去翻 Temporal、AWS Step Functions、LangGraph 这些框架的文档,会发现它们对“步骤执行”的描述惊人地一致:每个任务执行到某个阶段,都会停在一个确定的检查点上,等待下一步指令或外部事件;执行引擎可以随时被中断,再从最后一个检查点继续。

一个可长期运行的 Agent 状态机,状态定义一般包含:

状态含义核心特征
INITIALIZED任务刚创建,尚未开始已分配全局唯一的执行 ID
PLANNINGLLM 正在生成计划或推理下一步可能多次调用模型,结果未落库前不算完成
WAITING_TOOL_CALL已决定调用某个工具,等待工具返回记录工具名、入参、请求 ID
WAITING_EXTERNAL_EVENT需要等待用户/审批/上游系统事件没有超时上限或超时策略明确
RETRYING上一次执行失败,进入重试记录失败原因、重试次数、退避策略
TERMINATED_SUCCESS正常完成,产出最终结果结果已归档
TERMINATED_FAILED不可恢复失败,任务终止有失败快照和人工介入入口

设计时最容易犯的错是把状态定义得“太粗”:整个 Agent 只有一个RUNNING状态,那出问题你根本定位不到是哪一步。正确做法是状态要细化到可能的挂起点——凡是可能等待 IO、等待模型、等待用户反馈的位置,都值得一个独立状态。这样才能精确回答“这活儿现在搁在哪了”。

转折点、边界条件也要提前想清楚。比如 LLM 给出的工具调用参数不合法,这不算FAILED,应该回到PLANNING让模型重新生成;比如某工具连续重试三次仍失败,要不要换个工具?这些分支逻辑看似简单,但在状态机里每个分支都对应一条转移路径,路径越多,越需要结构化定义、测试覆盖。

2.2 状态存储的分层设计

状态存哪里,直接影响系统的恢复能力和并发行为。我目前落到项目里的方案是三层:

  • 会话/上下文层(瞬时态):当前轮对话的 prompt 拼接、最近的模型输出,放在内存缓存里就好,进程重启后丢弃也无所谓。
  • 检查点层(中间态):每个步骤执行完,把“执行到什么位置、已完成哪些工具调用、已收集哪些数据、剩余待办”写入持久化存储。这是崩溃恢复的关键,放 Redis 或数据库都行。
  • 记录层(终态):整个任务完成后的结果、最终对话历史、审计日志,写入数据仓库或文档库,长期保留。

先说检查点。我见过不少同学把整段对话历史直接塞进 Redis,认为这就是“持久化”了,其实不对。对话历史只是素材,真正需要保存的是执行进度。打个比方:你请了一个代理去帮你跑腿办事,你关心的是“事情办到哪一步了”,而不是“代理每一步说了什么话”。Agent 的状态机也一样——模型推理中间产生的那些草稿不重要,重要的是“已经调用了下单接口”这样的事实。所以检查点里至少要包含:

  • 当前状态,以及进入该状态的版本号
  • 已完成工具调用的结果摘要(不一定要全文,但要够恢复时用)
  • 待处理的队列(比如模型列了三个工具要调用,调完了两个,剩一个)
  • 执行环境的元信息(用的哪个模型、什么温度参数,保证恢复时上下文接近)

Redis 做检查点存储时,建议用 Hash 结构而不是 String。我之前贪方便直接SET key json,结果多实例并发写时互相覆盖,丢失了一大段进度。改成 HSET 后,每个字段单独更新,配合 WATCH 或 Lua 脚本做原子操作,问题才消除。

2.3 状态转移过程中最容易忽略的细节

定义好了状态和存储,转移逻辑还有几个细节容易被忽略。

第一,工具调用的请求幂等键必须进状态。假设 Agent 调用“发起退款”工具,请求发出去了,但网络超时。重试时如果拿同样的参数再发一次,用户就会被退两次款。正确做法是:在发起调用前生成一个全局唯一的幂等键(比如execution_id + step_id + tool_name),存进状态,工具侧按这个键做去重。这不是 Agent 特有的问题,任何分布式系统都有,但 Agent 的重试频率比普通接口高得多,出现概率更大。

第二,状态转移必须带上触发事件。不要只记录“当前是 WAITING_TOOL_CALL”,还要记录“是因为哪个调用返回了结果才转移到这里”。这其实就是事件溯源的思想,把状态变化做成 append-only 的日志流:tool_call_started、tool_call_succeeded、plan_generated、agent_paused_waiting_user。恢复的时候不需要靠猜,回放事件就能重建现场。

第三,要预留人工介入的入口。长期运行的任务,总会出现模型怎么都绕不出来的死胡同,比如连续重试五次同一个搜索、同一个审批被反复驳回等。设计时给状态机加一个WAITING_HUMAN状态,当重试耗尽或置信度过低时进入,把现场快照交给人工处理。别觉得这样“不够智能”,在生产环境里,可靠的人工兜底往往比模型自我修复更高效。

2.4 工具选型的核心标准

市面上能帮你搭这套东西的框架不少,我梳理下选型时最该看重的几条标准:

  • 是否原生支持持久化和恢复。LangGraph 有 checkpointer,Temporal 有 workflow 的确定性重放,Step Functions 天生按状态机建模。如果只是普通库,那你得自己补状态存储。
  • 是否提供幂等和重试语义。Temporal 的 Activity 自带重试、超时和幂等控制;Step Functions 的 Task 自带 Retry 和 Catch;自研的话,这些全得自己写。
  • 是否支持人工干预。生产级平台大多有“暂停流程”“人工审批”“终止执行”等操作,这对长期运行任务至关重要。
  • 生态和排查体验。出了故障,能不能看到执行到哪一步、每个步骤花了多久、失败原因是什么?Temporal 的 Web UI、Step Functions 的执行历史都做得不错,自研的话这部分成本极高。

我自己在重业务的场景倾向于用专门的编排平台(Temporal 或 Step Functions),轻量场景用 LangGraph + 自建存储。核心逻辑是一致的:执行引擎负责可靠,模型负责决策,二者的边界要划清楚。

3. 实操过程与核心环节实现

3.1 从零搭建可恢复 Agent 的总体流程

接下来给一个可落地的参考实现。场景是“订单售后 Agent”:用户提交售后申请,Agent 需要查询订单、校验权限、计算退款金额、发起退款、最后通知用户。流程会跨多个步骤,中间允许用户补充材料,也允许审批延迟。

整个骨架分五层:

  1. 状态定义层:枚举所有状态和事件。
  2. 事件日志层:所有状态转移都追加写入事件流(这里用 Redis Stream 演示)。
  3. 检查点存储层:当前状态和关键上下文写入 Redis Hash。
  4. 执行引擎层:从事件流里取任务,更新状态,调用工具,写回检查点。
  5. 恢复层:进程重启或任务中断后,从检查点恢复执行。

用到的东西:Python 3.10+、Redis、一个 LLM 客户端(这里用 OpenAI 风格的接口做示意)、一个模拟订单系统。代码层面我不会贴完整项目,会把关键片段都展开讲。

3.2 步骤一:状态定义与事件模型

先定义好枚举,这一步千万别偷懒,所有状态转移都以它为准:

# states.py from enum import Enum class AgentState(str, Enum): INITIALIZED = "initialized" PLANNING = "planning" WAITING_TOOL_CALL = "waiting_tool_call" WAITING_EXTERNAL_EVENT = "waiting_external_event" RETRYING = "retrying" TERMINATED_SUCCESS = "terminated_success" TERMINATED_FAILED = "terminated_failed" WAITING_HUMAN = "waiting_human" class AgentEvent(str, Enum): TASK_CREATED = "task_created" PLAN_GENERATED = "plan_generated" TOOL_CALL_STARTED = "tool_call_started" TOOL_CALL_SUCCEEDED = "tool_call_succeeded" TOOL_CALL_FAILED = "tool_call_failed" EXTERNAL_EVENT_RECEIVED = "external_event_received" HUMAN_APPROVED = "human_approved" HUMAN_REJECTED = "human_rejected" TASK_COMPLETED = "task_completed" TASK_ABORTED = "task_aborted"

所有事件统一用这样一个结构体:

# events.py from dataclasses import dataclass, field from datetime import datetime from typing import Any, Optional import uuid @dataclass class AgentEventRecord: event_id: str = field(default_factory=lambda: uuid.uuid4().hex) execution_id: str = "" event_type: str = "" payload: dict = field(default_factory=dict) timestamp: datetime = field(default_factory=datetime.utcnow) idempotency_key: str = ""

这里有个细节:idempotency_key不是可选项。每次工具调用生成一个,存到事件里,后续无论重试多少次,下游服务凭这个键做去重。我给它取成execution_id + step_id + tool_name,同一个执行任务中同一个 tool 的调用一定复用同一个键。

3.3 步骤二:事件日志与检查点存储

事件日志使用 Redis Stream,每个执行任务对应一个独立的 Stream key,事件追加写,消费组支持多个 worker 并发处理不同任务:

# event_store.py import redis import json class EventStore: def __init__(self, redis_client: redis.Redis): self.redis = redis_client def append(self, record: AgentEventRecord) -> None: key = f"agent_events:{record.execution_id}" value = json.dumps({ "event_id": record.event_id, "event_type": record.event_type, "payload": record.payload, "timestamp": record.timestamp.isoformat(), "idempotency_key": record.idempotency_key, }) self.redis.xadd(key, {"data": value}) def replay(self, execution_id: str) -> list: key = f"agent_events:{execution_id}" raw = self.redis.xrange(key) events = [] for _, item in raw: data = json.loads(item[b"data"].decode()) events.append(data) return events

检查点存储用 Redis Hash。每次步骤执行完成,原子地更新当前状态和关键上下文:

# checkpoint_store.py import redis import json class CheckpointStore: def __init__(self, redis_client: redis.Redis): self.redis = redis_client def save(self, execution_id: str, state: str, context: dict, version: int): key = f"agent_checkpoint:{execution_id}" # 用 Lua 脚本保证版本号递增,防止旧写入覆盖新写入 script = """ local current = tonumber(redis.call('HGET', KEYS[1], 'version') or '0') if tonumber(ARGV[2]) >= current then redis.call('HSET', KEYS[1], 'state', ARGV[1]) redis.call('HSET', KEYS[1], 'context', ARGV[3]) redis.call('HSET', KEYS[1], 'version', ARGV[2]) return 1 end return 0 """ self.redis.eval(script, 1, key, state, version, json.dumps(context)) def load(self, execution_id: str) -> dict | None: key = f"agent_checkpoint:{execution_id}" data = self.redis.hgetall(key) if not data: return None return { "state": data.get(b"state", b"").decode(), "context": json.loads(data.get(b"context", b"{}")), "version": int(data.get(b"version", b"0")), }

版本号这个细节我觉得值得展开讲讲。Redis 单键读写看起来没啥问题,但 Agent 可能被多个 worker 同时处理——比如两个 worker 同时消费到同一个任务的重试消息,各自执行了一步,写检查点时后写的覆盖先写的,进度就丢了。版本号 + Lua 脚本的乐观锁能保证只有版本更高的写入生效,代价很小,收益很大。

3.4 步骤三:执行引擎核心循环

执行引擎是整台机器的心脏。设计上采用事件驱动:主循环从事件流消费记录,根据当前状态执行对应动作,然后把新事件写回流,并更新检查点。核心循环长这样:

# engine.py import redis import json from .states import AgentState, AgentEvent from .events import AgentEventRecord from .event_store import EventStore from .checkpoint_store import CheckpointStore class AgentEngine: def __init__(self, llm_client, tool_registry, redis_client): self.llm = llm_client self.tools = tool_registry # 工具注册表:{"tool_a": callable} self.events = EventStore(redis_client) self.checkpoints = CheckpointStore(redis_client) def start(self, execution_id: str, initial_payload: dict): record = AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.TASK_CREATED.value, payload=initial_payload, ) self.events.append(record) self.checkpoints.save(execution_id, AgentState.PLANNING.value, initial_payload, version=1) def step(self, execution_id: str) -> str: cp = self.checkpoints.load(execution_id) if cp is None: raise ValueError(f"Checkpoint not found: {execution_id}") state = cp["state"] context = cp["context"] if state == AgentState.PLANNING.value: return self._do_planning(execution_id, context, cp["version"]) elif state == AgentState.WAITING_TOOL_CALL.value: return self._do_tool_call(execution_id, context, cp["version"]) elif state == AgentState.WAITING_EXTERNAL_EVENT.value: # 外部事件到来才继续 return state elif state == AgentState.TERMINATED_SUCCESS.value: return state else: return self._do_failure_handling(execution_id, context, cp["version"]) def _do_planning(self, execution_id, context, version): # 调用 LLM,得到下一步计划 plan = self.llm.chat(messages=context["messages"], tools=self.tools.list()) if plan.tool_calls: tool_call = plan.tool_calls[0] new_context = dict(context) new_context["pending_tools"] = [{ "tool_name": tool_call.function.name, "arguments": tool_call.function.arguments, "idempotency_key": f"{execution_id}:{self._step_count(context)}:{tool_call.function.name}", }] self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.PLAN_GENERATED.value, payload={"plan": plan.model_dump()}, )) self.checkpoints.save( execution_id, AgentState.WAITING_TOOL_CALL.value, new_context, version=version + 1, ) return "planning_done" else: # 没有工具调用,任务完成 self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.TASK_COMPLETED.value, payload={"final_answer": plan.content}, )) final_context = dict(context) final_context["final_answer"] = plan.content self.checkpoints.save( execution_id, AgentState.TERMINATED_SUCCESS.value, final_context, version=version + 1, ) return "completed" def _do_tool_call(self, execution_id, context, version): pending = context.get("pending_tools", []) if not pending: # 工具都调完了,回去继续规划 self.checkpoints.save( execution_id, AgentState.PLANNING.value, context, version=version + 1, ) return "tool_batch_done" tool_spec = pending.pop(0) tool_name = tool_spec["tool_name"] tool_func = self.tools.get(tool_name) idem_key = tool_spec["idempotency_key"] self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.TOOL_CALL_STARTED.value, payload=tool_spec, idempotency_key=idem_key, )) try: result = tool_func(**tool_spec["arguments"]) new_context = dict(context) new_context["pending_tools"] = pending new_context.setdefault("tool_results", []).append({ "tool_name": tool_name, "arguments": tool_spec["arguments"], "result": result, }) # 注意:同一个幂等键,结果只记录一次 self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.TOOL_CALL_SUCCEEDED.value, payload={"tool_name": tool_name, "result": result}, idempotency_key=idem_key, )) self.checkpoints.save( execution_id, AgentState.PLANNING.value if pending else AgentState.PLANNING.value, new_context, version=version + 1, ) return "tool_call_succeeded" except Exception as e: retry_count = context.get("retry_count", 0) + 1 if retry_count < 3: new_context = dict(context) new_context["retry_count"] = retry_count new_context["last_error"] = str(e) self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.TOOL_CALL_FAILED.value, payload={"tool_name": tool_name, "error": str(e), "retry_count": retry_count}, idempotency_key=idem_key, )) self.checkpoints.save( execution_id, AgentState.RETRYING.value, new_context, version=version + 1, ) else: # 超过重试,进入人工兜底 self.events.append(AgentEventRecord( execution_id=execution_id, event_type=AgentEvent.EXTERNAL_EVENT_RECEIVED.value, payload={"type": "human_review_required", "tool": tool_name, "error": str(e)}, )) self.checkpoints.save( execution_id, AgentState.WAITING_HUMAN.value, context, version=version + 1, ) return "tool_call_failed"

这套循环有几个特点值得说。

第一,LLM 只决定下一步计划,工具调用、重试、恢复都是引擎的事情。这样模型再不稳定,执行框架还是稳的。

第二,每个事件都带幂等键。即使同一个工具调用因为重试被重复消费,下游也能识别。实际运行中,并发消费、重复投递比想象中常见得多,这个键救过我好几次。

第三,重试次数放进了 context 而不是放在本地变量。重启后还能接着重试,而不是从零再来。

第四,所有分支都会更新检查点版本号,每一步都有据可查,谁改了状态、改了哪个版本,事后都能对上。

3.5 步骤四:恢复流程长什么样

有了检查点和事件日志,恢复逻辑其实很简单:

# recovery.py def recover(execution_id: str, engine: AgentEngine): cp = engine.checkpoints.load(execution_id) if cp is None: # 无检查点 = 从未开始,或者检查点过期被清理 raise RuntimeError(f"No checkpoint found for {execution_id}") state = cp["state"] if state in (AgentState.TERMINATED_SUCCESS.value, AgentState.TERMINATED_FAILED.value): return state # 已经终态,不需要恢复 if state == AgentState.WAITING_HUMAN.value: # 需要人工介入,不能自动继续 return state # 其余状态都可以从检查点继续推进 while True: next_state = engine.step(execution_id) if next_state in ( AgentState.TERMINATED_SUCCESS.value, AgentState.TERMINATED_FAILED.value, AgentState.WAITING_EXTERNAL_EVENT.value, AgentState.WAITING_HUMAN.value, ): break return next_state

恢复时最怕的就是“检查点有,但事件日志对不上”。我遇到过一次:进程在工具调用成功后、写检查点前崩了,事件日志里已有tool_call_succeeded,但检查点还停留在WAITING_TOOL_CALL。恢复后引擎会再调一次这个工具,因为幂等键相同,下游直接返回之前的成功结果,不会重复执行,所以最终一致了。这就是事件+检查点双写的意义:要么两边都一致,要么靠幂等键兜底。

3.6 步骤五:接入定时调度与外部事件

长期运行任务必然涉及“等待”。比如等待用户补充材料,或者等待审批人点击“同意”。有两种实现思路:

  • 轮询:定时器每隔一段时间扫描处于WAITING_EXTERNAL_EVENT的任务,查询外部系统状态。
  • 回调:外部系统状态变化时,通过 Webhook 回调到引擎,引擎推送一个EXTERNAL_EVENT_RECEIVED事件,再触发step()。

回调更及时,但需要暴露公网入口,也要考虑回调丢失。稳妥的做法是回调为主、定时扫描为辅。Redis Stream 的XREAD BLOCK天然支持阻塞等待新事件,非常适合等待外部事件的场景:

# 消费外部事件的 worker def external_event_worker(engine, redis_client, execution_id): stream_key = f"agent_events:{execution_id}" # 等待新事件,阻塞 30 秒 events = redis_client.xread({stream_key: "0"}, count=10, block=30000) if events: for _, items in events: for _, raw in items: data = json.loads(raw[b"data"]) if data["event_type"] in ( "external_event_received", "human_approved", "human_rejected", ): engine.step(execution_id) break

这里还有个大坑:外部事件和当前检查点可能不一致。比如用户已经补充完材料了,但检查点还停留在调用查询接口的状态,这时直接把external_event_received当成“继续推进”信号,可能跳过中间步骤。我在生产环境里加了一个校验:先检查当前状态是否真的能处理这个事件类型,不能的话先把事件入待处理队列,等主流程推进到对应状态再消费。

4. 基础设施可靠性的五条关键链路

4.1 队列与调度:别让任务裸奔

长任务最忌讳“裸奔”——没有队列,直接在线程池里跑,进程一重启全没了。生产环境至少要有一条持久化的任务队列,比如 Redis Stream 或 Kafka。任务创建时把执行 ID 推入队列,worker 消费后处理,处理完成要 ack。Redis Stream 的消费者组机制可以天然实现任务分发和 ack,Kafka 则利用 partition 保证同一执行 ID 的消息有序。

调度的频率也要斟酌。有些任务要等外部事件,在事件没来之前,不要盲目轮询,否则既浪费资源,又可能产生大量无意义的状态更新。设计原则是:能靠回调就少用轮询,能靠事件驱动就别写 sleep 循环。

4.2 编排引擎与确定性重放

如果你不想踩这么多底层实现的坑,直接上 Temporal 或 Step Functions 这类平台是更稳妥的选择。Temporal 的 workflow 是确定性代码:同样的事件序列,在任何机器上重放,结果都一致。这意味着恢复执行时不用真去重新调用工具,而是把以前的事件回放一遍,代码跑到底,状态自然恢复。

不过这类平台的学习曲线不低,尤其 Temporal,涉及 worker、namespace、activity heartbeat 等一堆概念。自研方案的优点是完全可控,缺点也很明显:所有边界情况都得自己处理。我的看法是:如果业务里 Agent 任务量大、并发高、价值敏感,直接上平台更划算;如果只是十几个并发任务的内部工具,自研轻量引擎也够用。

4.3 幂等设计是最不能省的一层

前面反复提到幂等键,这里专门讲一下。幂等设计有三个层级:

  1. 请求级幂等:每个工具调用带idempotency_key,下游服务按 key 去重。
  2. 状态转移幂等:同一个事件重复处理,不会导致状态重复转移。我的做法是在事件表里对(execution_id, event_id)建唯一索引,重复消费直接跳过。
  3. 任务重启幂等:恢复流程不能因为重启就重新触发副作用。检查点记录哪些工具已经成功,恢复时直接跳过已完成步骤。

第三点最容易漏。很多 Agent 框架默认重启后重跑全部流程,如果中途有“发短信通知”这类副作用,用户就会收到一堆重复消息。我们线上遇到过一次,后来在检查点里加了一个completed_side_effects列表,每次工具成功就追加,恢复时先查这个列表,已经完成的直接拿结果,不再调用。

4.4 可观测性:日志、追踪与审计

长任务的排查和普通接口不一样。普通接口出问题,看一个 trace 就能定位;Agent 出问题,要看你根本不知道任务跑到哪儿了。所以可观测性设计要围绕“执行现场”来做。

Temporal 和 Step Functions 都自带执行历史 UI,能看到每个步骤的输入输出。自研的话,我建议至少做三件事:

  • 结构化日志:每条日志带上execution_id、state、event_type,这样grep execution_id能串起整个生命周期。
  • OpenTelemetry 追踪:每次工具调用和 LLM 调用都生成一个 span,parent span 是整个任务的执行。这样能看到一次任务总共调用了多少次模型、哪些工具最慢、哪一步耗时最长。
  • 审计事件流:把状态转移、工具调用结果、人工审批记录全部落库。合规审计先不说,出纠纷时拿得出证据这条就值回成本。

4.5 降级与死信:留一条人工后路

最后一条链路是降级与人工介入。所有自动化系统都会遇到无法自动处理的场景,比如模型连续生成非法 JSON、工具依赖的下游系统宕机、需要用户提供凭证等。设计状态机时一定要有WAITING_HUMAN这个兜底状态,并且提供一套方便人工处理的界面:执行现场快照、已执行的工具结果、错误原因、当前上下文,一目了然。

别高估模型的自愈能力。我见过 LLM 在同一个错误上反复转圈十几次,就是不肯换个方案。这时候靠状态机里的“最大重试次数”直接终止并转人工,比让它继续烧钱高效得多。

5. 常见问题与排查技巧实录

5.1 故障速查表

我把实际遇到过的问题整理成一张表,排查顺序按出现频率从高到低排:

现象可能原因排查方法解决建议
Agent 任务卡住不动,状态停在上一步事件丢失或 worker 未 ack查看事件流消费组 pending 列表补投未 ack 消息,检查 worker 心跳
重启后重复执行工具检查点没记录已完成工具查看检查点completed_side_effects恢复时跳过已完成副作用
两个 worker 同时处理同一任务队列未按 execution_id 分区检查消费组配置用一致性哈希或用 Redis 分布式锁
状态被旧版本覆盖并发写检查点无版本控制查看 Redis Hash 的 version 字段用 Lua 脚本乐观锁
LLM 反复调用同一个工具模型陷入循环查看事件流中tool_call_started次数设最大工具调用次数,超出转人工
重试风暴:下游被打爆无退避策略查看日志时间戳间隔加指数退避和 jitter
外部事件到了但主流程没反应事件消费顺序不对查看事件流顺序先校验当前状态能否消费该事件
恢复后对话上下文对不上检查点未存 messages查看上下文快照同时持久化 messages 摘要和完整消息

5.2 最经典的坑:状态更新和事件写入的一致性

有一次线上任务,用户反馈“订单状态说已退款,但用户收到通知说退款失败”,两边对不上。一查,事件流里tool_call_succeeded写了,检查点也保存了,但实际工具调用因为网络超时压根没到达下游。原因是引擎在超时后把“假设成功”的写入顺序搞反了:先写了成功事件,再拿到工具侧确认。

正确顺序应该是:先向工具发出调用并同步等待确认结果,再写事件和检查点。如果是异步工具,要等确认回调。这个顺序我给团队立了铁律:工具结果的写入必须以“下游确认为准”,不能以“调用发出为准”。这个坑不踩一次,很难意识到有多致命。

5.3 排查时的高效姿势

Agent 问题的定位,先看事件流,再看检查点,最后看代码。事件流能告诉你发生了什么,检查点能告诉你当前停在哪,代码才能告诉你为什么。

我自己常用的排查命令(Redis CLI):

# 查看某个执行任务的所有事件 XREVRANGE agent_events:{execution_id} + - COUNT 50 # 查看检查点当前状态 HGETALL agent_checkpoint:{execution_id} # 查看队列中未处理的消息 XPENDING agent_queue:{queue_name} {group_name}

另外我强烈建议给每个工具调用加一个tool_call_id,打印日志时带上。有一次排查一个大模型反复调用搜索工具的问题,就是靠这个 ID 发现它十次调用里九次参数都一样——模型明显卡住了。定位之后加了“连续相同调用次数超限”的判断,转人工后问题迎刃而解。

5.4 关于框架的选型,说点实在的

虽然我前面给了自研方案,但如果你刚开始做,我还是推荐直接基于成熟框架起步:

  • LangGraph:上手快,自带 checkpointer,适合快速验证状态机设计思路。不好的地方是生产级的重试、幂等、可观测要自己补。
  • Temporal:稳如老狗,重试、超时、幂等、可观测全部内置,适合生产。学习曲线陡,但对长期运行任务来说,这些成本是值得的。
  • Step Functions:如果你已经深度用云服务,直接选它最省心,状态机和重试语义都是原生内置。

自研引擎适合的场景是:团队已经有很强的分布式系统功底,或者有特别定制化的需求(比如 Agent 状态和业务数据深度耦合)。否则,框架提供的可靠性能力远超自己造的轮子。

6. 最后分享一些个人体会

我在实际项目中踩过不少坑之后,最大的体会是:做 Agent 生产系统,别把重心放在让模型更聪明上,先把执行框架做扎实。模型今天会的,明天可能更新就忘了;但状态机、幂等、检查点这套东西,不管模型怎么变,都是可靠的底座。

从设计第一天就决定用“分布式状态机”思维来约束 Agent,团队里的沟通效率也高了。以前讨论方案,大家总纠结“这个 Agent 场景要不要用 LangChain 还是自研 prompt”,后来统一叫法,变成“这个任务从哪个状态开始、中间要经过哪些状态、哪些可以并行、哪些失败可以重试、哪些必须人工”,讨论立刻变得具体可落地。

最后再分享一个细节:写 Agent 代码前,先把状态转移图画出来,不需要画得很规范,哪怕手画一张草稿都行。这张图能帮你梳理清楚隐藏的分支——哪些地方会等待、哪些地方会失败、哪些地方需要人工介入。状态图定了,代码只是抄图的过程;状态图没定,代码全靠拍脑袋,出问题是迟早的事。

这套思路不止适用于 Agent,任何带外部依赖、长耗时、不可预测的执行流程,都值得用状态机重新建模一遍。哪怕最终没有把它彻底落地成完整的事件溯源系统,只要把“状态显式化、转移可视化、恢复可操作”这三点做到位,系统的可靠性就已经远超大多数同类项目了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询