核心思路:
如果是百万注册用户、10 万日活,要把 Agent 执行从 Web 进程中拆出来,采用持久化工作流 + 可水平扩展 Worker + 事件驱动架构。
API Gateway
↓
Run API / Auth / Quota
↓
Workflow Engine(Temporal / Cadence)
↓
Task Queue
↓
Agent Workers
├── LLM Worker
├── Tool Worker
├── RAG Worker
└── Delegation Worker
↓
PostgreSQL + Redis + Object Storage
↓
Event Bus / WebSocket Gateway / Audit
对于 Agent 场景,需要的中间件:
Temporal + PostgreSQL/MySQL + Kafka + Redis + S3/OSS
Temporal 负责:
- 父子 Workflow;
- 任务持久化;
- 重试;
- 超时;
- 崩溃恢复;
- 取消;
- 长时间暂停;
- 子任务等待和汇总。
Kafka 负责:
- 事件广播;
- 审计流;
- 指标流;
- 数据分析;
- WebSocket 事件分发。
不要让 Kafka 直接承担完整 Workflow 状态机,否则需要自己实现大量恢复、重试、定时器和幂等逻辑。
一、先正确理解规模
“百万用户、10 万 UV”不等于 10 万并发 Agent。
需要拆成几个指标:
DAU = 100,000
同时在线连接数 = 例如 5,000~20,000
Run 创建 QPS = 例如 50~500
LLM 调用并发 = 受模型供应商限制
工具调用并发 = 受数据库/API 限制
平均 Run 时长 = 5 秒、1 分钟还是 30 分钟
每日 Token 消耗 = 真正成本和容量约束
例如:
10 万 DAU
每人每天 3 个 Run
每天 30 万 Run
平均每个 Run 8 次 LLM 调用
每天 240 万次 LLM 调用
真正的瓶颈通常不是 CPU,而是:
- LLM API 并发限制;
- Token 预算;
- RAG 查询;
- 第三方工具限流;
- 数据库写入;
- 长连接推送;
- 单租户资源公平性。
所以调度系统必须同时做:
容量控制 + 优先级调度 + 租户配额 + LLM 限流 + 可靠执行
二、推荐的整体分层
1. API 层
API 只做:
鉴权
参数校验
配额检查
创建 Run
返回 run_id
不要在 HTTP 请求中直接执行 Agent。
POST /api/runs
→ 创建 run
→ 创建 workflow
→ 返回 202 + run_id
返回:
{
"run_id": "run_123",
"trace_id": "trace_456",
"status": "queued"
}
API 进程不负责:
- 调 LLM;
- 执行工具;
- 等待子 Agent;
- 重试;
- 长时间轮询。
这样 API 层才能无状态水平扩展。
2. Workflow 层
父 Agent 应该是一个持久化 Workflow:
AgentWorkflow
├── DecideActivity
├── ToolActivity
├── DelegateChildWorkflow
├── WaitChildResult
├── AggregateActivity
└── FinalizeActivity
伪代码:
@workflow.defn
class AgentWorkflow:
@workflow.run
async def run(self, task):
state = AgentState(task=task)
while not state.finished:
decision = await workflow.execute_activity(
decide,
state,
retry_policy=llm_retry_policy,
start_to_close_timeout=timedelta(minutes=2),
)
if decision.type == "final":
return await workflow.execute_activity(
finalize,
state,
)
if decision.type == "tool":
result = await workflow.execute_activity(
execute_tool,
decision,
retry_policy=tool_retry_policy,
)
state.append_tool_result(result)
if decision.type == "delegate":
child_result = await workflow.execute_child_workflow(
AgentWorkflow.run,
decision.child_task,
id=decision.child_run_id,
)
state.append_child_result(child_result)Workflow 的关键特性:
- Worker 挂掉后可以从历史状态恢复;
- 父 Workflow 可以等待子 Workflow;
- 等待期间不占用线程;
- 可以暂停数小时或数天;
- 可以接受取消信号;
- 可以查询当前进度;
- 可以继续执行未完成步骤。
三、如何解决跨进程执行
错误方式
不要这样:
background_tasks.add_task(run_agent)
因为:
- 进程重启后任务丢失;
- 多实例之间无法协调;
- 任务可能重复执行;
- 无法恢复执行位置;
- 没有可靠重试;
- 任务状态和内存状态不一致。
正确方式
API 创建持久化 Run,然后提交 Workflow:
POST /api/runs
↓
DB 写入 run(status=queued)
↓
Workflow Engine.start(run_id)
↓
Task Queue
↓
Worker 执行
Worker 可以部署多个:
agent-worker-1
agent-worker-2
agent-worker-3
...
agent-worker-N
不同能力使用不同队列:
agent-default
agent-llm
agent-rag
agent-tool
agent-delegation
agent-high-priority
例如:
LLM 任务 → llm-queue
PDF 解析 → document-queue
向量检索 → rag-queue
邮件发送 → side-effect-queue
这样可以单独扩容。
四、任务重试
不同错误必须使用不同策略,不能所有异常都重试。
1. 可重试错误
网络超时
HTTP 429
HTTP 500/502/503
临时数据库连接失败
Worker 暂时不可用
策略:
指数退避 + 抖动 + 最大次数
例如:
第 1 次:1 秒
第 2 次:2 秒
第 3 次:4 秒
第 4 次:8 秒
最大:30 秒
2. 不可重试错误
参数校验失败
Agent 不存在
跨租户委派
循环委派
权限不足
余额不足
内容安全拒绝
3. LLM 特殊策略
LLM 调用还要考虑:
单租户并发限制
模型级并发限制
Token 速率限制
请求成本
上下文长度
建议使用:
Retry Policy
+ Token Bucket
+ Per-tenant Semaphore
+ Model-level Rate Limiter
例如:
tenant-A:
最大并行 LLM 调用 = 20
每分钟 Token = 2,000,000
tenant-B:
最大并行 LLM 调用 = 5
每分钟 Token = 200,000
五、Worker 崩溃恢复
Worker 不应该自己决定“任务完成”。
正确流程:
Workflow Engine 持有任务历史
Worker 执行 Activity
Activity 成功 → 回报结果
Activity 崩溃 → Engine 超时后重试
Worker 进程死亡 → 任务被其他 Worker 接管
Activity 必须有:
start_to_close_timeout
heartbeat_timeout
retry_policy
对于长任务,Worker 发送 heartbeat:
while processing:
heartbeat(progress=0.5)
do_next_step()
如果 heartbeat 停止:
Workflow Engine 判定 Worker 失联
→ Activity 超时
→ 重新调度到其他 Worker
这比自己在 Redis 里写:
processing → timeout → retry
可靠很多。
六、分布式锁
分布式锁不应该用来包住整个 Agent Run。
错误:
拿 Redis 锁
→ 执行 10 分钟 Agent
→ 最后释放锁
这样会造成:
- 锁过期;
- Worker 崩溃后锁残留;
- 长任务阻塞;
- 锁续期复杂;
- 吞吐下降。
正确做法是:
1. 用数据库唯一约束做幂等
UNIQUE(tenant_id, idempotency_key)
2. 用 Workflow ID 防止重复启动
workflow_id = tenant_id + ":" + idempotency_key
3. 只在极短临界区使用锁
例如:
扣减配额
分配唯一资源
更新 Agent 状态
锁的持有时间尽量小于几十毫秒。
4. 外部副作用使用幂等键
例如发送邮件:
idempotency_key = run_id + ":send-email"
数据库表:
CREATE UNIQUE INDEX idx_effect
ON side_effects(run_id, effect_key);
七、幂等设计
分布式系统一定会出现重复执行,所以每个层次都要幂等。
1. 创建 Run 幂等
POST /api/runs
Idempotency-Key: abc-123
数据库:
UNIQUE(tenant_id, idempotency_key)
2. Activity 幂等
activity_id = run_id + ":step:" + step_number
3. 工具调用幂等
tool_call_id = run_id + ":" + sequence
4. 子 Agent 幂等
child_workflow_id =
parent_run_id + ":" + delegation_id
如果父 Workflow 因网络抖动重复发起委派,系统应返回已有子 Run,而不是创建两个子任务。
5. 副作用幂等
高风险操作必须使用:
审批状态
幂等键
操作日志
补偿机制
例如:
发送短信
创建订单
生成合同
修改外部系统
不能简单依赖“重试不会重复”。
八、父任务暂停和恢复
父任务等待子 Agent 时,不应该占用线程。
状态机:
RUNNING
↓
WAITING_FOR_DELEGATION
↓
CHILD_RUNNING
↓
CHILD_SUCCEEDED
↓
RUNNING
↓
SUCCEEDED
父 Workflow:
child_id = await start_child_workflow(...)
await workflow.wait_condition(
lambda: state.child_results.get(child_id) is not None
)
child_result = state.child_results[child_id]
state.messages.append(child_result)
如果使用 Temporal:
ChildWorkflowExecutionStub child =
Workflow.newChildWorkflowStub(
AgentWorkflow.class
);
Promise<AgentResult> result =
Async.function(child::run, childTask);
AgentResult childResult = result.get();
等待时:
- 不占用 JVM 线程;
- Worker 可被回收;
- Workflow 状态由引擎保存;
- 重启后继续等待。
如果子任务可能等待人工审批:
WAITING_FOR_HUMAN_APPROVAL
可以暂停几天,收到 Signal 后继续。
九、多个子任务汇总
父 Agent 可能同时委派:
设计 Agent
安全 Agent
排期 Agent
预算 Agent
可以并行执行:
children = await asyncio.gather(
delegate("design", design_task),
delegate("safety", safety_task),
delegate("schedule", schedule_task),
delegate("budget", budget_task),
)
但生产上建议使用持久化子 Workflow:
parent workflow
├── child design
├── child safety
├── child schedule
└── child budget
父任务维护:
state.children = {
"design": "run_design",
"safety": "run_safety",
"schedule": "run_schedule",
"budget": "run_budget",
}
汇总时:
results = [
child.output
for child in completed_children
]
final = await llm.complete(
parent_context=state.messages,
child_results=results,
)
需要定义部分失败策略:
全部成功 → 正常汇总
部分成功 → 汇总成功结果,标记缺失结果
关键子任务失败 → 父任务失败
非关键子任务超时 → 降级继续
例如:
安全审核失败 → 不允许最终发布
视觉设计失败 → 可以使用默认模板
天气查询失败 → 输出“天气信息暂不可用”
十、超时和取消传播
1. 超时层级
HTTP 请求超时:30 秒
单次 LLM 调用:60 秒
单个工具:15 秒
单个子 Agent:5 分钟
整个父 Run:30 分钟
每层都要有 deadline:
deadline = min(
parent_deadline,
tool_deadline,
tenant_deadline
)
2. 取消流程
POST /api/runs/{id}/cancel
处理:
Run 设置 cancel_requested=true
↓
发送 Workflow cancellation
↓
父 Workflow 收到取消
↓
取消等待中的子 Workflow
↓
取消可取消的 Activity
↓
写入 run.cancelled
↓
写入 child run.cancelled
外部 HTTP 请求已经发出时,通常无法真正撤回:
等待当前调用自然结束
或在客户端 timeout
结果丢弃
Run 最终标记 cancelled
工具代码要定期检查:
if cancellation_requested():
raise CancelledError()
十一、长时间运行任务
长任务不要依赖:
一个 HTTP 请求
一个线程
一个进程内变量
一个 Redis key
应该拆成多个 Activity:
解析文档
→ 分块
→ 向量化
→ 写入索引
→ 检索
→ 子 Agent 分析
→ 人工审批
→ 最终生成
每一步都有:
状态
输入
输出
重试次数
耗时
错误
长任务状态:
RUNNING
WAITING_FOR_CHILD
WAITING_FOR_APPROVAL
PAUSED
RESUMING
SUCCEEDED
前端不要保持一个长 HTTP 请求,而是:
POST /api/runs → 202
GET /api/runs/{id}
GET /api/runs/{id}/events?after_sequence=N
WebSocket / SSE
十二、任务队列设计
建议按任务特征拆分队列:
agent.interactive
agent.batch
agent.llm
agent.rag
agent.document
agent.side_effect
agent.delegation
优先级:
P0:用户正在等待的交互任务
P1:普通任务
P2:批处理
P3:离线评测
调度时不能只使用全局 FIFO,还要防止大租户霸占队列。
可以采用:
Weighted Fair Queue
+ Per-tenant concurrency limit
+ Priority queue
+ Token budget
例如:
每租户最大并发 Run = 20
每用户最大并发 Run = 5
每 Agent 最大并发 Run = 10
全局最大 LLM 并发 = 500
令牌桶:
租户预算
↓
队列准入
↓
LLM 并发信号量
↓
模型调用
十三、审计和追踪
每个请求必须有:
request_id
run_id
trace_id
parent_run_id
root_run_id
tenant_id
user_id
agent_id
事件至少包括:
run.queued
run.started
run.paused
run.resumed
run.cancel.requested
run.cancelled
run.succeeded
run.failed
llm.invoked
llm.completed
llm.failed
tool.invoked
tool.completed
tool.failed
delegation.started
delegation.ended
delegation.failed
rag.retrieved
approval.requested
approval.completed
事件表:
CREATE TABLE run_events (
id UUID PRIMARY KEY,
tenant_id VARCHAR(64) NOT NULL,
trace_id VARCHAR(128) NOT NULL,
run_id VARCHAR(128) NOT NULL,
parent_run_id VARCHAR(128),
sequence BIGINT NOT NULL,
event_type VARCHAR(64) NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMP NOT NULL,
UNIQUE(run_id, sequence)
);
日志和事件要区分:
日志:
调试信息,短期保留,可能采样
事件:
审计事实,长期保存,不可随意修改
指标:
聚合数据,用于监控和报警
LLM 原始请求和返回结果不建议无限期存储,应该:
- 默认脱敏;
- 限时保存;
- 加密;
- 按租户权限访问;
- 生产环境关闭完整 payload;
- 只存安全摘要和 token/latency。
十四、数据库设计建议
PostgreSQL
存:
tenants
users
agents
runs
run_events
agent_invocations
idempotency_records
quotas
approvals
Redis
适合:
短期缓存
限流计数
在线连接
WebSocket Pub/Sub
热状态
租户并发信号量
不要把 Redis 作为唯一任务事实来源。
Kafka
适合:
事件流
审计流
指标流
异步通知
数据分析
S3/OSS
存:
PDF/DOCX 原文件
大模型长输出
大对象上下文
导出文档
事件归档
数据库只存 URI 和摘要,不要把几百 MB 文档直接放在 Run 表。
十五、扩容方式
API 层
无状态水平扩展:
api-1
api-2
api-3
...
通过 Load Balancer 分发。
Worker 层
按队列独立扩展:
interactive-worker: 20
llm-worker: 50
rag-worker: 20
document-worker: 10
batch-worker: 5
WebSocket 层
WebSocket Gateway 独立部署:
client
↓
WebSocket Gateway
↓
Redis Pub/Sub / Kafka
不能让每个 Agent Worker 直接维护用户长连接。
自动扩缩容
Kubernetes HPA/KEDA 根据:
队列长度
任务等待时间
LLM 并发
CPU
内存
错误率
例如:
queue_wait_p95 > 10s
→ 增加 Worker
queue_wait_p95 < 1s 且持续 10 分钟
→ 减少 Worker
十六、一个实际的父子 Workflow 方案
POST /runs
↓
创建 parent_run
↓
启动 ParentWorkflow
↓
ParentWorkflow 调用 LLM
↓
LLM 返回 collaborate_with_agent
↓
校验目标 Agent
↓
创建 ChildWorkflow
↓
ParentWorkflow 状态变为 WAITING_FOR_DELEGATION
↓
ChildWorkflow 投递到 child-agent 队列
↓
Child Worker 执行
↓
ChildWorkflow 写入 output
↓
ParentWorkflow 收到 child result
↓
父消息历史追加 tool_result
↓
ParentWorkflow 再次调用 LLM
↓
最终输出
数据关系:
ParentWorkflow:
run_id = P
trace_id = T
status = WAITING_FOR_DELEGATION
ChildWorkflow:
run_id = C
parent_run_id = P
trace_id = T
status = RUNNING
最终:
P.output = 父 Agent 汇总结果
C.output = 子 Agent 原始结果
父 Agent 不应该直接读取子 Agent 的完整内部消息,只读取:
{
"child_run_id": "C",
"status": "SUCCEEDED",
"output": "...",
"citations": [],
"error_code": null
}
这样可以控制上下文大小,也避免把子 Agent 的内部 prompt、思维过程泄露给父 Agent。
十七、推荐落地顺序
不要一开始就上百万用户架构。建议分阶段:
阶段 1:可靠单机
FastAPI
PostgreSQL
Redis
持久化 Run
持久化 RunEvent
幂等
重试
取消
阶段 2:多 Worker
API 与 Worker 分离
Redis Streams / RabbitMQ
Worker 崩溃恢复
队列优先级
租户限流
阶段 3:持久化 Workflow
Temporal
父子 Workflow
长任务暂停
子任务汇总
超时传播
人工审批
阶段 4:大规模化
Kafka
Kubernetes
KEDA
独立 WebSocket Gateway
多区域部署
数据分片
租户配额
成本治理
最终建议
对于百万用户、10 万 UV,推荐架构是:
Spring Boot / FastAPI API
+
Temporal Workflow
+
PostgreSQL
+
Kafka
+
Redis
+
Kubernetes Workers
+
OpenTelemetry
核心原则:
API 无状态
Workflow 有状态
Worker 可崩溃
任务可重试
操作有幂等
父子可恢复
事件可审计
租户有限流
副作用可补偿
如果只需要实现当前项目的生产演进版本,不必一次上所有组件。最有价值的第一步是:
把 BackgroundTasks 替换为持久化任务队列,
再把父子协作改成持久化 Workflow。
这样就能真正解决:
- 跨进程执行;
- Worker 崩溃恢复;
- 父任务暂停/恢复;
- 子任务结果汇总;
- 取消传播;
- 长时间运行;
- 可靠审计。