Agent · 生产化(Production)

从 demo 到生产隔着五道坎: 流式、持久化、幂等、限流降级、成本归因 — 每一道都是真金白银的教训

生产 agent 骨架: 网关 → 队列 → 无状态 Durable Worker → 模型/工具层; 人在环与观测挂在主干上 API 网关 鉴权 · 用户级限流 SSE 透传不缓冲 任务队列 异步削峰 · 重试 优先级 + 死信 DLQ Durable Worker checkpoint · 断点恢复 无状态: 任意实例可接管 模型 / 工具层 旗舰/轻量分级路由 写操作幂等 + 硬上限 事故现场 — 无 checkpoint 的滚动发布 发版重启瞬间, 200 个进行中任务全部蒸发 代价: 双倍 token 重跑 + 用户重等 10 分钟 HITL 审批分支 pending → approve / 超时拒绝 批准后从挂起点继续 全链路观测 成本/延迟/错误率 三告警 按用户与 agent 下钻 稳态运营六件套 Steady-State Ops 流式 SSE 首字 0.6s; nginx 记得 proxy_buffering off 预算熔断 任务/用户双维度 $ 上限, 超限立停并汇报 降级链 旗舰超时 → 轻量兜底; 限流带优先级排队 成本大盘 按 用户/agent/task 三维归因出账 幂等与重试 消息 claim + 业务幂等键, 重试永远安全 灰度与回滚 prompt 即代码: 版本化 + 10% 灰度 + 秒回滚 心跳: 长任务每 15s ping 一次, 防止被 LB 判死掐断连接; SSE 必须关代理缓冲 容量: Little's Law L = λ × W — 500 任务/s × 平均 20s = 10000 并发容量, worker 数按这个配 多租户: 每租户独立 token 桶配额; 互相挤占即告警, 谁家失控熔断谁家 事故现场: prompt 改动全量直发 → 全站质量回退 3 小时才发现 — 版本化 + 灰度 + eval 门禁是保险丝 发布检查单: eval 通过 → 灰度 10% → 观测 30 分钟 → 全量; 回滚永远比修复快

骨架无状态(机制视角)

  • • worker 不存状态: 任务状态全在 DB/Redis
  • • 步骤级 checkpoint: 崩溃从断点续跑
  • • 队列削峰: 流量尖刺不再直接打穿模型层

安全带(行为视角)

  • • 消息 at-least-once → 幂等双保险
  • • 预算熔断: 任务/用户双维度上限
  • • 降级链: 旗舰挂了轻量顶上, 不裸奔报错

账与闸(生产价值)

  • • 成本三维归因: 用户/agent/task 各一本账
  • • prompt 版本化 + 10% 灰度 + eval 门禁
  • • 容量不拍脑袋: Little's Law 算并发

💡 一句话理解

demo 里的 agent 跑在你眼前、内存里、一次成功就散场; 生产的 agent 跑在没人看着的凌晨三点: 用户会关掉页面(SSE 要流式)、发版会重启进程(状态要外置)、网络会抖(重试要幂等)、对手会刷你(限流要按用户)、模型会超时(降级链要备好)、老板会看账单(成本要归因)。一句话: demo 证明它能跑, 生产化证明它死了也能爬起来接着跑。

🧠 必知必会 必考 & 必会

Streaming SSE
agent 一轮 10 秒起步, 不流式就是白屏 10 秒。SSE 把 token 与状态事件分开推, 首字延迟从秒级降到亚秒。
event: text   data: "退款"      # 文字 token 直接上屏
event: status data: "调用查询工具..."  # 工具状态条
# 关键: nginx proxy_buffering off, 否则白屏依旧
Durable Execution
长任务拆成步骤状态机, 每步先落盘再执行。崩溃重启后从第一个未完成步骤继续——Temporal 的核心思想。
steps: [plan, search, draft, review, submit]
status: done/running/pending   # 存 DB, 不存内存
# 重启 → 跳过 done, 从 pending 继续
At-least-once 与幂等
队列的消息语义是"至少一次": 重复投递是常态不是异常。消息 claim + 业务幂等键双保险, 重试才安全。
claim(msg.id)   # ON CONFLICT DO NOTHING → 重复忽略
idem_key        # 业务级: 同单同额同键
# 没有 idempotency 的重试 = 随机重复扣款
HITL 审批流
高危动作挂起生成审批单, 人批准后由调度器继续。审批有超时, 默认动作是拒绝。
pending(1h) → approve: 从挂起点继续
           → timeout: rejected_by_timeout
