044AI Agent阅读记录

从零开始 AI Agent 实战(十五):机器人回群、可观测性与部署

把 Agent 接回钉钉群完成闭环,并建立 Trace、成本估算、限流降级与 Docker Compose 部署。

乌漆嘛黑和 Ahri
第 044 期

机器人回群、可观测性与部署

「今天 Agent 很慢」

这是上线后你会收到的第一条反馈,而且它不带任何可用信息。没有 trace 时,能做的只有猜:是检索慢了?重排慢了?模型首 token 慢了?数据库连接池耗尽了?还是那个用户的网络问题?

猜的代价很具体。团队通常会先给向量检索加缓存——听起来最像瓶颈。做完发现没用。

Preview
一次 2.8 秒请求的 span 树

这条 trace 一眼给出答案:检索三路加起来 180ms,而 llm.answer 一个人占了 1.7 秒。想让它变快,要动的是输出长度和模型选择,不是向量索引。那 180ms 就算降到 0,用户也感觉不到。

这一篇的顺序是:先让系统可观测,再基于观测做限流和降级,最后部署。反过来做的话,你会在没有数据的情况下做优化决策。

Trace:span 划分与属性

OpenTelemetry 的 FastAPI 自动埋点能给你 HTTP 层的 span,但 Agent 的关键路径全在业务代码里,要手动划:

tracer = trace.get_tracer('ai_agent_guide')

async def handle_chat(message: str, user: User) -> ChatResult:
    with tracer.start_as_current_span('chat.request') as span:
        span.set_attribute('conversation.id', conversation_id)
        span.set_attribute('user.role', user.role)          # 角色可以,用户名不行
        with tracer.start_as_current_span('retrieval') as rspan:
            rspan.set_attribute('retrieval.rewritten', bool(rewritten))
            rspan.set_attribute('retrieval.top_k', 5)
            chunks = await retriever.hybrid(rewritten or message, user)
            rspan.set_attribute('retrieval.result_count', len(chunks))
        ...

span 的粒度原则:能独立成为优化对象的操作各占一个 spanvector.searchfulltext.search 分开,因为它们的优化手段完全不同;llm.planllm.answer 分开,因为一个是选工具(输出很短),一个是生成回答(输出很长),延迟特征差一个数量级。

每个 span 该带的属性:request_idconversation_idmodeltokens_intokens_out。不带的:用户问题原文、文档正文、工具参数、Authorization 头。

脱敏要在采集侧做,不能在展示侧做。 采集时就不写进去,才叫脱敏;写进去之后在 UI 上隐藏,那些数据仍然在 trace 后端的存储里,仍然会进备份,仍然可能被任何有查询权限的人看到。

request_id 要一路透传到前端,让用户报障时能直接给出它。「大概下午三点多」和 req_8f3a2c 之间的排查成本差着几个小时。

指标:四类,不要更多

chat_requests = meter.create_counter('agent.chat.requests')
first_token = meter.create_histogram('agent.chat.first_token_ms')
tool_errors = meter.create_counter('agent.tool.errors')
token_usage = meter.create_counter('agent.llm.tokens')

四类指标:

  • 成功率与拒答率。 拒答率单独一个指标,而且它升高不一定是坏事——可能是第 12 篇的阈值在正常工作。要和引用正确率一起看:拒答率升、引用正确率也升,说明系统变诚实了;拒答率升、引用正确率没变,才是检索退化了。
  • 首 token 延迟与总延迟。 流式场景下这两个数的意义完全不同。首 token 决定用户感知的「响应快不快」,总延迟决定「答完要等多久」。只报总延迟会掩盖首 token 的退化。
  • 工具错误与审批等待。tool_nameerror_code 分维度。审批等待时长的分布能告诉你审批超时时间设得合不合理。
  • Token 与成本。 下一节单独说。

指标的维度(label)要控制基数。按 tool_nameerror_codemodel 分是安全的;按 user_idconversation_id 分会让时序数据库爆掉——那些是 trace 该干的事,不是指标。

流式成本:需要 tiktoken 兜底

非流式请求的 usage 字段直接给出 token 数。流式响应默认不返回 usage,OpenAI 的解法是 stream_options: {include_usage: true}——但很多 OpenAI-compatible 服务端没实现这个参数。有的忽略它,有的直接报 400。

