跳到主内容
快讯直播
AI智模界
教程

Temporal 编排长时 AI Agent:断点续跑与人工审批

这篇能做出什么

先看成果。你要做的是一个"竞品调研 Agent":用户丢进来一个主题,Agent 自动拆关键词、抓网页、写草稿,然后把草稿推给人审批;审批人说"通过"就发布,说"不行,补一下定价数据"就打回去改,改完再送审。

这个流程听起来简单,但真正跑起来会遇到一堆现实问题:

  • 一次调研要跑十几分钟到几小时,中间要等审批人三天后才回话;
  • 你的进程重启、机器重启、部署发版,正在跑的任务不能丢;
  • 调用 LLM 和抓网页会遇到 429、502、超时,需要重试但不能重复扣费、重复发布;
  • 审批人回复的时候,那个等了三天的进程早就不在了。

普通写法的解法是:起一个数据库存状态,写一个 while True 轮询,自己记 stage 字段,自己写重试和去重。写到后面你会发现,你 80% 的代码在写调度器,20% 在写 Agent。

Temporal 把这一层接过去了:它把工作流的每一步都记进事件历史(Event History),进程崩了重启后,代码会从头"重放",已经完成的步骤直接从历史里读结果,不会真的再跑一遍。这就是"断点续跑"的实现方式。人工审批用 Signal 实现,等三天就是一个持久化定时器,不占任何进程资源。

跑完本文,你会得到一个能演示"杀掉 Worker 再启动,任务继续跑"的完整工程骨架。

前置条件清单

动手前确认这几件事:

1. 一个能跑 Docker 的机器,用来起本地 Temporal 开发服务。

2. Python 3.10 及以上(具体支持版本以官方文档为准)。本文用 Python SDK,Temporal 也有 TypeScript / Go / Java / .NET / Ruby SDK,概念完全一致,换语言只是换语法。

3. 一个 LLM 的 API Key(OpenAI、Anthropic、或国内的任意一家都行)。本文把它封装成 call_llm(),不绑定具体厂商。

4. 基础的 async/await 认知。Temporal 的 Python SDK 是异步的,但用法和普通 asyncio 很像。

5. 一个搜索或抓取接口,用来拿网页内容。没有的话用任何能返回文本的 HTTP 接口替代即可。

先建立一个心智模型,这五个词后面会反复出现:

概念一句话解释
Workflow你写的编排逻辑。必须确定性:同样的输入和历史,跑出来必须一样
Activity所有脏活:调 LLM、抓网页、写数据库、发通知。可重试、可超时、可心跳
Worker一个常驻进程,向 Temporal 领任务并执行 Workflow / Activity
Task QueueWorker 和任务的匹配队列。不同任务用不同队列,别互相堵
Signal从外部异步塞进 Workflow 的事件,人工审批就靠它

记住一条铁律:Workflow 里不碰网络、不读系统时间、不用随机数,这些都放进 Activity。 违反了这条,重放时会报非确定性错误。

分步骤

第 1 步:起一个本地 Temporal 服务

开发阶段最省事的方式是用 Temporal CLI 的开发服务器:

```bash

具体安装方式与命令参数以官方文档为准

temporal server start-dev

```

它会同时拉起服务端和 Web UI,UI 默认在本地 8233 端口(端口号以官方文档为准)。打开 UI,你能看到所有工作流实例、每个实例的事件历史、正在重试的活动——这个界面在后面排错时会救你的命。

不想装 CLI 就用 Docker Compose 起服务端加数据库,配置文件参考官方仓库的示例。

第 2 步:搭项目骨架

```bash

mkdir ai-agent-temporal && cd ai-agent-temporal

python -m venv .venv && source .venv/bin/activate

pip install temporalio httpx

再按你的 LLM 厂商装对应的 SDK,版本以官方文档为准

```

目录结构:

```

ai-agent-temporal/

├── activities.py # 所有脏活

├── workflows.py # 编排逻辑

├── worker.py # 常驻进程

├── client.py # 启动任务 / 发审批信号

└── llm.py # LLM 调用的薄封装

```

第 3 步:定义数据结构和 Activity

先写共享的数据结构。注意 PageDigest 里只有摘要,没有网页全文——这是一个刻意的设计,后面讲坑的时候会说为什么。

