跳到主要内容

Agent 可观测性:会话追踪、打分与在线评估

30 秒导读: 前面几章讲的是"一条 LLM 请求怎么被截获、算成本、投队列、加工成一行结构化日志"(见 0304)。本章讲的是在这行日志之上再长出三样 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.tslib/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;拿到浮点直接 throwScoreStore 里再兜一层,非整数值 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)
PostHogHandlerPostHog 产品分析请求头带 posthogApiKeycaptureEvent(PostHogHandler.ts:53)
SegmentLogHandlerSegment CDP组织配了 segment 集成POST api.segment.io/v1/track(SegmentLogHandler.ts:82)
LytixHandlerLytix 评估平台请求头带 lytixKeyPOST {host}/v1/metrics/modelIO(LytixHandler.ts:61)

这就是"旁挂"的字面意思:数据主体仍是那一条请求日志,外发只是把它的一份拷贝抄送出去,不改变主链路、不阻塞落库(发送都在 handleResults 里、日志已落库之后批量做)。

6. 为什么是"派生/旁挂"而不是另建管道

把三样能力放一起看,设计取舍就清楚了。核心一句:Helicone 死守"单一事实源"——所有 agent 观测能力都长在 request_response_rmt 这一行日志上,绝不复制第二份主数据。

能力手法靠什么复用主日志
会话追踪派生(物化视图)session 信息本就在 properties 里,视图只是抽键+过滤成快查表
同步打分/在线评估内联(责任链一环)落库前写进 processedLog.scores,随主行一次写入
事后打分回读重插旁路独立队列,但落点仍是主表 scores 列,靠 ReplacingMergeTree 去重
外发(webhook/PostHog/…)旁挂(责任链尾部)复用同一个 context,抄送整行,不改主链路

这么做的好处很直接:

  1. 不重复存储、不数据漂移。 只有一张主表是真相,会话表是它的视图、分数是它的字段。不会出现"会话表说花了 $3、主表说花了 $5"这种对不上。
  2. 加能力 = 加一环/加一个视图,不动主链路。 想接个新分析平台,就在责任链尾部再挂一个 handler;想加个新派生视图,写一条 CREATE MATERIALIZED VIEW。主管道(03/04 章)一行不用改。
  3. 延迟解耦。 迟到的分数走独立队列 + 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)。
  • 在线评估默认全采样但会抽样降载:sampleRate 默认 100(每条都评),高流量下不设采样会把评估器(尤其 LLM 评估器)成本顶上去;handler 里靠 Math.random()*100 > sampleRate 做闸(OnlineEvalHandler.ts:57)。
  • 会话名查询有硬窗口:getSessionNames 写死只看最近 150 天、LIMIT 50(SessionManager.ts:252)。

8. 代码地图(导航索引)

按符号名可 grep 定位;行号 as-of sourceCommit

主题文件符号
会话物化视图定义clickhouse/migrations/schema_50_sessions_mv.sqlsession_rmt_mv
会话表 DDL(主键含 session)clickhouse/migrations/schema_49_sessions.sqlsession_rmt
会话历史回填(30 天)clickhouse/migrations/schema_51_sessions_backfill.sqlINSERT INTO session_rmt
会话聚合(成本/时长/请求数)valhalla/jawn/src/managers/SessionManager.tsgetSessionsgetMetricsgetSessionNames
会话反馈/标签(复用 property/tags)valhalla/jawn/src/managers/SessionManager.tsupdateSessionFeedbackupdateSessionTag
会话 API 路由valhalla/jawn/src/controllers/public/sessionController.tsSessionController
OTEL trace → 主队列valhalla/jawn/src/managers/traceManager.tsconsumeTracessendLogToKafka
trace 入口valhalla/jawn/src/controllers/public/traceController.tsTraceController.logTrace
主表 scores 列(Map)clickhouse/migrations/schema_30_request_response_versioned_merge_tree.sqlscoresidx_scores_key
打分入口 APIvalhalla/jawn/src/controllers/public/requestController.tsaddScores(:281)
打分队列/延迟/DLQvalhalla/jawn/src/managers/score/ScoreManager.tsaddBatchScoresgetDefaultDelayMsmapScores
打分消费入口valhalla/jawn/src/lib/consumer/consumeMiniBatchScores.tsconsumeMiniBatchScores
Kafka→打分消息解析valhalla/jawn/src/lib/consumer/helpers/mapKafkaMessageToScoresMessage.tsmapKafkaMessageToScoresMessage
回读—合并—重插valhalla/jawn/src/lib/stores/ScoreStore.tsputScoresIntoClickhouse
在线评估环valhalla/jawn/src/lib/handlers/OnlineEvalHandler.tsOnlineEvalHandler.handle
在线评估配置读取valhalla/jawn/src/lib/stores/OnlineEvalStore.tsgetOnlineEvalsByOrgIdhasOnlineEvals
责任链装配顺序valhalla/jawn/src/managers/LogManager.tsauthHandler.setNext(...)(:104)
落库时写 scoresvalhalla/jawn/src/lib/handlers/LoggingHandler.tsLoggingHandler(:565)
外发 handlersvalhalla/jawn/src/lib/handlers/{Webhook,PostHog,SegmentLog,Lytix}Handler.ts各自 handle / handleResults

同组其它章:index · 01 边缘代理与 AI 网关 · 02 捕获算成本投队列 · 03 消费与责任链 · 04 三处落库