跳到主要内容

三处落库:ClickHouse 分析引擎 + S3 大对象 + Postgres 元数据

30 秒导读: 上一章(03)把队列消息加工成了一条结构化日志。 这一章讲这条日志最终怎么落地。答案是拆成三份写到三个地方:能过滤、能聚合的分析字段进 ClickHouse 的一张宽表;又大又不常查的请求/响应正文进 S3(自建部署是 Minio);必须事务、 必须强一致的少量元数据留在 Postgres。三路并发写,谁都不等谁。核心看点是那张 ClickHouse 主表 request_response_rmt 的 schema 取舍——排序键、去重引擎、跳数索引、TTL——这些决定了 dashboard 为什么快。


1. 这一章讲什么:为什么一条日志要分三处

先给结论,再解释。一条 LLM 调用日志天生是"混合体":

数据成分例子特点该用什么库
分析维度/指标org、provider、model、tokens、cost、延迟、properties小、要过滤、要聚合、写多读多列式分析库(ClickHouse)
请求/响应正文几十 KB 到几 MB 的 prompt 和 completion JSON大、只在"看单条详情"时才读对象存储(S3 / Minio)
强一致的账务元数据prompt 版本、组织 onboarding 状态少、要事务、要 join关系库(Postgres / Supabase)

为什么不能一处装下? 三种成分的读写模式互相打架:

  • 把几 MB 的正文塞进 ClickHouse 主表,会拖慢每一次扫描聚合(哪怕这次查询根本不看正文)。
  • 把 tokens、cost 这类要 GROUP BY provider SUM(cost) 的字段放进 Postgres 行存,亿级行的聚合会跪。
  • 把 prompt 版本这种要事务 join 的东西放进 ClickHouse,它没有真正的事务和外键。

所以 Helicone 的选择是按访问模式分家:让每种数据去它最擅长的引擎。代价是写入端要做一次"分拣", 这就是本章主角 LoggingHandler 的活。

本章聚焦存储引擎与 schema 决策,不重复 03 章讲的责任链编排; sessions / scores 的物化视图放在 05 章


2. 顶层全景:一条日志的三向落库

LoggingHandler 是责任链的最后一环。它先把 HandlerContext 映射成三套目标结构,再三路并发写出去。

怎么读下面这张图:从上往下是时间顺序,底部三个框是并发发生的(Promise.all),不是先后。

一条已加工好的日志 (HandlerContext)


┌────────────────────────────┐
│ LoggingHandler.handle() │ ← 先做映射,攒进 batchPayload
│ · mapRequest / mapResponse│
│ · mapRequestResponseCH │
│ · mapS3Records │
│ · 决定 storageLocation │ ← 存哪、正文放哪,这里定
└────────────┬───────────────┘
│ handleResults() → Promise.all([...])
┌────────────┼────────────────────────────┐
▼ ▼ ▼
┌────────────┐ ┌──────────────────┐ ┌────────────────────┐
│ Postgres │ │ S3 / Minio │ │ ClickHouse │
│ LogStore │ │ uploadToS3() │ │ logToClickhouse() │
│ │ │ │ │ │
│ prompt 输入 │ │ 超大请求/响应正文 │ │ request_response_ │
│ org 状态 │ │ (>10MB 时) │ │ rmt 宽表 + 指标 │
└────────────┘ └──────────────────┘ └────────────────────┘
强一致元数据 冷的大对象 热的分析数据

三个落库出口的一句话职责:

出口代码落到哪装什么
insertLogBatchstores/LogStore.tsPostgresprompt 输入、组织 onboarding 标记
uploadToS3LoggingHandler.uploadToS3S3/Minio请求+响应正文(仅超阈值时)
logToClickhouseLoggingHandler.logToClickhouseClickHouserequest_response_rmt 主表 + 缓存指标

三路的发起在一处,file:line 直接看:

const [pgResult, s3Result, chResult] = await Promise.all([
this.logStore.insertLogBatch(this.batchPayload),
this.uploadToS3(),
this.logToClickhouse(),
]);

—— valhalla/jawn/src/lib/handlers/LoggingHandler.ts:292-296,handleResults()。三路任一出错就整体返回对应错误(pgError/s3Error/chError),交给上游决定重试还是进死信。


3. 存哪里?写入前的"三档"决策