所以成本统计必须有兜底:

def compute_cost(usage: dict | None, prompt: str, answer: str,
                 pricing: Pricing) -> CostRecord:
    if usage:
        tokens_in = usage['prompt_tokens']
        tokens_out = usage['completion_tokens']
        estimated = False
    else:
        tokens_in = estimate_tokens(prompt, pricing.model)
        tokens_out = estimate_tokens(answer, pricing.model)
        estimated = True
    return CostRecord(
        tokens_in=tokens_in, tokens_out=tokens_out,
        usd=(tokens_in / 1e6 * pricing.input_usd
             + tokens_out / 1e6 * pricing.output_usd),
        model=pricing.model, pricing_version=pricing.version,
        estimated=estimated,
    )

三个字段容易被漏掉,但缺了就没法回头分析:

estimated 标记这条记录是实测还是估算。混在一起算月度总成本,你不知道误差有多大。tiktoken 对中文的估算通常偏低 10% 左右,因为它的 BPE 是按英文语料训的。

pricing_version 记录用的是哪版价格表。模型降价或涨价之后,历史记录不该被重算——那会让上个月的成本曲线突然变形。

model 必须落库。同一个请求链路里可能用了三个模型:改写用小模型、重排用小模型、回答用大模型。混成一个数字,优化方向就错了:

# 不要这样
total_cost += cost

# 要这样
costs['rewrite'] += rewrite_cost
costs['rerank'] += rerank_cost
costs['answer'] += answer_cost
costs['embedding'] += embedding_cost

分开记之后,一个常见的发现是:重排的成本占比远超预期。第 11 篇那个逐条打分的实现,输入是 20 条 chunk 拼起来的一万字符,每次问答都要跑一遍。它可能比主回答还贵。

Embedding 的成本要单独看,因为它的计费模式不同——只在索引时发生,是一次性投入,不随问答量增长。把它混进单问成本会让那个数字失真。

限流、超时与降级

限流采用 Redis 固定窗口计数器,按用户和 IP 双维度落地:

async def check_rate_limit(user_id: str, ip: str) -> None:
    for key, limit, window in (
        (f'rl:user:{user_id}', 20, 60),
        (f'rl:ip:{ip}', 60, 60),
    ):
        count = await redis.incr(key)
        if count == 1:
            await redis.expire(key, window)
        if count > limit:
            raise RateLimited(retry_after=await redis.ttl(key))

按 IP 限流是为了防未登录路径和批量注册;按用户限流才是主要防线。固定窗口实现极轻量,单次 INCR 即可判断;若需防范窗口边界突发流量,可平滑升级为滑动窗口或基于 Lua 的令牌桶。返回 429 时带上 Retry-After,让前端知道该等多久。

模型调用要设三个超时,不是一个:

timeout = httpx.Timeout(connect=3.0, read=90.0, pool=5.0)

connect 短,连不上要快速失败;read 长,因为流式响应的整体时长本来就长;额外还要有一个首 token 超时——连上了但 30 秒不吐第一个字,通常意味着上游排队严重,这时该降级而不是继续等。

降级路径要显式写出来,而且每一级都不能静默:

async def answer_with_fallback(query: str, user: User) -> Answer:
    candidates = []
    try:
        candidates = await asyncio.wait_for(retriever.hybrid(query, user), 1.5)
        return await asyncio.wait_for(rerank_and_answer(candidates, query), 25)
    except TimeoutError:
        if candidates:
            logger.warning('rerank_timeout_fallback', extra={'query_hash': h(query)})
            return await answer_from_top_chunks(candidates[:3], query)   # 跳过重排
        raise ServiceUnavailable('retrieval_timeout')
    except UpstreamUnavailable:
        raise ServiceUnavailable('model_temporarily_unavailable')

降级的两条纪律:

不能静默切换到权限更宽的模型或数据源。 主模型不可用时返回明确的「暂时无法回答」,比用一个没接权限过滤的备用链路给出答案要好得多。可用性不能拿正确性换。

每次降级都要打日志和指标。 降级如果不可见,它就会变成常态——系统一直在降级运行,报表上的延迟很漂亮,而没人知道回答质量已经掉了。

Docker Compose 拓扑