```python

activities.py

from dataclasses import dataclass

import httpx

from temporalio import activity

from llm import call_llm # 你自己的 LLM 封装,读环境变量拿 API Key

@dataclass

class ResearchBrief:

topic: str

audience: str

@dataclass

class PageDigest:

url: str

title: str

summary: str # 只传摘要,不传全文

@dataclass

class Draft:

title: str

body: str

sources: list[str]

@dataclass

class DraftRequest:

brief: ResearchBrief

digests: list[PageDigest]

```

接下来是最关键的一步:LLM 调用和网页抓取全部写成 Activity。

```python

@activity.defn

async def plan_queries(brief: ResearchBrief) -> list[str]:

activity.logger.info("规划检索词: %s", brief.topic)

text = await call_llm(

system="你是调研助手。输出 3 个检索关键词,每行一个,不要编号。",

user=f"主题:{brief.topic}\n受众:{brief.audience}",

)

return [line.strip() for line in text.splitlines() if line.strip()][:3]

@activity.defn

async def fetch_page(query: str) -> PageDigest:

长任务打心跳,Worker 掉线时能更快被发现并重试

activity.heartbeat("fetching")

async with httpx.AsyncClient(timeout=20, follow_redirects=True) as client:

换成你真实的搜索 / 抓取接口

url = f"https://example.com/search?q={query}"

resp = await client.get(url)

resp.raise_for_status()

raw = resp.text[:20000]

activity.heartbeat("summarizing")

summary = await call_llm(

system="把下面内容压缩成 200 字内的中文摘要,保留关键数字,不要编造。",

user=raw,

)

return PageDigest(url=url, title=query, summary=summary)

@activity.defn

async def write_draft(req: DraftRequest) -> Draft:

material = "\n\n".join(f"[{d.url}]\n{d.summary}" for d in req.digests)

body = await call_llm(

system="基于给定材料写一份中文调研草稿,注明每条结论的来源链接。",

user=f"主题:{req.brief.topic}\n受众:{req.brief.audience}\n\n材料:\n{material}",

)

return Draft(title=req.brief.topic, body=body,

sources=[d.url for d in req.digests])

@activity.defn

async def notify_reviewer(draft: Draft) -> None:

发邮件 / 发飞书卡片 / 写待办表,带上 workflow_id 让审批人能回调

activity.logger.info("待审批草稿: %s", draft.title)

@activity.defn

async def revise_draft(draft: Draft, feedback: str) -> Draft:

body = await call_llm(

system="根据审阅意见修改草稿,保留来源标注。",

user=f"原文:\n{draft.body}\n\n意见:{feedback}",

)

return Draft(title=draft.title, body=body, sources=draft.sources)

@activity.defn

async def publish_report(draft: Draft) -> str:

这里必须幂等!重试时不能发出两份

建议用 workflow_id 做去重键,见"常见坑"第 2 条

return "https://example.com/reports/xxx"

```

第 4 步:写 Workflow,用 Signal 接人工审批

这是全文的核心。注意三个地方:Signal 怎么收审批、wait_condition 怎么等、重试策略怎么配。