# agent 无常驻权限, 拿到的只是审批结果
预算熔断
任务级与用户级双闸: 单任务 $0.5、单用户日 $5, 超限立停并把中间结论汇报给用户, 而不是悄悄烧穿。
if spent > 0.50 or user_daily > 5.00: 熔断
# 熔断话术: "已花费 $0.52, 当前结论是..., 是否继续?"
模型路由与降级
分类/抽取/路由类任务用轻量模型(成本 ~1/15), 复杂推理才用旗舰; 旗舰超时/过载自动降级兜底。
classify → haiku-class   # 简单活
deep-research → sonnet-class
# 降级链: 旗舰 → 轻量 → 排队, 不裸奔 500
用户级限流
按用户而非 IP 限流(NAT 后一个 IP 一群人)。Redis 令牌桶原子扣减, 超限 429 + Retry-After。
bucket(user_id): capacity 20, refill 5/s
# Lua 原子操作; 超限明确告知, 不默默排队
会话状态外置
worker 无状态: messages/step/cost 全在 Redis/DB。任意实例可接管任意任务, 发布/扩缩容不丢任务。
redis.set("agent:{task_id}", state, ex=86400)
# worker 崩了 → 另一个实例 load_state 继续
prompt 版本化
prompt 是生产代码: 进 git、带版本号、改动走评审; 每个版本 eval 通过才可发布。
prompts/v14.md  v15.md
# 错: 改 const 字符串直接上生产
灰度与回滚
新版本 prompt/模型先 10% 流量, 观测指标 30 分钟, 异常秒级切回。回滚速度决定事故时长。
bucket = hash(user_id) % 100
version = "v15" if bucket < 10 else "v14"
# ROLLOUT 配置热更新: 0 → 秒级全量回滚
成本归因
每个 run 记 user/agent/task 三维, SQL 聚合出账。没有归因的成本优化都是瞎砍。
GROUP BY user_id → 谁烧的钱
GROUP BY agent   → 哪个 agent 烧的
# 优化永远打最大头
容量规划
Little's Law: 并发数 = 到达率 × 平均耗时。500 任务/s × 20s = 10000 并发, worker 与连接池照此配, 不是拍脑袋。
L = λ × W = 500 × 20 = 10000 并发
# λ 翻倍 → L 翻倍; 只加服务器不加连接池 = 白加

🏭 生产实战 real world

场景 1 · SSE 流式接口: 首字 0.6 秒

FastAPI + Anthropic 流式, 文字 token 上屏、工具调用转状态事件, nginx 关缓冲。

@app.post("/agent/chat")
async def chat(req: ChatReq):
    async def gen():
        async with client.messages.stream(...) as stream:
            async for event in stream:
                if event.type == "content_block_delta":
                    yield sse({"type": "text", "data": event.delta.text})
        yield sse({"type": "done"})
    return StreamingResponse(gen(), media_type="text/event-stream")
# nginx: proxy_buffering off; 首字 8s → 0.6s

场景 2 · 步骤级 checkpoint: 断点续跑

五步长任务跑一半赶上发版。每步先落盘再执行, 重启后跳过已完成步骤。

PLAN = ["plan", "search", "draft", "review", "submit"]

async def run_task(task_id: str):
    for seq, step in enumerate(PLAN):
        row = await db.fetchrow(
            "SELECT status FROM task_steps WHERE task_id=%s AND seq=%s",
            (task_id, seq))
        if row and row["status"] == "done":
            continue                    # 断点续跑: 完成的跳过
        await set_status(task_id, seq, "running")
        result = await execute(step)     # 每个步骤都幂等
        await set_status(task_id, seq, "done", result)

场景 3 · 幂等双保险: 重试永远安全

队列至少投递一次, 重复消息是常态。消息级 claim + 业务幂等键, 两层都不能省。

async def handle(msg):
    claimed = await db.execute(
        "INSERT INTO msg_claims(id) VALUES (%s) ON CONFLICT DO NOTHING",
        (msg.id,))
    if claimed == 0:
        return                          # 重复投递, 直接忽略
    await process(msg)                  # 工具内部还有业务幂等键
# refund 工具: key=refund:{order}:{amount} → 重试不重扣

场景 4 · HITL 审批 worker: 超时默认拒绝

审批单 1 小时无人处理自动过期拒绝, 宁可让用户再点一次, 不可默认放行。

async def approval_worker():
    while True:
        rows = await db.fetch(
            "SELECT id, task_id FROM approvals"
            " WHERE status='pending'"
            "   AND created_at < now() - interval '1 hour'")
        for r in rows:
            await db.execute(
                "UPDATE approvals SET status='expired' WHERE id=%s", r["id"])
            await notify_task(r["task_id"], "rejected_by_timeout")
        await asyncio.sleep(60)

场景 5 · 预算熔断: 双维度上限

失控任务一晚烧 $400 的教训。任务级+用户日级双闸, 熔断时带中间结论汇报。

class BudgetExceeded(Exception): ...

