数据截至 (上游 commit 67df07b8d807)
Agent 可观测性:会话追踪、打分与在线评估
30 秒导读: 前面几章讲的是"一条 LLM 请求怎么被截获、算成本、投队列、加工成一行结构化日志"(见 03、04)。本章讲的是在这行日志之上再长出三样 agent 级能力:把散落的请求串成一次"会话"、给每条请求挂上"分数"、用评估器自动打分并把数据外发。关键是:这三样全都没有另建管道——要么从主日志表物化派生,要么旁挂在同一条消费责任链上。
1. 这是什么(零基础也能懂)
从"单条请求"到"一次 agent 运行"
一次 agent 跑起来,底下往往是十几次乃至上百次 LLM 调用:规划一步、调工具一步、反思一步……如果你的观测台只能一条一条看请求日志,你根本看不出"这一整轮 agent 干了什么、花了多少钱、在哪一步崩的"。
Agent 可观测性(agent observability)要补的正是这个缺口。它围绕三个问题:
| 你想知道的 | 对应能力 | 一句话 |
|---|---|---|
| 这一整轮 agent 是哪些请求组成的? | 会话追踪(Session) | 把同一次运行的多条请求归到一个 session_id 下 |
| 这条(或这次会话)回答得好不好? | 打分/反馈(Scores) | 给请求挂上数值/布尔分,人打或程序打 |
| 能不能不用人工、自动判好坏? | 在线评估(Online Eval) | 日志经过时用评估器自动打分,并把数据外发到分析工具 |
用起来什么样
会话追踪对用户几乎零成本——只多传两个请求头:
POST https://oai.helicone.ai/v1/chat/completions
Helicone-Auth: Bearer sk-helicone-...
Helicone-Session-Id: "run-2f9c... ← 这一轮 agent 的唯一 id
Helicone-Session-Name: "researcher" ← 这类 agent 的名字
...正常的 OpenAI 请求体...
打分则是事后补一刀,对着某条请求 id 追加分数:
POST /v1/request/{requestId}/score
{ "scores": { "helpfulness": 8, "contains_pii": false } }
就这么两个动作。剩下的"会话怎么聚合、分数落到哪张表、评估器什么时候跑",全在服务端 jawn 里完成。
一句话直觉
会话/评估不是新数据库,是主日志的"视图"和"批注"。 把主表 request_response_rmt(每行一条请求,见 04 章)当"流水账":
- 会话 = 按
Helicone-Session-Id分组的一个"派生视图"; - 分数 = 在流水账那一行上补写的一个
scores字段; - 在线评估 = 日志流过时顺手算出分数、再把整行数据抄送给外部工具。
这就是本章的主线:派生 + 旁挂,绝不另起炉灶。
2. 顶层全景(它大概怎么转)
一张图看清"派生"与"旁挂"
先说怎么读这张图:中间竖线是第 03/04 章讲过的主链路(请求 → 责任链 → 主表),本章的三样能力全部挂在它的左右两侧——上面是"顺流内联"(评估),下面是"旁挂外发",右边是"事后派生"。
┌───────────────────────────────────────────┐
一条请求日志 ───▶ │ 消费责任链 (LogManager 责任链, 见03章) │
│ │
│ ... → OnlineEvalHandler ──┐(顺流:算分内联) │
│ ↓ 把分数写进 processedLog │
│ LoggingHandler ─────┼──▶ 主表 │
│ ↓ │ request_ │
│ PostHog/Lytix/ │ response_rmt │
│ Webhook/Segment ────┘(旁挂:整行外发) │
└───────────────┬─────────────────────────────┘
│
┌───────────────┴───────────────┐
(事后派生) │ │ (事后补分)
▼ ▼
物化视图 session_rmt_mv 独立打分队列 helicone-scores-prod
↓ 过滤+改写 ↓ consumeMiniBatchScores
会话表 session_rmt 回读主表 → 合并 scores → 重插
(schema_50/51) (ScoreStore, ReplacingMergeTree 去重)
部件一句话职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
session_rmt_mv | 物化视图:主表新行只要带 session 头就抄进会话表 | clickhouse/migrations/schema_50_sessions_mv.sql |
session_rmt | 会话专用表,主键含 session_id,按会话查很快 | clickhouse/migrations/schema_49_sessions.sql |
SessionManager | 服务端按 session 分组聚合(成本/时长/请求数) | valhalla/jawn/src/managers/SessionManager.ts |
TraceManager | 收 OTEL trace,拆成 span 转成日志投回主队列 | valhalla/jawn/src/managers/traceManager.ts |
scores 列 | 主表上的 Map(String, Int64) 字段,存每条请求的分 | schema_30/41/49 ...merge_tree.sql |
ScoreManager / ScoreStore | 事后补分:队列消费 → 回读主行 → 合并分数重插 | valhalla/jawn/src/managers/score/ScoreManager.ts、lib/stores/ScoreStore.ts |
OnlineEvalHandler | 责任链上的一环:日志流过时按配置自动打分 | valhalla/jawn/src/lib/handlers/OnlineEvalHandler.ts |
| Webhook/PostHog/Segment/Lytix Handler | 责任链尾部:把整行数据外发到外部评估/分析工具 | lib/handlers/{Webhook,PostHog,SegmentLog,Lytix}Handler.ts |
3. 核心原理之一:会话追踪 = 从主表物化派生
它要解决的小问题
同一次 agent 运行的几十条请求,散在主表 request_response_rmt 里。若每次"看一次会话"都去主表全扫、按 properties['Helicone-Session-Id'] 分组,既慢又贵(主表还扛着全量流量)。
思路:两个请求头 → 一个 properties 键 → 一张派生表
会话信息根本不是单独字段,它就藏在每条请求的 properties Map 里。用户传的 Helicone-Session-Id / Helicone-Session-Name 两个头,在 03 章的加工阶段被塞进 properties。会话能力要做的,只是把带这两个键的行,派生进一张查询更快的专用表。
主表 request_response_rmt 会话表 session_rmt
┌───────────────────────────┐ ┌──────────────────────────┐
│ properties: │ 物化视图 │ session_id (独立列) │
│ {Helicone-Session-Id: X, │ ───────▶ │ session_name(独立列) │
│ Helicone-Session-Name:Y}│ 抽键改写 │ + 其余字段原样带过来 │
│ ...其余 40 个字段 │ │ 主键含 session_id → 查得快 │
└───────────────────────────┘ └──────────────────────────┘
每来一行,带 session 头的就被抄一份过去
真实实现:一个 WHERE 就是全部魔法
物化视图(materialized view,ClickHouse 里"插入触发的增量派生表")的定义几乎全是把主表列原样搬,只多做两件事——把 Map 里的键抽成独立列 + 只放带 session 头的行:
-- clickhouse/migrations/schema_50_sessions_mv.sql:1 (session_rmt_mv)
CREATE MATERIALIZED VIEW session_rmt_mv TO session_rmt AS
SELECT
properties['Helicone-Session-Id'] AS session_id, -- 抽键成列
properties['Helicone-Session-Name'] AS session_name,
response_id, latency, status, ... , scores, ... -- 其余原样搬
FROM request_response_rmt
WHERE (has(properties, 'Helicone-Session-Id')) -- 只要带 session 头的
派生的目标表 session_rmt 之所以"按会话查得快",在于它的主键就是拿 session 排序的,而主表主键不是:
-- clickhouse/migrations/schema_49_sessions.sql:41 (session_rmt)
ENGINE = ReplacingMergeTree(updated_at)
PRIMARY KEY (organization_id, session_name, session_id, request_id)
ORDER BY (organization_id, session_name, session_id, request_id)
关键细节:物化视图只对"未来"的行生效。 建视图前已有的历史数据不会自动进来,所以要配一支回填脚本,把最近 30 天补进去:
-- clickhouse/migrations/schema_51_sessions_backfill.sql:2
INSERT INTO session_rmt
SELECT properties['Helicone-Session-Id'] AS session_id,
properties['Helicone-Session-Name'] AS session_name, *
FROM request_response_rmt
WHERE request_created_at > earliest_date - INTERVAL 30 DAY
AND has(mapKeys(properties), 'Helicone-Session-Id')
服务端只做聚合,不碰"归属"
SessionManager 里没有任何"把请求分配到会话"的逻辑——归属早在物化视图那步定死了。它只负责按 session 分组算指标:成本、时长(首末请求时间差)、请求数,全是对着 properties['Helicone-Session-Id'] 做 GROUP BY:
// managers/SessionManager.ts:382 (getSessions)
GROUP BY properties['Helicone-Session-Id'], properties['Helicone-Session-Name']
// 时长 = dateDiff('second', min(request_created_at), max(request_created_at))
对外这些方法挂在 SessionController(controllers/public/sessionController.ts:73)的 POST /v1/session/* 路由上:query(列会话)、metrics/query(指标直方图)、name/query(按名聚合)。
两个"会话级"写操作值得单独点出——它们也不新建存储,而是复用既有机制:
updateSessionFeedback(SessionManager.ts:485):给整个会话点赞/踩,做法是找到会话的第一条请求,给它追加一个Helicone-Session-Feedback属性——反馈也是一条 property。updateSessionTag(SessionManager.ts:554):往独立的tags表插一行,entity_type = SESSION。
旁支:OTEL trace 也是"投回主队列"
如果 agent 用的是 OpenTelemetry/Traceloop 那套埋点,数据从 POST /v1/trace/log(traceController.ts:107)进来。TraceManager.consumeTraces(traceManager.ts:187)把每个 OTEL span 拆开——gen_ai.prompt.* 拼成 messages、gen_ai.completion.* 拼成 choices、traceloop.association.properties.Helicone-* 还原成 Helicone 属性——再包成一条普通日志投回主队列:
// managers/traceManager.ts:159 (sendLogToKafka)
await kafkaProducer.sendMessages([kafkaMessage], "request-response-logs-prod");
注意投的就是主日志 topic request-response-logs-prod(和 02 章边缘代理投的同一个)。所以 OTEL trace 不是第二条管道,而是"翻译成主日志格式后汇入主管道"——会话追踪对它自然也就免费生效了。
4. 核心原理之二:打分 = 主表上一个 Map 字段 + 一条补分旁路
分数存在哪:不是新表,是主行上的一列
分数没有独立的分数表。主表 request_response_rmt 从一开始就带了一个 scores 列,类型是字符串→整数的 Map:
-- clickhouse/migrations/schema_30_request_response_versioned_merge_tree.sql:23
`scores` Map(LowCardinality(String), Int64) CODEC(ZSTD(1)),
INDEX idx_scores_key mapKeys(scores) TYPE bloom_filter(0.01) GRANULARITY 1,
INDEX idx_scores_value mapValues(scores) TYPE bloom_filter(0.01) GRANULARITY 1
一条请求的所有分就是这个 Map 的键值对({"helpfulness": 8, "contains_pii-hcone-bool": 0})。两个布隆过滤器索引让"按分数键/值筛请求"也不慢。分数是整数——布尔被折成 1/0,浮点直接被拒(见下)。
难点:分数常常"迟到"
打分有两种时机:
- 同时到:请求日志加工时分数已经算出(在线评估就是这种,见第 5 节)——这种分数在
LoggingHandler写主行时顺手一起写进scores字段(lib/handlers/LoggingHandler.ts:565)。 - 事后到:人工审完、或离线评估器几分钟后才出分,这时那条请求早已落库。
第二种是难点。主表用的是 ReplacingMergeTree(靠 updated_at 保留最新版的引擎),没有"就地改一个字段"的能力。所以补分的唯一办法是:回读那一整行 → 把新分并进它的 scores → 整行重新插入,让引擎自己按 updated_at 去重留新。
补分旁路:一条和主链路平行的独立队列
事后补分走的是独立的 topic helicone-scores-prod,不挤主日志管道。整条路是:
POST /v1/request/{id}/score requestController.ts:281 (addScores)
↓ ScoreManager.addScores
↓ 发到独立队列 helicone-scores-prod (ScoreManager.ts:138 sendScoresMessage)
↓ (还会挂一个默认 10 分钟的延迟再发一次,等日志先落库)
消费: mapKafkaMessageToScoresMessage → 解析成 HeliconeScoresMessage[]
↓ consumeMiniBatchScores → new ScoreManager → handleScores
↓ ScoreStore.putScoresIntoClickhouse
1. 回读主表这些 (request_id, org_id) 的现有行
2. combinedScores = 旧 scores ∪ 新 scores
3. 整行带新 scores 重新插入 → ReplacingMergeTree 去重留最新
为什么发两次、还默认延迟 10 分钟?因为补分可能比它要打分的那条请求还先到(分数队列和日志队列各跑各的)。延迟给日志留出落库时间:
// managers/score/ScoreManager.ts:59 (getDefaultDelayMs)
return process.env.NODE_ENV === "production" ? 10 * 60 * 1000 : 0; // 10 分钟
消费入口很薄,就是把 Kafka 消息交给 ScoreManager.handleScores:
// lib/consumer/consumeMiniBatchScores.ts:19
const scoresManager = new ScoreManager({ organizationId: "" });
await scoresManager.handleScores({ batchId: miniBatchId, ... }, messages);
真实实现:回读—合并—重插
ScoreStore.putScoresIntoClickhouse 是补分的核心。它先把这些请求的现有整行从主表捞回来,再把新旧分数并集,然后整行重插:
// lib/stores/ScoreStore.ts:117 (putScoresIntoClickhouse)
const combinedScores = {
...(row.scores || {}), // 旧分
...newVersion.mappedScores.reduce(...) // 新分(只收整数,非整数打日志跳过)
};
// 然后把 row 的全部 40+ 字段原样带上、只替换 scores,插回 request_response_rmt
两个容易踩的坑:
- 只收整数。 打分入口
mapScores(ScoreManager.ts:18)把布尔转成 1/0;拿到浮点直接throw。ScoreStore里再兜一层,非整数值console.log跳过(ScoreStore.ts:121)。 - 布尔分改名。 布尔分在存进 Map 前,键会被加后缀
-hcone-bool(ScoreManager.ts:200),这样前端才知道0/1该显示成"是/否"而不是数字。
5. 核心原理之三:在线评估 = 责任链上顺手打分 + 尾部外发
它要解决的小问题
前面的打分要么靠人、要么靠外部离线跑。在线评估(online evaluation)想做到:日志流过服务端时,就地用一个评估器(LLM 或代码)自动给它打分,不用用户再调一次接口。
思路:塞进已有的责任链,而不是新开一个 job
03 章讲过,jawn 处理每条日志走的是一条责任链(chain of responsibility):认证 → 取 body → 加工 → 落库 → 外发,一环调 super.handle 把 context 传给下一环。在线评估就是这条链上的一环 OnlineEvalHandler,插在"落库"之前:
LogManager.ts:104 责任链装配(顺序即执行序)
authHandler
→ rateLimit → s3Reader → requestBody → responseBody → prompt
→ onlineEvalHandler ← 在这里算分,写进 processedLog.request.scores
→ loggingHandler ← 落库时把上一步的分一起写进主表 scores 列
→ posthog → lytix → webhook → segment ← 尾部:把整行外发出去
因为它在 loggingHandler 之前,算出的分数只要写进 context.processedLog.request.scores,就会被紧接着的落库环顺手写进主表——评估分和请求日志同一次写入,零额外往返。
真实实现:抽样 → 跑评估器 → 写回 context
OnlineEvalHandler.handle 的逻辑:先用带缓存的查询问"这个组织有没有配在线评估"(没有就直接放行,不做任何事),有才逐个评估器跑:
// lib/handlers/OnlineEvalHandler.ts:29
const hasOnlineEvals = await cacheResultCustom(
"has-online-evals-" + orgId,
async () => await onlineEvalStore.hasOnlineEvals(orgId), kvCache); // 缓存 1 分钟
if (hasOnlineEvals.data === false) return await super.handle(context); // 无配置→放行
每个评估器先过抽样率和属性过滤两道闸(省钱:不必每条都评),再调评估器打分,最后把分写回 context:
// lib/handlers/OnlineEvalHandler.ts:46 / :95 / :121
const sampleRate = Number((onlineEval.config as any)?.["sampleRate"] ?? 100);
if (Math.random() * 100 > sampleRate || ...) continue; // 抽样 + 排除实验/评估自身流量
const result = await evaluatorManager.runLLMEvaluatorScore({ ... }); // 跑评估器
const scoreName = getFullEvaluatorScoreName(onlineEval.evaluator_name);
context.processedLog.request.scores[scoreName] = result.data?.score ?? 0; // 写回 → 随日志落库
配置从哪来?OnlineEvalStore(lib/stores/OnlineEvalStore.ts:49)从 Postgres 的 online_evaluators join evaluator 表读出评估器模板(LLM prompt 模板或代码模板)。注意存储分工:评估器"是什么"配在 Postgres,评估"产出的分"落进 ClickHouse 主表——和 04 章"元数据在 Postgres、日志在 ClickHouse"的分工一致。
尾部:同一条链把整行外发给外部工具
责任链末尾的四个 handler,是把加工好的整行数据旁挂外发给第三方评估/分析平台——一样复用 context,不重新查库。它们的共性都是:没有该组织的对应配置就直接 super.handle 放行,配了才收集、最后在 handleResults 里批量发:
| Handler | 外发到 | 触发条件 | 出口 |
|---|---|---|---|
WebhookHandler | 用户自建 webhook | 组织配了 webhook | 带请求/响应 body + 成本/token 元数据 + S3 签名 URL(WebhookHandler.ts:87) |
PostHogHandler | PostHog 产品分析 | 请求头带 posthogApiKey | captureEvent(PostHogHandler.ts:53) |
SegmentLogHandler | Segment CDP | 组织配了 segment 集成 | POST api.segment.io/v1/track(SegmentLogHandler.ts:82) |
LytixHandler | Lytix 评估平台 | 请求头带 lytixKey | POST {host}/v1/metrics/modelIO(LytixHandler.ts:61) |
这就是"旁挂"的字面意思:数据主体仍是那一条请求日志,外发只是把它的一份拷贝抄送出去,不改变主链路、不阻塞落库(发送都在 handleResults 里、日志已落库之后批量做)。
6. 为什么是"派生/旁挂"而不是另建管道
把三样能力放一起看,设计取舍就清楚了。核心一句:Helicone 死守"单一事实源"——所有 agent 观测能力都长在 request_response_rmt 这一行日志上,绝不复制第二份主数据。
| 能力 | 手法 | 靠什么复用主日志 |
|---|---|---|
| 会话追踪 | 派生(物化视图) | session 信息本就在 properties 里,视图只是抽键+过滤成快查表 |
| 同步打分/在线评估 | 内联(责任链一环) | 落库前写进 processedLog.scores,随主行一次写入 |
| 事后打分 | 回读重插旁路 | 独立队列,但落点仍是主表 scores 列,靠 ReplacingMergeTree 去重 |
| 外发(webhook/PostHog/…) | 旁挂(责任链尾部) | 复用同一个 context,抄送整行,不改主链路 |
这么做的好处很直接:
- 不重复存储、不数据漂移。 只有一张主表是真相,会话表是它的视图、分数是它的字段。不会出现"会话表说花了 $3、主表说花了 $5"这种对不上。
- 加能力 = 加一环/加一个视图,不动主链路。 想接个新分析平台,就在责任链尾部再挂一个 handler;想加个新派生视图,写一条
CREATE MATERIALIZED VIEW。主管道(03/04 章)一行不用改。 - 延迟解耦。 迟到的分数走独立队列 + 10 分钟延迟,和主日志各跑各的,互不阻塞。
7. 边界与局限(诚实)
- 物化视图只对未来行生效:建视图后的历史会话得靠回填脚本补,且脚本写死了只补最近 30 天(
schema_51_sessions_backfill.sql:15)。更早的会话进不了会话表。 - 分数必须是整数:浮点分会被
mapScores直接throw、或在ScoreStore里静默跳过(ScoreStore.ts:121)。想存 0.87 这种,只能自己先 ×100 转成整数。 - 补分靠"回读整行重插":代价是每次补分都要把那一整行(含 body)读回来再写一遍;注释里也留了 TODO 说去重是"hand rolling"、想改用
FINAL(ScoreStore.ts:58)。