services:
  web:
    build: ./web
    ports: ['3000:80']
  api:
    build: ./api
    command: uv run uvicorn ai_agent_guide.main:app --host 0.0.0.0
    env_file: .env
    depends_on: [postgres, redis]
    healthcheck:
      test: ['CMD', 'curl', '-f', 'http://localhost:8000/health']
      interval: 10s
      timeout: 3s
      retries: 5
  worker:
    build: ./api                    # 同一个镜像
    command: uv run arq ai_agent_guide.worker.WorkerSettings
    env_file: .env
    depends_on: [postgres, redis]
  postgres:
    image: pgvector/pgvector:pg16
    volumes: ['pgdata:/var/lib/postgresql/data']
  redis:
    image: redis:7-alpine
volumes:
  pgdata:
Preview
服务拓扑与发布顺序

apiworker 用同一个镜像、不同 command。这保证它们的依赖版本永远一致——两个镜像分开构建,迟早会出现 worker 用旧版解析器、api 用新版 schema 的情况。

depends_on 只保证启动顺序,不保证服务已就绪。Postgres 容器启动到能接受连接有几秒差距,api 在这期间连接会失败。要在应用层重试:

async def wait_for_db(max_attempts: int = 30) -> None:
    for attempt in range(max_attempts):
        try:
            async with engine.connect() as conn:
                await conn.execute(text('select 1'))
            return
        except OperationalError:
            await asyncio.sleep(1)
    raise RuntimeError('database_unreachable')

有状态和无状态的区分决定了运维操作的边界。webapiworker 可以随时重启、并行多份、直接换镜像回滚;postgresredis 带数据卷,回滚要单独考虑。

发布与回滚

顺序是固定的:

  1. 执行向前兼容的数据库迁移(只加列、加表、加索引)
  2. 启动新版本的 api 和 worker
  3. 健康检查通过后切流
  4. 观察 15 分钟

关键规则:回滚只回滚镜像,不回滚已执行的迁移。

这条规则反过来约束了迁移的写法。「删掉一个列」不能一次做完,要拆成两次发布:

发布 N:   代码不再读写 old_column(列还在)
发布 N+1: 迁移删掉 old_column

这样发布 N 出问题时,回滚到 N-1 的镜像仍然能工作,因为列还在。如果一次做完,回滚后旧代码会去读一个已经不存在的列,直接 500。

改列类型同理,要经过「加新列 → 双写 → 回填 → 切读 → 删旧列」这条路。这很啰嗦,但它是唯一能安全回滚的路径。

第 13 篇的 checkpoint 表在这里要单独提一句:它和业务表在同一个 postgres 里,但备份和清理策略不同。checkpoint 表增长很快(每个节点执行写一次),过期会话的检查点没有保留价值。配一个定期清理,保留最近 7 天,或只保留 approvals 里还有 pending 的那些 thread_id

回到群里:机器人闭环

这个系统的语料来自群聊,但到目前为止答案只出现在 Web 界面上。用户得离开正在提问的地方,去另一个页面问同一个问题——这一步流失掉的人比你想的多。

最后一公里是把 Agent 接回群里。钉钉的 Stream 模式不需要公网回调地址,本地起一个长连接就能收消息:

async def on_bot_message(event: BotMessage) -> None:
    if not event.is_at_bot:                     # 只处理明确 @ 机器人的消息
        return
    async with tracer.start_as_current_span('bot.message') as span:
        span.set_attribute('conversation.ref', event.conversation_ref)
        result = await agent.answer(
            event.text,
            user=await resolve_user(event.sender_ref),   # 映射到内部身份
            session_key=f'dingtalk:{event.conversation_ref}:{event.sender_ref}',
        )
        await bot.reply(event, render_for_im(result))

四个和 Web 端不同的地方:

只处理明确 @ 的消息。 机器人本来也只能收到这些——这是平台的限制,不是你的选择。

身份要映射。 第 6 篇的 RBAC 依赖内部用户,而群里来的是 senderRef。映射不到就按最低权限处理,绝不因为「他在群里」就放行内部文档。

会话键带上会话引用。 同一个人在不同群里的追问是不同的上下文;共用一个会话会让第 12 篇的多轮字段收集串台。

