用 Temporal 给 Agent 长任务加持久化与可恢复(崩溃不丢进度)
你的 Agent 跑一个长任务:调了三个外部 API、读了两次数据库、刚要写结果——进程崩了。重启后,前面三步白做,还可能出现「钱扣了但单没下」的半成品状态。普通脚本对此束手无策:没有重试、没有进度记录、没有续跑。Temporal 用「持久执行(Durable Execution)」把多步 Agent 任务变成可恢复的 Workflow——进程崩了、API 超时了,它都能从断点接着跑。本教程用 Python SDK 跑通最小可复现样例。
先搞懂:持久执行是什么?
Temporal 把「业务步骤」拆成两层:Workflow(编排顺序,状态自动持久化)和 Activity(真正干活的外部调用,可重试)。每执行一步,状态就被记下来;Worker 崩了,就从最后一步续跑,不会重头来。
| 层 | 职责 | 能否崩后续跑 |
|---|---|---|
| Workflow | 编排步骤、状态、信号 | 能,状态持久化 |
| Activity | 调 API / 写库 / 跑模型 | 能,自动重试 |
Step 1:为什么 Agent 需要持久执行
# 一个会「半途而废」的脚本
resize_compute(db) # 第1步成功
update_load_balancer(db) # 第2步成功
update_dns(db) # 崩在这里 → 前两步无法回滚、无记录
普通脚本没有进度记录、不会自动重试、不会回滚。长任务只要跨多个外部调用,就必须自己写对账与补偿逻辑——这正是 Temporal 替你做的事。
Step 2:起 Temporal Dev Server
pip install temporalio
# 另开一个终端,启动本地 Dev Server(默认端口 7233)
temporal server start-dev
Step 3:写 Activity(把外部调用放进来)
LLM 调用、HTTP 请求、数据库写入都放 Activity;把底层 SDK 的自动重试关掉,交给 Temporal 统一重试。
from dataclasses import dataclass
from openai import AsyncOpenAI
from temporalio import activity
@dataclass
class LLMRequest:
model: str
instructions: str
input: str
@activity.defn
async def call_llm(req: LLMRequest) -> str:
client = AsyncOpenAI(max_retries=0) # 重试交给 Temporal
resp = await client.chat.completions.create(
model=req.model,
messages=[{"role":"system","content":req.instructions},
{"role":"user","content":req.input}],
timeout=15,
)
return resp.choices[0].message.content
非确定性代码别进 Workflow:Workflow 可能被重放(replay),里面不能有用随机、当前时间、直连外部 API 等非确定性操作;这些一律放 Activity。
Step 4:写 Workflow(编排 Agent 循环)
Workflow 用 workflow.execute_activity 调 Activity,并给每个步骤设超时与重试策略。官方示例的 Agent 循环骨架如下:
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy
with workflow.unsafe.imports_passed_through():
from activities import call_llm
@workflow.defn
class AgentWorkflow:
@workflow.run
async def run(self, goal: str) -> str:
messages = [{"role":"user","content": goal}]
policy = RetryPolicy(
initial_interval=timedelta(seconds=2),
maximum_attempts=5,
)
while True:
r = await workflow.execute_activity(
call_llm,
LLMRequest(model="gpt-4o-mini",
instructions="你是助手。",
input=str(messages)),
start_to_close_timeout=timedelta(seconds=60),
retry_policy=policy,
)
messages.append({"role":"assistant","content": r})
if "DONE" in r:
return r
# 否则继续循环,状态全程持久化
Step 5:写 Worker 与启动器
# worker.py
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from workflows.agent import AgentWorkflow
from activities import call_llm
async def main():
client = await Client.connect("localhost:7233")
worker = Worker(client, task_queue="agent-q",
workflows=[AgentWorkflow], activities=[call_llm])
await worker.run()
# start.py
async def main():
client = await Client.connect("localhost:7233")
h = await client.execute_workflow(
AgentWorkflow.run, "调研 A 公司并给摘要",
id="agent-1", task_queue="agent-q")
print(await h)
序列化:Workflow 输入输出需可序列化;用 Pydantic 模型时给 Worker/Client 配 pydantic_data_converter,否则会报序列化错误。
Step 6:加人审与信号
长任务可中途等待人工信号(如「批准继续」)再往下走,用 workflow.wait_for_signal 实现,不占线程、不丢状态。
@workflow.defn
class ApprovalAgent:
@workflow.run
async def run(self, goal: str):
# 做一半,等人工批准
await workflow.wait_for_signal("approve")
# 批准后继续后续 Activity
...
wait_for_signal 等人点确认,再执行写操作——天然支持 human-in-the-loop,且中途崩溃不影响已确认状态。Step 7:生产化提醒
· Workflow 内状态保持紧凑,只存恢复所需,勿存无限增长日志
· 超长任务用 Continue-as-New 避免历史过大
· 复杂任务拆 Child Workflow,隔离失败域
· 绝不在 Workflow 写随机/时间/直连外部等非确定逻辑
· Activity 设合理超时与重试上限,避免无限重试
重放一致性:Workflow 代码改了之后,旧正在运行的执行会按新代码重放,务必保证新代码对相同输入产生相同步骤序列,否则重放失败。变更前先读完 Temporal 关于确定性Workflow的文档。
常见问题速查
| 现象 | 原因与解决 |
|---|---|
| 连不上 7233 | Dev Server 没起,或地址/端口错 |
| 序列化报错 | 用了 Pydantic 但未配 pydantic_data_converter |
| 重放失败 | Workflow 含非确定逻辑,移到 Activity |
| 一直重试 | Activity 未设最大重试,或外部永久失败需人工 |