在真正写之前,LoggingHandler 要先决定这条日志的正文放哪、甚至存不存。这一步很关键,因为它 直接决定 ClickHouse 主表里 storage_location 这一列的值,以及要不要给 S3 派活。

三档判断,依次短路(白话在右):

命中条件storageLocation含义
免费额度超了 且 不是 PTB 计费not_stored_exceeded_free正文一律不存(省钱),只留指标
正文体积 ≤ 10 MBclickhouse正文小,直接塞进 CH 主表,连 S3 都不碰
否则(> 10 MB)s3正文太大,踢去 S3,CH 里正文列留空

对应源码:

const isPTB = context.message.heliconeMeta.isPassthroughBilling;
if (context.message.heliconeMeta.freeLimitExceeded && !isPTB) {
context.storageLocation = "not_stored_exceeded_free";
} else {
context.storageLocation =
size && size <= S3_MIN_SIZE_THRESHOLD ? "clickhouse" : "s3";
}

—— LoggingHandler.ts:179-185。阈值 S3_MIN_SIZE_THRESHOLD = 10 * 1024 * 1024(:118)。 PTB(passthrough billing,直通计费)例外:哪怕超额度,正文也要存 S3 留作账务凭证(:178-179 注释)。

为什么这么设计? 小正文进 ClickHouse 省一次 S3 往返(读详情时不用再签 URL 去拉),大正文进 S3 避免污染主表扫描。这就是"按体积分流":

正文体积

≤10MB├──────────► 直接进 ClickHouse 的 request_body/response_body 列

>10MB├──────────► 进 S3,CH 列留空,读详情时按 key 回源

超额免费├──────────► 都不存,只留 tokens/cost/延迟等指标

映射正文文本的逻辑呼应了这三档——clickhouse 档直接 JSON.stringify 塞进列,s3 档走 heliconeRequestToMappedContent 生成预览再截断,not_stored_exceeded_free 档返回空串:见 requestResponseTextFromContext(),LoggingHandler.ts:641-678


4. ClickHouse 主表:request_response_rmt 深挖

这是整个可观测系统的心脏表。dashboard 上几乎每个过滤、每张图、每次聚合,最终都打在这一张表上。 它的建表语句(37 行)把所有 schema 取舍浓缩在一处:clickhouse/migrations/schema_41_request_response_replacing_merge_tree.sql

下面逐个决策拆开讲——每一条都直接影响"查得快不快、存得省不省、数据对不对"。

4.1 引擎:ReplacingMergeTree(updated_at)——用"后写覆盖"实现幂等去重

要解决的小问题: 队列可能重投,同一条 request_id 会被写两次;而且用户后来还能给一条已落库的 请求追加 property / score(见 03 章的反馈回填)。同一逻辑行会有多个物理版本,怎么办?

思路: 不做"先删后插"(ClickHouse 删除很贵),而是每次都只管往里插,让存储引擎在后台合并时 按排序键去重、保留 updated_at 最大的那一版。这就是 ReplacingMergeTree(updated_at) 的语义。

ENGINE = ReplacingMergeTree(updated_at)

—— schema_41...sql:34。追加 property 的实现正是"读出最新一版 → 改 properties → 整行重插", 靠更大的 updated_at 胜出:见 VersionedRequestStore.putPropertyIntoClickhouse, stores/request/VersionedRequestStore.ts:92-166

坑(必须知道): ReplacingMergeTree 的去重是最终一致——只在后台 merge 后才生效。所以查询时若 要保证读到去重后的唯一行,得显式带 FINAL 或按 updated_at DESC 取最新(store 里读单条就是 ORDER BY updated_at DESC LIMIT 1,VersionedRequestStore.ts:70-74)。

4.2 分区:PARTITION BY toYYYYMM(request_created_at)——按月切块

PARTITION BY toYYYYMM(request_created_at)

—— schema_41...sql:35。每个自然月一个分区目录。好处有二:

  • 查询裁剪:dashboard 默认查"最近 N 天",按月分区能让引擎直接跳过无关月份的数据块。
  • 运维批量:删旧数据、导出、TTL 清理都能以"整月分区"为粒度,DROP PARTITION 是元数据操作,极快。

4.3 排序键:ORDER BY (organization_id, provider, model, user_id, request_created_at, request_id)

这是最影响查询性能的一行。ClickHouse 的主键即排序键(此表 PRIMARY KEYORDER BY 同列), 它决定数据在磁盘上的物理排布,也决定稀疏索引怎么建。