async def check_budget(user_id: str, task_spent: float):
    daily = float(await redis.get(f"spend:{user_id}:{today()}") or 0)
    if task_spent > 0.50 or daily > 5.00:
        raise BudgetExceeded(f"task=${task_spent:.2f} daily=${daily:.2f}")
# 熔断话术: "已花费 $0.52, 当前中间结论是..., 是否继续?"

场景 6 · 模型路由与降级链

分类任务也用旗舰模型, 月账单 $12k。按任务类型路由 + 超时降级, 成本降 70%。

def route_model(task) -> str:
    if task.kind in {"classify", "extract", "route"}:
        return "haiku-class"            # 成本 ≈ 旗舰 1/15
    return "sonnet-class"

async def call_with_fallback(messages, model):
    try:
        return await call(messages, model=model, timeout=30)
    except (TimeoutError, Overloaded):
        return await call(messages, model="haiku-class")  # 兜底

场景 7 · 用户级限流: Redis 令牌桶

按 IP 限流把整个公司挡在门外。改按用户, Lua 原子扣减, 超限明确 429。

LUA = """
local tokens = redis.call('get', KEYS[1]) or ARGV[1]
tokens = math.min(tonumber(tokens) + (ARGV[3]-last), ARGV[1])
if tokens >= 1 then ... return 1 else return 0 end
"""  # 容量 20, 速率 5/s, 原子执行

async def allow(user_id: str) -> bool:
    ok = await redis.eval(LUA, keys=[f"rl:{user_id}"],
                          args=[20, 5, int(time.time())])
    return ok == 1   # False → 429 + Retry-After

场景 8 · prompt 灰度: 10% 起步秒回滚

新 prompt 不再全量直发。哈希分桶 10% 灰度, 配置热更新, 回滚零发布成本。

ROLLOUT = {"v15": 10}     # 热更新配置: 10 → 0 即秒级回滚

def pick_version(user_id: str) -> str:
    bucket = hash(user_id) % 100
    return "v15" if bucket < ROLLOUT.get("v15", 0) else "v14"
# 纪律: 每个版本 eval 通过才允许加灰度比例
# 同一用户永远同一版本(hash 稳定), 体验不漂移

场景 9 · 状态外置: worker 无状态化

K8s 扩缩容随机杀 pod。任务状态全在 Redis, 任意实例接管任意任务。

async def save_state(task_id: str, state: dict):
    await redis.set(f"agent:{task_id}", json.dumps(state), ex=86400)

async def load_state(task_id: str) -> dict:
    raw = await redis.get(f"agent:{task_id}")
    return json.loads(raw) if raw else new_state(task_id)
# 发布/扩缩容/OOM 重启: 任务都从 checkpoint 接着跑

场景 10 · 运营大盘: 一条 SQL 出账

老板要"昨天谁花了多少、成功率如何"。聚合 SQL 直出 Top 20。

-- 近 24h 按用户聚合: 成本 / 成功率 / P99 延迟
SELECT user_id,
       count(*)                                            AS tasks,
       round(sum(cost_usd)::numeric, 2)                    AS cost,
       round(100.0 * avg((status = 'done')::int), 1)       AS success_pct,
       percentile_cont(0.99) WITHIN GROUP
         (ORDER BY duration_s)                             AS p99_s