```python

workflows.py

import asyncio

from datetime import timedelta

from temporalio import workflow

from temporalio.common import RetryPolicy

沙箱环境里,把带 IO 的模块放行,避免导入期报错

with workflow.unsafe.imports_passed_through():

from activities import (

ResearchBrief, Draft, DraftRequest, plan_queries, fetch_page,

write_draft, notify_reviewer, revise_draft, publish_report,

)

@workflow.defn

class ResearchAgentWorkflow:

def __init__(self) -> None:

self._decision: bool | None = None # None 表示还没人审批

self._feedback: str = ""

self._stage: str = "init"

---------- 外部交互入口 ----------

@workflow.signal

def submit_review(self, approved: bool, feedback: str = "") -> None:

"""审批人调用这个信号。先到也没关系,会被缓冲。"""

self._decision = approved

self._feedback = feedback

@workflow.query

def stage(self) -> str:

"""给前端 / 运维查进度用,不会改变状态"""

return self._stage

---------- 主流程 ----------

@workflow.run

async def run(self, brief: ResearchBrief) -> str:

self._stage = "plan"

queries = await workflow.execute_activity(

plan_queries,

brief,

start_to_close_timeout=timedelta(minutes=2),

retry_policy=RetryPolicy(maximum_attempts=5),

)

self._stage = "research"

并发抓取。Workflow 里用 asyncio.gather 是安全的

digests = await asyncio.gather(*[

workflow.execute_activity(

fetch_page,

q,

start_to_close_timeout=timedelta(seconds=90),

heartbeat_timeout=timedelta(seconds=30),

retry_policy=RetryPolicy(

initial_interval=timedelta(seconds=2),

backoff_coefficient=2.0,

maximum_interval=timedelta(minutes=2),

maximum_attempts=6,

non_retryable_error_types=["InvalidURL"],

),

)

for q in queries

])

self._stage = "draft"

draft = await workflow.execute_activity(

write_draft,

DraftRequest(brief=brief, digests=list(digests)),

start_to_close_timeout=timedelta(minutes=10),

retry_policy=RetryPolicy(maximum_attempts=3),

)

---------- 审批循环 ----------

rounds = 0

while True:

self._stage = "awaiting_approval"

await workflow.execute_activity(

notify_reviewer, draft,

start_to_close_timeout=timedelta(seconds=30),

)

try:

这一步可以挂三天,但 Worker 随时可以重启

await workflow.wait_condition(

lambda: self._decision is not None,

timeout=timedelta(days=3),

)

except asyncio.TimeoutError:

self._stage = "expired"

return "审批超时,草稿未发布"

approved = bool(self._decision)

feedback = self._feedback

self._decision = None # 重置,为下一轮做准备

self._feedback = ""

if approved:

break

rounds += 1

if rounds >= 3:

self._stage = "rejected"

return f"连续 {rounds} 轮未通过,已终止"

self._stage = f"revising_{rounds}"

draft = await workflow.execute_activity(

revise_draft, draft, feedback,

start_to_close_timeout=timedelta(minutes=10),

retry_policy=RetryPolicy(maximum_attempts=3),

)

self._stage = "publishing"

url = await workflow.execute_activity(

publish_report, draft,

start_to_close_timeout=timedelta(minutes=2),

retry_policy=RetryPolicy(maximum_attempts=5),

)

self._stage = "done"

return url

```

几个设计要点:

  • wait_condition 不是轮询。 它注册一个持久化定时器,不消耗 CPU,进程可以完全不存在。审批人三天后回话,Worker 起来接着跑。
  • Signal 早于 wait_condition 到达也没问题。 Temporal 会按事件顺序重放,_decision 在进入 wait_condition 前就已被赋值,条件立刻为真。
  • non_retryable_error_types 用来区分"值得重试"和"重试也没用"。 URL 本身格式错,重试六次还是错,直接失败更快。

第 5 步:启动 Worker

```python

worker.py

import asyncio

from temporalio.client import Client

from temporalio.worker import Worker

from activities import (

plan_queries, fetch_page, write_draft,

notify_reviewer, revise_draft, publish_report,

)

from workflows import ResearchAgentWorkflow

async def main() -> None:

client = await Client.connect("localhost:7233")

worker = Worker(

client,

task_queue="ai-agent",

workflows=[ResearchAgentWorkflow],

activities=[

plan_queries, fetch_page, write_draft,

notify_reviewer, revise_draft, publish_report,

],

)

print("Worker 已启动,Ctrl+C 停止")

await worker.run()

if __name__ == "__main__":

asyncio.run(main())

```

```bash

python worker.py

```

task_queue 这个名字是 Worker 和任务之间的对接口,两边写错一个字,任务就会永远卡在"待处理"。

第 6 步:发起一个任务并验证

```python

client.py

import asyncio

import uuid

from temporalio.client import Client

from activities import ResearchBrief

from workflows import ResearchAgentWorkflow

async def main() -> None:

client = await Client.connect("localhost:7233")

brief = ResearchBrief(topic="三家主流云厂商的 AI 推理定价", audience="采购负责人")

handle = await client.start_workflow(

ResearchAgentWorkflow.run,

brief,

id=f"research-{uuid.uuid4().hex[:8]}", # Workflow ID 就是幂等键

task_queue="ai-agent",

execution_timeout=None, # 长任务不要设太短的执行超时

)

print("已启动:", handle.id)

查进度(Query 是只读的,随时可调)

print("当前阶段:", await handle.query(ResearchAgentWorkflow.stage))

if __name__ == "__main__":

asyncio.run(main())

```

跑起来之后,去 UI 里看这个实例,你会看到事件历史一条条增长,停在 awaiting_approval

现在故意杀掉 Worker(Ctrl+C),等一分钟,再 python worker.py 重新启动。回到 UI,工作流依然是活的。这时发审批信号:

```bash

参数格式以官方文档为准

temporal workflow signal \

--workflow-id research-xxxxxxxx \

--name submit_review \

--input '{"approved": false, "feedback": "缺少定价对比表,请补充"}'

```

