这篇能做出什么
先看成果。你要做的是一个"竞品调研 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 Queue | Worker 和任务的匹配队列。不同任务用不同队列,别互相堵 |
| 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,你不再写"一个会自己跑的循环",而是写"一份可被反复重放的流程图"。想清楚每一步失败了该重试几次、超时算多久、要不要人来拍板,这些想清楚了,代码反而是最短的那部分。