FROM agent_runs
WHERE started_at > now() - interval '24 hours'
GROUP BY user_id
ORDER BY cost DESC
LIMIT 20;

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 无流式白屏 — 症状: 用户等 15 秒白屏后流失. 原因: 整轮完成才返回. 正解: SSE 分事件推送。
# 错: result = await agent.run(...); return result
# 对: StreamingResponse(gen(), media_type="text/event-stream")
坑 2 · 崩溃从头重跑 — 症状: 发版后任务重新执行, 重复扣款. 原因: 无 checkpoint. 正解: 步骤状态机 + 断点续跑。
# 错: 进程一挂从 step0 重跑
# 对: status=done 跳过, 从 pending 继续
坑 3 · 状态存进程内存 — 症状: 扩缩容/重启任务蒸发. 原因: messages/step 在 dict 里. 正解: Redis/DB 外置。
# 错: TASKS = {}                    # → pod 重启清零
# 对: redis.set("agent:{tid}", json.dumps(state))
坑 4 · 重试非幂等工具 — 症状: 网络抖动一次, 用户被退款两次. 原因: 重试语义遇上无幂等的写操作. 正解: 业务幂等键 + 唯一约束。
# 错: retry(refund(o, amt))        # → 重复退款
# 对: key=f"refund:{o}:{amt}" ON CONFLICT DO NOTHING
坑 5 · 无预算上限 — 症状: 一夜烧穿 $400, 早上才发现. 原因: 循环无费用闸门. 正解: 任务/用户双维度熔断。
# 错: while True: run(...)       # → 无限烧钱
# 对: if spent > 0.50: raise BudgetExceeded
坑 6 · 按 IP 限流 — 症状: 整个公司用户一起 429. 原因: NAT 后共享出口 IP. 正解: 按用户/租户令牌桶。
# 错: bucket(ip)                 # → 一人触顶全员躺枪
# 对: bucket(user_id)             # Lua 原子扣减
坑 7 · prompt 全量直发 — 症状: 一次改动全站质量回退 3 小时. 原因: 无灰度无版本. 正解: 版本化 + 10% 灰度 + eval 门禁。
# 错: SYSTEM_PROMPT = "新写法" 直接部署
# 对: ROLLOUT v15=10% → 观测 → 全量
坑 8 · 模型超时无降级 — 症状: 上游过载时全站 500. 原因: 单模型无兜底. 正解: 降级链 + 排队优先级。
# 错: await call(model="旗舰")   # → 超时即 500
# 对: except Overloaded: call("轻量") 兜底
坑 9 · 长任务无心跳 — 症状: 跑 5 分钟的 SSE 连接总在 60s 被掐. 原因: LB 只认数据流, 静默即判死. 正解: 每 15s 发心跳注释帧。
# 错: 中间 3 分钟无输出 → 连接被 LB 掐断
# 对: 每 15s yield ": keepalive\n\n"(SSE 注释帧)
坑 10 · SSE 被代理缓冲 — 症状: 代码里明明流式, 用户还是白屏到底. 原因: nginx 默认 proxy_buffering on. 正解: 关缓冲 + X-Accel-Buffering: no。
# 对: location /agent/ { proxy_buffering off; }
# 对: 响应头 X-Accel-Buffering: no
坑 11 · 敏感日志落盘 — 症状: 审计发现日志里有用户对话全文+密钥. 原因: debug 日志忘关上生产. 正解: 分级日志 + redact + 日志留存策略。
# 错: log.debug(raw_messages) 全量落盘
# 对: 生产 WARN 级 + redact() 出口
坑 12 · 多租户无配额 — 症状: 大客户一次导出把共享配额吃光, 其他客户全 429. 原因: 无租户级配额隔离. 正解: 每租户独立桶 + 抢占告警。
# 错: 全局一个限流桶
# 对: bucket(tenant_id) + 独立预算与告警
坑 13 · 全量旗舰模型 — 症状: 月账单 $12k, 80% 是分类任务烧的. 原因: 不分任务全用旗舰. 正解: 任务路由 + 降级链。
# 错: classify 也用 sonnet-class
# 对: 路由 haiku-class → 成本降 ~70%
坑 14 · 审批无超时 — 症状: 三天前的审批单被误批, 上下文早已变化. 原因: 审批单永不过期. 正解: expires_in + 执行前复核参数 hash。
# 错: pending 永久等待
# 对: 1h 过期拒绝 + args_sha 复核一致才执行
坑 15 · 成本不归因 — 症状: 账单翻倍, 不知道砍谁. 原因: 只有总量. 正解: user/agent/task 三维出账。
# 错: "这个月花了 $9k"(完)
# 对: GROUP BY agent → researcher 62% → 定点优化
坑 16 · 恢复后重放副作用 — 症状: 断点续跑后用户收到两封邮件. 原因: 恢复时不确定哪些副作用已发生. 正解: 每步记录副作用凭证, 恢复时核对。
# 错: 只记 status, 不记"邮件已发"
# 对: payload 里存 sent_msg_id, 恢复时查证再决定
坑 17 · 队列无死信 — 症状: 一条毒消息重试无限次, worker 全在空转. 原因: 无 DLQ 与最大重试. 正解: 重试 3 次进 DLQ + 告警。
# 错: 失败重新入队无限循环
# 对: max_retries=3 → DLQ + PagerDuty 告警
坑 18 · 任务无优先级 — 症状: 高峰期实时问答排在批量导出后面等 2 分钟. 原因: 队列 FIFO 一把抓. 正解: 分优先级队列, 交互>批量。
# 错: 一个队列 FIFO          # → 批量任务堵死交互
# 对: interactive / batch 两队列, 交互优先
坑 19 · 无金丝雀回滚 — 症状: 出事只能回滚整个服务, 恢复 20 分钟. 原因: prompt/模型配置不可独立回滚. 正解: 配置热更新 + 一键切版。
# 错: 回滚 = 重新部署整个服务
# 对: ROLLOUT 热更 10→0, 秒级切回 v14
坑 20 · 容量拍脑袋 — 症状: 促销流量一来, worker 池和连接池同时耗尽. 原因: 没算过并发需求. 正解: Little's Law 定容量, 压测验证。
# 错: "先上 10 个 worker 试试"
# 对: L=λ×W=500×20=10000 并发 → 按需配+压测验证