PRIMARY KEY (organization_id, provider, model, user_id, request_created_at, request_id)
ORDER BY (organization_id, provider, model, user_id, request_created_at, request_id)

—— schema_41...sql:36-37。列的顺序是精心排的,遵循"最常用来过滤的列放最左"的最左前缀原则:

位置为什么排这
1organization_id每个查询必带的租户隔离条件,放最左,一刀切掉别家数据
2-4provider / model / user_iddashboard 高频过滤维度,低基数(适合排序聚簇)
5request_created_at时间范围过滤,配合月分区做二次裁剪
6request_id兜底唯一性,保证 ReplacingMergeTree 的去重粒度到单请求

直觉: 把它想成一本按"组织→厂商→模型→用户→时间"层层排好序的电话簿。任何"某组织、某模型、 最近一周"的查询,都能顺着排序前缀二分定位,而不是全表扫。

4.4 Map 列 + bloom_filter 跳数索引——让"按自定义标签过滤"也快

小问题: 用户会打任意自定义标签(properties,如 {"env":"prod","feature":"chat"})和评分 (scores)。这些 key 是动态的,不可能每个都建成一列。

思路: 用 ClickHouse 的 Map 类型一列装下所有键值对;再对 Map 的 keys/values 建 bloom_filter 跳数索引(skip index),让"某 property 等于某值"的查询能跳过大部分数据块。

`properties` Map(LowCardinality(String), String) CODEC(ZSTD(1)),
`scores` Map(LowCardinality(String), Int64) CODEC(ZSTD(1)),
...
INDEX idx_properties_key mapKeys(properties) TYPE bloom_filter(0.01) GRANULARITY 1,
INDEX idx_properties_value mapValues(properties) TYPE bloom_filter(0.01) GRANULARITY 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,

—— schema_41...sql:19-30

跳数索引 ≠ 传统索引: 它不指向具体行,而是给每个数据块存一个概率摘要,回答"这个块可能含有 X 吗"。bloom filter 有假阳性(0.01 = 1% 误报率)但无假阴性——说"没有"就一定没有,于是可以放心 跳过整块。key 和 value 分别建索引,properties['env']='prod' 这类查询两头都能剪枝。

4.5 正文列:ngram bloom 索引 + 3 个月 TTL

小正文(≤10MB 那档)会直接进这两列。为支持"全文搜 prompt 里含某词"且不至于永久占地,做了两件事:

`request_body` String DEFAULT ''
TTL toDateTime(request_created_at) + toIntervalMonth(3),
`response_body` String DEFAULT ''
TTL toDateTime(request_created_at) + toIntervalMonth(3),
...
INDEX idx_request_body_bloom request_body TYPE ngrambf_v1(4, 1024, 1, 0) GRANULARITY 1,
INDEX idx_response_body_bloom response_body TYPE ngrambf_v1(4, 1024, 1, 0) GRANULARITY 1

—— schema_41...sql:21-32

  • ngrambf_v1(4,...) = 4-gram bloom filter:把正文切成 4 字符片段建 bloom,让 LIKE '%关键词%' 这类子串搜索也能靠跳数索引剪块,而不是逐块扫字符串。
  • 3 个月 TTL:只有正文列过期清空(列级 TTL,DEFAULT ''),分析字段(tokens、cost、延迟) 永久保留。所以三个月后你还能看历史用量曲线,但点进单条已看不到 prompt 原文。这是"分析价值长、 正文价值短"的成本取舍。

4.6 一次插入怎么发出去

映射产物是 RequestResponseRMT 结构(db/ClickhouseWrapper.ts:304 起),攒成批后由 VersionedRequestStore.insertRequestResponseVersioneddbInsertClickhouse("request_response_rmt", ...) 写入,格式 JSONEachRow,并开 async_insert=1(服务端攒批)+ wait_end_of_query=1(等落盘确认, 避免"返回 200 后才报错")。见 ClickhouseWrapper.ts:46-63


5. 主表的增量演化史(schema_42 之后)

schema_41 是主表的"出生证",但真实生产字段是一版版 ALTER 加出来的。读这些迁移能看清 Helicone 产品是怎么长出来的——每加一列背后都是一个新功能。按需查下表(全部 ALTER TABLE request_response_rmt):

