从 demo 到生产隔着五道坎: 流式、持久化、幂等、限流降级、成本归因 — 每一道都是真金白银的教训
demo 里的 agent 跑在你眼前、内存里、一次成功就散场; 生产的 agent 跑在没人看着的凌晨三点: 用户会关掉页面(SSE 要流式)、发版会重启进程(状态要外置)、网络会抖(重试要幂等)、对手会刷你(限流要按用户)、模型会超时(降级链要备好)、老板会看账单(成本要归因)。一句话: demo 证明它能跑, 生产化证明它死了也能爬起来接着跑。
event: text data: "退款" # 文字 token 直接上屏 event: status data: "调用查询工具..." # 工具状态条 # 关键: nginx proxy_buffering off, 否则白屏依旧
steps: [plan, search, draft, review, submit] status: done/running/pending # 存 DB, 不存内存 # 重启 → 跳过 done, 从 pending 继续
claim(msg.id) # ON CONFLICT DO NOTHING → 重复忽略 idem_key # 业务级: 同单同额同键 # 没有 idempotency 的重试 = 随机重复扣款
pending(1h) → approve: 从挂起点继续
→ timeout: rejected_by_timeout
# agent 无常驻权限, 拿到的只是审批结果if spent > 0.50 or user_daily > 5.00: 熔断 # 熔断话术: "已花费 $0.52, 当前结论是..., 是否继续?"
classify → haiku-class # 简单活 deep-research → sonnet-class # 降级链: 旗舰 → 轻量 → 排队, 不裸奔 500
bucket(user_id): capacity 20, refill 5/s # Lua 原子操作; 超限明确告知, 不默默排队
redis.set("agent:{task_id}", state, ex=86400) # worker 崩了 → 另一个实例 load_state 继续
prompts/v14.md v15.md
# 错: 改 const 字符串直接上生产bucket = hash(user_id) % 100 version = "v15" if bucket < 10 else "v14" # ROLLOUT 配置热更新: 0 → 秒级全量回滚
GROUP BY user_id → 谁烧的钱
GROUP BY agent → 哪个 agent 烧的
# 优化永远打最大头L = λ × W = 500 × 20 = 10000 并发 # λ 翻倍 → L 翻倍; 只加服务器不加连接池 = 白加
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
五步长任务跑一半赶上发版。每步先落盘再执行, 重启后跳过已完成步骤。
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)
队列至少投递一次, 重复消息是常态。消息级 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} → 重试不重扣
审批单 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)
失控任务一晚烧 $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, 当前中间结论是..., 是否继续?"
分类任务也用旗舰模型, 月账单 $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") # 兜底
按 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
新 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 稳定), 体验不漂移
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 接着跑
老板要"昨天谁花了多少、成功率如何"。聚合 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;
# 错: result = await agent.run(...); return result # 对: StreamingResponse(gen(), media_type="text/event-stream")
# 错: 进程一挂从 step0 重跑 # 对: status=done 跳过, 从 pending 继续
# 错: TASKS = {} # → pod 重启清零 # 对: redis.set("agent:{tid}", json.dumps(state))
# 错: retry(refund(o, amt)) # → 重复退款 # 对: key=f"refund:{o}:{amt}" ON CONFLICT DO NOTHING
# 错: while True: run(...) # → 无限烧钱 # 对: if spent > 0.50: raise BudgetExceeded
# 错: bucket(ip) # → 一人触顶全员躺枪 # 对: bucket(user_id) # Lua 原子扣减
# 错: SYSTEM_PROMPT = "新写法" 直接部署 # 对: ROLLOUT v15=10% → 观测 → 全量
# 错: await call(model="旗舰") # → 超时即 500 # 对: except Overloaded: call("轻量") 兜底
# 错: 中间 3 分钟无输出 → 连接被 LB 掐断 # 对: 每 15s yield ": keepalive\n\n"(SSE 注释帧)
# 对: location /agent/ { proxy_buffering off; } # 对: 响应头 X-Accel-Buffering: no
# 错: log.debug(raw_messages) 全量落盘 # 对: 生产 WARN 级 + redact() 出口
# 错: 全局一个限流桶 # 对: bucket(tenant_id) + 独立预算与告警
# 错: classify 也用 sonnet-class # 对: 路由 haiku-class → 成本降 ~70%
# 错: pending 永久等待 # 对: 1h 过期拒绝 + args_sha 复核一致才执行
# 错: "这个月花了 $9k"(完) # 对: GROUP BY agent → researcher 62% → 定点优化
# 错: 只记 status, 不记"邮件已发" # 对: payload 里存 sent_msg_id, 恢复时查证再决定
# 错: 失败重新入队无限循环 # 对: max_retries=3 → DLQ + PagerDuty 告警
# 错: 一个队列 FIFO # → 批量任务堵死交互 # 对: interactive / batch 两队列, 交互优先
# 错: 回滚 = 重新部署整个服务 # 对: ROLLOUT 热更 10→0, 秒级切回 v14
# 错: "先上 10 个 worker 试试" # 对: L=λ×W=500×20=10000 并发 → 按需配+压测验证