Worker 会立刻从 wait_condition 醒来,进入 revising_1,改完草稿再次送审。再发一次 {"approved": true},任务走到 publish_report 并返回链接。

你也可以在 client.py 里用代码发信号,前端审批按钮就是这么实现的:

```python

handle = client.get_workflow_handle("research-xxxxxxxx")

await handle.signal(ResearchAgentWorkflow.submit_review, True, "")

```

常见坑与排错

1. 在 Workflow 里写了非确定性代码。

典型症状:报 NonDeterminismError,或者重放时结果和第一次不一样。禁止项包括 datetime.now()random.random()uuid.uuid4()、直接发 HTTP 请求、直接读文件。对应替代:workflow.now()workflow.random()workflow.uuid4(),以及——把 IO 全部挪进 Activity。

2. Activity 不幂等,重试造成重复副作用。

Temporal 保证"至少执行一次",不保证"只执行一次"。所以 publish_report 重试时可能发两份。做法是:把 Workflow ID 或一个业务唯一键传进 Activity,写入前先查重(数据库唯一索引 / 幂等键),命中就直接返回上次结果。

3. start_to_close_timeout 设太短。

长上下文推理动辄几分钟,超时了 Temporal 会当成失败去重试,于是你会看到同一个 prompt 被反复调用、反复计费。给 LLM 类 Activity 留足余量,抓取类 Activity 配合心跳。

4. 把大对象塞进 Workflow 参数。

网页全文、base64 图片、几万字的向量,都别直接传。事件历史会膨胀,还可能撞上 payload 大小上限,而且每次重放都要反序列化一遍。正确做法:Activity 把大内容写到对象存储或数据库,只返回一个引用 ID 加摘要——本文的 PageDigest 就是这个思路。

5. 所有任务共用一个 Task Queue。

一个跑了二十分钟的 LLM Activity 会占住 Worker 的并发槽位,把轻量的通知 Activity 一起堵死。按任务粒度拆队列,给慢任务单独配 Worker 池和并发上限。

6. 改了 Workflow 代码,旧实例重放失败。

已经在跑的实例,其历史是按旧代码产生的。直接改逻辑会让重放对不上。生产环境要用 SDK 提供的版本化 / patching 机制,或者对存量实例用 Continue-As-New 开新的一轮。这块以官方文档为准。

7. 用信号做审批,但没重置状态。

本文在读取 _decision 之后立刻把它置回 None,否则第二轮 wait_condition 会瞬间返回上一次的结果,变成死循环。

排错三板斧:

  • 工作流卡住不动 → 先查 Worker 起没起、task_queue 名字对不对、Activity 有没有注册进 activities=[...]
  • ApplicationError: Activity task failed → 去 UI 点开那次 Activity 的失败原因和重试次数,maximum_attempts 用完就会冒到工作流层。
  • 信号发了没反应 → 用 handle.query(...) 看状态,确认工作流还活着、wait_condition 的条件真的能被这个信号满足。

下一步建议

骨架跑通之后,按这个顺序往上加:

1. 把 CLI 审批换成真实入口。 飞书 / 钉钉 / Slack 的交互卡片回调里调 handle.signal(...),审批人点一下就完成闭环。审批页面上再挂一个 handle.query(stage) 做进度条。

2. 看看 Update 能力。 相比 Signal 的单向通知,Update 能让审批人拿到返回值,更适合"提交表单—返回处理结果"的交互。具体语义以官方文档为准。

3. 用 Child Workflow 拆多 Agent。 一个"调研 Agent"可以派生出"定价子 Agent""舆情子 Agent",各自独立重试、独立超时,父工作流只负责汇总。

4. 用 Continue-As-New 做常驻 Agent。 历史事件攒到一定数量就开新一轮,保持实例轻量,适合"每天早上自动跑一遍"这类周期性任务。

5. 接可观测性。 Workflow 和 Activity 都带上 OpenTelemetry 追踪,把 LLM 的 token 消耗、耗时、失败率接到你的监控面板上。

6. 规划生产部署。 本地 start-dev 不能上生产,需要持久化的数据库后端;托管服务或自建集群的选型,以官方文档为准。

最后提醒一句心态上的转变:用 Temporal 写 Agent,你不再写"一个会自己跑的循环",而是写"一份可被反复重放的流程图"。想清楚每一步失败了该重试几次、超时算多久、要不要人来拍板,这些想清楚了,代码反而是最短的那部分。

AI 生成本文由 AI 基于公开信息自动生成,仅供参考。