迁移加了什么背后的功能
schema_42从旧表 request_response_versioned 回填全量数据主表切换的数据搬迁
schema_43prompt_cache_write_tokens / prompt_cache_read_tokens提示词缓存计量
schema_44prompt_audio_tokens / completion_audio_tokens多模态语音 token
schema_47cache_reference_id / cache_enabled缓存命中溯源
schema_52cost UInt64 DEFAULT 0把成本落库(见下方精度处理)
schema_75storage_location LowCardinality(String) DEFAULT 's3'记录本条正文存哪(呼应 §3)
schema_76size_bytes UInt32 DEFAULT 0正文体积,做分流决策与统计
schema_78reasoning_tokens Int64 DEFAULT 0推理模型的思考 token

依据:clickhouse/migrations/schema_43_cache_tokens.sqlschema_44_audio_tokens.sqlschema_47_cache_to_request_response_rmt.sqlschema_52_add_cost_to_request_response_rmt.sqlschema_75_storage_location.sqlschema_76_size.sqlschema_78_reasoning_tokens.sql

cost 列的精度技巧(值得学): 成本是很小的浮点(一次调用可能 $0.0000123),浮点在聚合里会累积 误差,而 ClickHouse 对整数聚合又快又准。所以 Helicone 不存浮点,而是把美元乘以 10⁹ 存成 UInt64 定点整数:

const cost = Math.round(rawCost * COST_PRECISION_MULTIPLIER);

—— LoggingHandler.ts:519,COST_PRECISION_MULTIPLIER = 1_000_000_000(packages/cost/costCalc.ts:8)。 读出来再除回去。一句话:用定点整数换聚合的精度与速度。

另外 §3 说过的缓存命中还会写第二张表 cache_metrics,引擎是 AggregatingMergeTree,列用 SimpleAggregateFunction(sum, ...) 预聚合各种"节省的 token/延迟",按 (org, date, hour, request_id) 排序——专门为"缓存帮你省了多少"这类累加查询而生。见 schema_48_cache_metrics.sqlmapCacheMetricCH()(LoggingHandler.ts:589-639)。


6. S3 / Minio:大对象的家

超过 10MB 的正文(以及 PTB 计费的正文)走 S3。自建部署用 Minio 顶替,协议兼容,所以代码里只有一个 S3Client(shared/db/s3Client.ts),靠 endpoint 环境变量切换云端还是本地。

为什么正文不进 ClickHouse 主查询路径? 一句话:让主表保持"瘦"。主表每次聚合都要扫过所有行的 存储块,正文哪怕不被 SELECT 也会拖慢合并、膨胀磁盘。把大对象挪到 S3,主表只留一个能回源的 key, "看单条详情"时才按需签 URL 去拉——冷数据冷读,不拖累热路径

关键设计点:

代码说明
key 布局getRequestResponseKey,s3Client.ts:352-354organizations/{orgId}/requests/{requestId}/request_response_body,天然按租户隔离
写入压缩store(),s3Client.ts:285-328正文 gzip 压缩后上传(ContentEncoding: gzip),压不动才存原文
请求响应合体uploadToS3,LoggingHandler.ts:332-338request 和 response 打包成一个 JSON 对象存一个 key,一次读取拿全
读取回源RequestResponseBodyStore.getRequestResponseBody,stores/request/RequestResponseBodyStore.ts:16-35先签一个有效期 1 天的 GET signed URL,再 fetch 拉回、JSON.parse
限流保护getLimiter/putLimiter,s3Client.ts:22-36用 Bottleneck 把并发压在 S3 配额内(读 5500/s、写 3500/s)

注意一个小分工:uploadToS3 里若某条 record 的 location === "clickhouse",直接跳过 S3 上传 (LoggingHandler.ts:321-324)——因为那条正文已经进 CH 主表了,不需要 S3 再存一份。


7. Postgres / Supabase:留下来的"元数据"

历史上 Helicone 把 request/response 也写 Postgres,现在这些表已删(LoggingHandler.ts:49 的 "Legacy type definitions for deleted tables" 注释、VersionedRequestStore.ts:52 的 DEPRECATED 注释都在 交代这段迁移)。如今 Postgres 侧的写入只剩两件必须事务/强一致的小事:

写什么代码为什么留在 Postgres
prompt 输入(prompts_2025_inputs)LogStore.processPromptInputsBatch,stores/LogStore.ts:60-115要先 join 校验 version_id 存在才插,是关系型事务活
组织 onboarding 标记LogStore.insertLogBatch,stores/LogStore.ts:22-38首次接入时把 org 标记为 has_integrated,幂等 UPDATE