渲染要换一套。 群消息不支持折叠引用面板,把引用压成两行尾注,超长回答截断并附 Web 链接。审批按钮在 IM 里是交互卡片,不是 Web 的按钮组件——但 approval_required 事件本身不用改,这是第 7 篇把审批做成协议而不是 UI 的回报。

Preview
从群消息到语料再回到群里的闭环

图里那个回路是这个专栏真正的终点:群里的提问变成语料,Agent 用语料回答,回答本身又成为群里的新消息,下一次采集把它收进来。

这里有一个必须写进文档的边界,第 8 篇也提过:机器人只能收到满足触发条件的新消息,拿不到历史。 它是增量来源,不能替代数据库采集。如果哪天有人提议「用机器人收集就够了,不用做那套解密和 WAL 重放」,答案是不行——机器人上线之前的所有讨论,它一条都看不到。

还有一条运营上的:机器人答错会被全群看到。Web 端答错只有一个人尴尬,群里答错是公开的。所以群内回复的拒答阈值应该比 Web 端更保守,宁可说「这个我不确定,建了工单」。

压测与最终评测

docker compose up -d
uv run locust -f tests/load/locustfile.py --headless -u 50 -r 5 -t 5m
uv run python -m ai_agent_guide.eval --provider real --repeat 3

压测要用 Mock provider 跑一遍,再用真实模型跑一遍小规模的。Mock 版本测的是你自己的代码——连接池、锁竞争、序列化开销;真实版本测的是外部依赖的限流表现。混在一起你分不清瓶颈在哪边。

--repeat 3 是因为真实模型有随机性。跑一次的分数不能作为基线,取三次的中位数。

Preview
从第 4 篇到第 15 篇的评测曲线

这条曲线是第 4 篇先建评测基线的全部理由。几个值得读的地方:

第 5 到第 7 篇几乎是平的。 换 PostgreSQL、加权限、做审批幂等,一分没涨。但没有它们,这个系统不能上线。这是工程投入里最容易被忽视的一类——它们不改善指标,只是让系统合法。

第 10 篇是最大的一次跳升。 群聊语料从「逐条 embedding」换成线程重建之后,长尾问题第一次有了答案来源。这也是全专栏投入产出比最高的一次改动。

第 11 篇跳升第二大,也是最大的一次延迟代价。 混合召回加重排,同时把延迟从 1.1 秒推到 2.4 秒。这个取舍必须显式做,不能只看左边那条线。

第 13 篇不涨分。 用 LangGraph 重写买的是可恢复性,不是质量。如果当时期待它提升回答质量,那就是选错了工具。

第 15 篇延迟降回 1.9 秒而分数没掉。 限流、降级、缓存的收益就在这里。

图里的数字是示意基线。真正的价值在于你自己的仓库能画出同样一条曲线,每个点都对应一次 commit——任何指标下降都能追溯到具体改动。

最终质量门槛

  • 自动化测试全部通过,包括第 12 篇的注入 fixture
  • 仓库自带语料上 Recall@5 ≥ 0.80,引用正确率 ≥ 0.85
  • 关键工具成功率 ≥ 0.90,高风险操作越权次数为 0
  • 跨来源冲突的用例上,回答分组呈现两个来源且各自带时间
  • docker compose up 能完整启动,健康检查通过,迁移可重复执行
  • 任意请求能通过 request_id 找到 trace、工具审计记录和成本记录
  • 审批挂起期间重启服务,resume 后不丢消息、不重复写入
  • README 的命令能从空目录复现到 part-15

十五篇之后

这个系统现在是:可评测的(第 4 篇的基线一直在跑)、可审计的(每次工具调用都有记录)、可恢复的(审批跨重启存活)、可部署的(一条命令起全套),而且语料会自己长大(群里每天的新讨论都是下一批 chunk)。

比「调用过一个大模型 API」更值得写进简历的,是能讲清其中的取舍:为什么第 5 篇才换数据库、为什么权限要在检索侧和工具侧各拦一次、为什么群聊不能用文档那套切分、为什么两个来源冲突时不让模型仲裁、为什么重排用 LLM 而不是 cross-encoder、为什么 LangGraph 只买到可恢复性。

这些取舍在面试里比技术清单有用得多,因为它们无法从文档里背出来。

回到 专栏目录,可以按顺序对照每一次演进——每篇的验收标准都是可执行的,跑一遍比读一遍收获大。