两者都在一个 db.tx(...) 事务里完成(LogStore.ts:16-52)。prompt 输入插入前会先查出已存在的 版本 id,过滤掉指向不存在版本的脏数据再批量插——这种"先校验外键再写"正是 ClickHouse 做不了、只能 留给 Postgres 的原因。库结构类型定义见 supabase/(db/database.types.ts)。


8. 巧妙之处(可借鉴)

  • 按访问模式分家,而不是按业务分家。 同一条日志被拆到三个引擎,依据是"这份数据怎么被读",不是 "它属于哪个功能"。热的分析字段、冷的大正文、强一致的元数据各去所长——LoggingHandler.handleResults 的三路 Promise.all(:292-296)。
  • 成本用定点整数存。 美元 × 10⁹ 存 UInt64,躲开浮点聚合误差,换来又快又准的求和(:519)。
  • 同一列可以有独立 TTL。 正文列 3 个月过期、指标列永久,靠列级 TTL 一表两治(schema_41:21-24)—— 不用把冷热数据拆成两张表。
  • 动态标签用 Map + bloom 跳数索引。 不为每个自定义 key 建列,却仍能对任意标签快速过滤(:19-30)。
  • 去重靠引擎不靠代码。 ReplacingMergeTree 让"重投 + 后续回填"天然幂等,应用层只管无脑追加插入。

9. 边界与局限(诚实)

  • ClickHouse 去重是最终一致。 merge 之前同一 request_id 可能有多行;读单条要 ORDER BY updated_at DESCFINAL,否则会读到旧版本(VersionedRequestStore.ts:70-74)。
  • 正文 3 个月后消失。 过了 TTL,单条详情里看不到 prompt/response 原文,只剩指标(schema_41:21-24)。
  • S3 上传错误被吞。 uploadToS3await Promise.all 后有一句 TODO: How to handle errors here? 且直接 return ok(LoggingHandler.ts:351-355)——单条正文上传失败不会让整批失败,但也就丢了。
  • 10MB 阈值是粗估。 分流用的 size 来自 tryToGetSize 的粗略估算(LoggingHandler.ts:168-173), 不是精确字节数,临界值附近可能判错档。
  • 超免费额度不存正文。 省钱策略下这类请求点进去看不到正文(PTB 计费除外),:180-181

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

主题文件路径符号名
三路并发落库编排valhalla/jawn/src/lib/handlers/LoggingHandler.tshandleResults / Promise.all
存哪里的三档决策valhalla/jawn/src/lib/handlers/LoggingHandler.ts:179storageLocation / S3_MIN_SIZE_THRESHOLD
映射成 CH 行valhalla/jawn/src/lib/handlers/LoggingHandler.ts:476mapRequestResponseVersionedCH
正文文本三档处理valhalla/jawn/src/lib/handlers/LoggingHandler.ts:641requestResponseTextFromContext
成本定点整数valhalla/jawn/src/lib/handlers/LoggingHandler.ts:519COST_PRECISION_MULTIPLIER
CH 主表建表clickhouse/migrations/schema_41_request_response_replacing_merge_tree.sqlrequest_response_rmt
CH 插入封装valhalla/jawn/src/lib/db/ClickhouseWrapper.ts:46dbInsertClickhouse
CH 主表写入/回填valhalla/jawn/src/lib/stores/request/VersionedRequestStore.tsinsertRequestResponseVersioned / putPropertyIntoClickhouse
缓存指标表clickhouse/migrations/schema_48_cache_metrics.sqlcache_metrics
S3/Minio 客户端valhalla/jawn/src/lib/shared/db/s3Client.tsS3Client.store / getRequestResponseKey
S3 正文回源读取valhalla/jawn/src/lib/stores/request/RequestResponseBodyStore.tsgetRequestResponseBody
Postgres 元数据写入valhalla/jawn/src/lib/stores/LogStore.tsinsertLogBatch / processPromptInputsBatch

上一章: 03 消费与责任链 —— 一条队列消息怎么被加工成这里落库的结构化日志。 下一章: 05 Agent 可观测性:会话追踪、打分与在线评估 —— sessions/scores 的物化视图与评估。 回到: Helicone 全景与导读