跳到主要内容

数据截至 (上游 commit 0261ea4f33d4)

服务端存储:ClickHouse 怎么吞下乱序到达的 trace 与 span

30 秒导读: 上一章讲 SDK 怎么把数据发出去(02-ingestion-pipeline),这一章讲数据落地。核心矛盾只有一句话:一个 span 的「开始」和「结束」是两条独立的 HTTP 请求,它们可能乱序到达;而 ClickHouse 是个列存分析库,根本没有高效的行级 UPDATE。 Opik 的解法是把「更新」也写成 INSERT——用 INSERT INTO … SELECT 新行 LEFT JOIN 旧行,逐字段做三段式合并,再靠 ReplacingMergeTree 收尾去重。

本章路径约定(为了行内引用不至于长得没法读):

  • Java 路径相对 apps/opik-backend/src/main/java/com/comet/opik/,例:domain/SpanDAO.java:174
  • SQL 迁移相对 apps/opik-backend/src/main/resources/liquibase/db-app-analytics/migrations/,例:000001_init_script.sql:24
  • 章末「代码地图」给完整路径。

1. 这一章要解决的问题(零基础也能懂)

先看一个具体场景。 你的 agent 调了一次 LLM,@track 装饰器(见 01-tracing-decorator)在函数进入时造了一个 span,函数返回时才知道输出和耗时。

于是 SDK 发出的不是一条消息,而是两条:

时刻SDK 消息打到哪个后端端点带了什么
函数进入CreateSpanMessagePOST /v1/private/spans/batchid、name、type、start_time、input
函数返回UpdateSpanMessagePATCH /v1/private/spans/{id}end_time、output、usage、cost

依据:sdks/python/src/opik/message_processing/messages.py:157CreateSpanMessage)与 :189UpdateSpanMessage);分发见 sdks/python/src/opik/message_processing/processors/online_message_processor.py:57,62。后端两个端点在 api/resources/v1/priv/SpansResource.java:229(批量创建)与 :268(单条 PATCH)。

问题来了:这两条消息谁先到,不保证。 create 走的是攒批队列,update 走的是单条队列,网络重试、限流退避、多消费者线程都会打乱顺序。所以后端必须能处理三种到达序:

  1. create 先到、update 后到(正常);
  2. update 先到、create 后到(乱序);
  3. 同一条 create 到了两次(SDK 重试)。

而 ClickHouse 帮不上忙。 它是列存分析库,ALTER TABLE … UPDATE 是重写整个数据块的异步 mutation,绝不能拿来做每秒几万次的热路径写入。

一句话直觉: Opik 把数据库当成一本只能往后追加的账本。改一笔账不是去涂改旧行,而是再记一行新的,读的时候只认最新那行。而「新行怎么和旧行合并」这件事,被塞进了 INSERT 语句自己的 SQL 里。


2. 顶层全景(一条 span 从 HTTP 到磁盘)

怎么读这张图: 从左到右是一次写请求的路径;下方两个方框是两套完全不同的存储,各管各的。

SDK / OTLP Dropwizard + Jersey 存储
┌──────────┐ HTTP ┌──────────────────────────────┐
│ 批量上报 │─────────▶│ ① 资源层 SpansResource │
│ PATCH 更新│ │ 限流 @RateLimited → 429 │
└──────────┘ └───────────────┬──────────────┘

┌──────────────────────────────┐ ┌─────────────┐
│ ② 服务层 SpanService │─────▶│ Redis 分布式锁│
│ 按 id 上锁、判断走哪条写路径 │ │ + 限流计数器 │
└───────────────┬──────────────┘ └─────────────┘

┌──────────────────────────────┐ ┌─────────────┐
│ ③ DAO 层 SpanDAO │─────▶│ MySQL 元数据 │
│ 渲染 SQL 模板 + R2DBC 异步 │ │ projects 等 │
└───────────────┬──────────────┘ └─────────────┘

┌──────────────────────────────┐
│ ④ ClickHouse spans / traces │
│ ReplacingMergeTree 后台去重 │
└──────────────────────────────┘

各部件一句话职责:

部件干什么在哪个文件
资源层JAX-RS 端点、鉴权、限流注解、OTLP 反序列化api/resources/v1/priv/SpansResource.javaTracesResource.java
服务层加分布式锁、探测行是否已存在、选择写路径、把 ClickHouse 报错翻译成 409domain/SpanService.javadomain/TraceService.java
DAO 层持有全部 SQL 常量(StringTemplate 模板),用 R2DBC 异步执行domain/SpanDAO.javadomain/TraceDAO.java
过滤下推把 REST 的 filters 查询参数编译成 ClickHouse WHERE 片段api/filter/FiltersFactory.javadomain/filter/FilterQueryBuilder.java
存储装配R2DBC 连接工厂 + ClickHouse v2 客户端的 Guice 绑定infrastructure/db/DatabaseAnalyticsModule.java

两套数据库的分工(这是理解全局的关键):

存储装什么访问方式代表表
MySQL(state DB)低频变更的元数据,需要真 UPDATE 和唯一约束JDBI 同步 SQL 对象projectsdatasetspromptsfeedback_definitionsalerts
ClickHouse(analytics DB)高频海量的观测数据R2DBC 响应式 + SQL 模板tracesspansfeedback_scoresdataset_itemsexperiments
Redis按 span id 的分布式锁、限流令牌桶、流Redisson无表

注意一个容易看错的地方:datasets 在 MySQL,dataset_items 在 ClickHouse——数据集本身是元数据,数据集里的样本是分析数据。DAO 风格也随之不同:domain/ProjectDAO.java:25 用 JDBI 注解式接口,domain/SpanDAO.java 则是手写 SQL 模板 + R2DBC。


3. 数据模型:表引擎与主键的设计意图

3.1 一句话:ORDER BY 同时是三样东西

看初始建表脚本 000001_init_script.sql:24-25

) ENGINE = ReplacingMergeTree(last_updated_at)
ORDER BY (workspace_id, project_id, trace_id, parent_span_id, id);

在 ClickHouse 的 MergeTree 家族里,ORDER BY 这一个声明同时承担三个角色:

  1. 磁盘排序键 —— 数据按这个顺序物理排列;
  2. 稀疏主键索引 —— 查询能按前缀做范围裁剪;
  3. 去重键 —— ReplacingMergeTree 认为「排序键相同的行 = 同一行的多个版本」。

ReplacingMergeTree(last_updated_at) 括号里的 last_updated_at版本列:后台合并时,排序键相同的多行只保留 last_updated_at 最大的那一行。

这就是「用 INSERT 做 UPDATE」能成立的地基:同一个 span 写十次,磁盘上短暂有十行,最终收敛成一行。

3.2 主键前缀为什么长这样

初始脚本里五张核心表的排序键(000001_init_script.sql):

ORDER BY行号
spans(workspace_id, project_id, trace_id, parent_span_id, id)25
traces(workspace_id, project_id, id)42
feedback_scores(workspace_id, project_id, entity_type, entity_id, name)58
dataset_items(workspace_id, dataset_id, source, trace_id, span_id, id)74
experiments(workspace_id, dataset_id, id)85

三条一致的设计规律:

  • workspace_id 永远是第一列。 Opik 是多租户 SaaS,每个查询都必然带 workspace_id = :workspace_id。放在最前面意味着一个租户的数据在磁盘上物理连续,跨租户的数据块直接被主键索引裁掉。
  • 第二列是"容器"project_iddataset_id),继续按最常用的第二维裁剪。
  • 末尾必须是能唯一确定一行的东西。 因为去重键就是主键——如果末尾不唯一,两个不同的 span 会被合并成一个。feedback_scores 的末尾是 (entity_type, entity_id, name):同一个实体上同名的打分只留最新一条,这正是「重新打分」该有的语义。

主键即去重键,是个双刃剑。 好处是零额外成本;代价是主键里的列一旦写错就无法靠更新修复——写进去的是另一行,不是同一行的新版本。第 4 节的哨兵值机制就是为这件事兜底的。

3.3 从 MergeTree 切到 Replicated(000017)

单机版够用,集群版不够——多副本之间需要通过 ZooKeeper/Keeper 协调复制。迁移 000017_change_tables_to_replicated.sql 把所有表换成 Replicated* 引擎:

ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/${ANALYTICS_DB_DATABASE_NAME}/spans', '{replica}', last_updated_at)
ORDER BY (workspace_id, project_id, trace_id, parent_span_id, id)

依据:000017_change_tables_to_replicated.sql:191-192(spans)、:223-224(traces)。语义没变——排序键、版本列一字不改,只是加了副本协调路径。

切换手法值得单独看一眼,因为它是零停机换引擎的标准套路

① CREATE TABLE spans1 (…新引擎…) -- 建影子表
② ALTER TABLE spans1 ATTACH PARTITION -- 把数据分区「挂」过去(不拷贝文件)
tuple() FROM spans
③ ALTER TABLE spans DETACH PARTITION -- 老表卸下分区
④ DROP TABLE spans SYNC -- 删空壳
⑤ RENAME TABLE spans1 TO spans -- 影子表顶上

依据:000017_change_tables_to_replicated.sql:196-201ATTACH PARTITION … FROM 走的是硬链接而非数据拷贝,所以 TB 级表也能秒级完成。

3.4 读的时候怎么办:不用 FINAL,手写去重

ReplacingMergeTree 的去重是后台合并时发生的,不保证及时。所以刚写完立刻查,可能读到同一个 id 的多行。

ClickHouse 提供了 FINAL 关键字强制查询时去重,但它很慢。Opik 只在小表上用(domain/SpanDAO.java:1637feedback_scores FINAL),核心的 spans/traces 查询一律手写去重

ORDER BY (workspace_id, project_id, trace_id, parent_span_id, id) DESC, last_updated_at DESC
LIMIT 1 BY id

依据:domain/SpanDAO.java:1078-1079(分页查询的 spans_deduped CTE)、:806(按 id 单查)、:785-786SELECT_PARTIAL_BY_ID)。

LIMIT 1 BY id 是 ClickHouse 特有语法:按 id 分组,每组按当前 ORDER BY 取第一行。配合 last_updated_at DESC,效果就是"每个 id 只留最新版本"。这段模式在两个 DAO 里出现了十几次,是全代码库最高频的惯用法。


4. 核心机制:用 INSERT 做 UPDATE

这是本章的主菜。

4.1 思路:把合并逻辑写进 INSERT 语句

要解决的小问题: 新来一条 PATCH(只带 end_time 和 output),数据库里可能已经有一行(带 name、input),也可能什么都没有。怎么在一条语句里搞定?

直觉: 既然要 INSERT 一行完整的数据,那就让这一行的每个字段值由「新数据 + 旧数据」现场算出来。SQL 表达式完全够用:

-- 示意,非源码:一条语句同时完成「读旧行 + 合并 + 写新行」
INSERT INTO spans (id, name, end_time,)
SELECT
new.id,
if(old.name != '', old.name, new.name) AS name, -- 旧的有就用旧的
if(old.end_time IS NOT NULL, old.end_time, new.end_time) AS end_time
FROM (SELECT :id AS id, :name AS name,) AS new -- 新数据来自绑定参数
LEFT JOIN (SELECT * FROM spans WHERE id = :id LIMIT 1) AS old
ON new.id = old.id

LEFT JOIN 是关键:旧行不存在时,old 的所有列取类型默认值(String 取 '',Nullable 取 NULL),条件自然全部落到「用新值」,无需分支。一条 SQL 同时覆盖了 insert 和 update 两种情况。

4.2 真实实现:multiIf 三段式

真源码在 domain/SpanDAO.java:174-352(常量 INSERT),domain/TraceDAO.java:283-427 是同构写法。核心是每个字段一段 multiIf

multiIf(
LENGTH(CAST(old_span.trace_id AS Nullable(String))) > 0
AND notEquals(old_span.trace_id, new_span.trace_id), leftPad('', 40, '*'),
LENGTH(CAST(old_span.trace_id AS Nullable(String))) > 0, old_span.trace_id,
new_span.trace_id
) as trace_id,

依据:domain/SpanDAO.java:213-217。三段的含义:

┌─ 旧值存在 且 与新值冲突 ──▶ 写入「哨兵值」(40 个星号,故意撑爆列宽)
字段 ──┼─ 旧值存在(不冲突)──────▶ 保留旧值
└─ 其余(旧行不存在)───────▶ 采用新值

但不是每个字段都有第一段。 只有身份类字段才带冲突哨兵:

字段有哨兵分支?合并规则行号
project_id冲突 → 哨兵;否则旧值优先SpanDAO.java:207-211
trace_id同上:208-212
parent_span_id同上(哨兵额外 CAST 成 FixedString(19):213-217
name / input / output / metadata / model旧值非空则保留旧值,否则取新值:218-245
end_time / ttft旧值非 NULL 则保留,否则取新值:230-233:292-295
type / source旧值不是 'unknown' 则保留:222-225:296-299
tags / usage旧值非空集合则保留:262-269
last_updated_by / truncation_threshold无条件取新值:282-283

为什么身份字段和内容字段规则不同? 身份字段(这个 span 属于哪个 trace、哪个项目、哪个父 span)在 span 生命周期内是不变量。旧值和新值不一致,只有一个可能:客户端传错了,或者两个不同的 span 撞了 id。这不该悄悄"合并",该报错。而内容字段本来就是逐步补齐的,旧值优先只是防止后到的 POST 重试覆盖掉已经写好的 PATCH 结果。

4.3 哨兵值 leftPad('', 40, '*') 到底干了什么

这是全代码库最反直觉的一处设计,值得说清楚。

leftPad('', 40, '*') 生成一个 40 个字符的字符串。而 trace_idproject_id 的列类型是 FixedString(36)000001_init_script.sql:10-11)。40 > 36。

所以这个值根本写不进去——ClickHouse 会抛 TOO_LARGE_STRING_SIZE,整条 INSERT 失败。

这不是 bug,是故意的。服务层守在异常出口,按错误消息里的列名把它翻译成业务错误:

if (ex instanceof ClickHouseException
&& ex.getMessage().contains("TOO_LARGE_STRING_SIZE")
&& ex.getMessage().contains("_CAST(trace_id, FixedString(36))")) {
return failWithConflict(TRACE_ID_MISMATCH);
}

依据:domain/SpanService.java:316-327handleSpanDBError),parent_span_id 分支特意匹配 CAST(leftPad(FixedString(19)。三种冲突分别映射成 PROJECT_AND_WORKSPACE_NAME_MISMATCHPARENT_SPAN_IS_MISMATCHTRACE_ID_MISMATCH,最终返回 HTTP 409。

妙在哪: ClickHouse 没有约束、没有触发器、没有 RAISE。要在一条纯 SQL 里表达"这种情况必须失败",就得借用类型系统当断言——用一个类型装不下的值做毒丸。代价是错误信息靠字符串匹配识别,很脆(源码里 //TODO: refactor to implement proper conflict resolution 出现在 SpanDAO.java:357:431SpanService.java:250,作者自己也标了)。

4.4 三种到达顺序,走三条不同的 SQL

服务层先探测数据库现状,再决定用哪条 SQL。怎么读这张图:从上往下是判断顺序。

POST /spans (create) PATCH /spans/{id} (update)
│ │
▼ ▼
getPartialById(id) 只查 start_time getOnlySpanDataById(id)
│ │
┌────┴─────┐ ┌────┴──────┐
│ │ │ │
没查到 查到了 查到了 没查到
│ │ │ │
│ start_time 是 epoch? │ │
│ │ │ │ │
│ 是 否 │ │
│ │ │ │ │
▼ ▼ ▼ ▼ ▼
INSERT INSERT 直接返回 id UPDATE PARTIAL_INSERT
(幂等,忽略重复 POST)

依据:domain/SpanService.java:189-205insertSpan)与 :202-228update)。

start_time == epoch 这个判断是整套机制的暗号。 PATCH 先到、行不存在时走的 PARTIAL_INSERT,它在新行里把 start_time 硬编码成 1970 纪元:

toDateTime64('1970-01-01 00:00:00.000', 9) as start_time,

依据:domain/SpanDAO.java:572。这行"半成品"因此被打上了可识别的标记。后到的 POST 读到 epoch,就知道"这只是个 PATCH 留下的壳,我该把真数据补进去";读到真实时间,就知道"完整的 span 已经在了,我是重试,直接返回"。

四条写 SQL 的分工:

SQL 常量触发场景合并优先级行号
BULK_INSERT批量创建端点不合并,直接写新行SpanDAO.java:100-167
INSERT单条创建 / 补全 PATCH 留下的壳旧值优先(+ 身份字段哨兵):169-347
UPDATEPATCH 且行已存在逐字段 <if(x)> :x <else> x <endif>,传了就换:353-420
PARTIAL_INSERTPATCH 但行不存在(乱序)新值优先,其次旧值(+ 哨兵):432-600

对比 INSERTPARTIAL_INSERT 的同一个字段就能看出反转:

-- INSERT(POST 补全):旧值优先,别覆盖 PATCH 已写好的
multiIf(LENGTH(old_span.input) > 0, old_span.input, new_span.input) as input -- :234-237

-- PARTIAL_INSERT(PATCH 先到):新值优先,PATCH 带来的就是最新的
multiIf(LENGTH(new_span.input) > 0, new_span.input,
LENGTH(old_span.input) > 0, old_span.input,
new_span.input) as input -- :475-479

两条路径的优先级正好相反,合起来才让乱序收敛到同一个结果。

4.5 快路径:BULK_INSERT 不做任何合并

批量端点是吞吐主战场,它走的是完全不同的 SQL:

INSERT INTO spans( id, project_id,)
SETTINGS log_comment = '<log_comment>'
FORMAT Values
<items:{item | ( :id<item.index>, :project_id<item.index>,) <if(item.hasNext)>,<endif>}>

依据:domain/SpanDAO.java:100-167。三点差异:

  1. 没有 LEFT JOIN、没有 multiIf。 纯粹的多值 INSERT,一次写 N 行,靠 ReplacingMergeTree 事后去重。
  2. FORMAT Values 是 ClickHouse 的写入快路径,避免服务端解析成完整 SELECT 计划。源码注释特意警告:last_updated_at 必须在客户端格式化成字符串字面量,若在 SQL 里写 now64(6) 函数调用会破坏这条快路径(domain/SpanDAO.java:1876-1881,OPIK-5694);同理 cost 必须绑 BigDecimal 本身而非其字符串形式(:1770-1774)。
  3. last_updated_at 显式绑定,且整批共用一个时间戳Instant nowForBatch:1726),保证同批次内版本时间一致。

对比之下,INSERT / UPDATE / PARTIAL_INSERT 的列清单里根本没有 last_updated_at——交给列默认值 DEFAULT now64(6)(迁移 000025_reduce_span_last_updated_at_to_micros.sql:5)。写入时刻自动成为版本号,后写的天然赢。

代价与边界: 快路径没有字段级合并,所以同一个 id 在同一批里出现两次,只有最后写入那次的全部字段存活——没传的字段不会从旧行"继承"过来。服务层因此在进批之前先做了应用层去重(domain/SpanService.java:394dedupSpans)。


5. 查询侧:怎么把宽表查快

写完只是一半。spans 表在生产里是 TB 量级(迁移注释直接写了 ~19.5TB/1.42B rows000082_materialize_spans_skip_indexes.sql:3),查询侧有一整套针对性优化。

5.1 宽列瘦身:一个字段派生出五个

input / output 存的是完整 JSON,可能几百 KB。列表页只需要预览,不需要全文。于是围绕这两列陆续长出一堆派生列:

迁移加了什么类型解决什么
000024_add_output_input_pre_computed_columns_to_spans_and_traces.sql:4-12input_length / output_length / metadata_lengthMATERIALIZED length(input)前端要显示"是否被截断",不必读原文就能判断
000041_add_truncated_input_output_to_spans_and_traces_tables.sql:4-13truncation_threshold(默认 10001)+ truncated_input / truncated_outputMATERIALIZED if(length(x) >= 阈值, substring(x,1,阈值), x)列表页只读前 10KB,磁盘 IO 直接砍数量级
000054_add_slim_input_output_columns.sql:5-11input_slim / output_slim普通 String 列,服务端写入时算好保结构的 JSON 瘦身(截断不破坏 JSON 语法)

MATERIALIZED 列的好处:写入时不用传,查询时不用算,ClickHouse 在插入时计算并按列存储,读的时候和普通列一样便宜。

5.2 两段式分页:先查窄列定页,再回表取宽列

分页查询 SELECT_BY_PROJECT_IDdomain/SpanDAO.java:868)的骨架,源码自带的说明在 :810-820

① spans_deduped 窄列扫描:把 input/output/metadata 从 SELECT 里剔掉
(:990) 做过滤 + LIMIT 1 BY id 去重


② page_ids 只留排序键和 id,ORDER BY + LIMIT/OFFSET → 拿到这一页的 20 个 id
(:1024)


③ page_wide WHERE id IN (page_ids):只为这 20 行读宽列
(:1036)


④ 最终 SELECT LEFT JOIN comments / feedback_scores,套截断正则

**核心收益:**扫描百万行时不碰宽列,只有最终那一页(20 行)才付宽列的 IO。

一个例外由 sort_needs_wide 标志控制:如果用户就是要按 input 排序,那阶段 ① 不能丢掉 input,否则排不出来(infrastructure/FilterUtils.java:44addSortNeedsWideFlag,用在 SpanDAO.java:1045)。

5.3 Skip index:给非主键列的粗筛

主键索引只对 ORDER BY 前缀有效。想按 thread_idsourcecreated_at 过滤,就得靠 skip index(数据跳过索引:为每个数据块记摘要,查询时整块跳过)。

索引类型迁移
traces.thread_idbloom_filter(0.01) —— 等值查询000077_add_thread_id_skip_index_to_traces.sql:5-6
spans.source / traces.source补做 MATERIALIZE INDEX000082_materialize_spans_skip_indexes.sql:5-7000081_materialize_traces_skip_indexes.sql
spans/traces.created_atlast_updated_atminmax GRANULARITY 1 —— 范围查询000088_add_created_at_last_updated_at_skip_indexes_to_spans_and_traces.sql:4-15

两个运维细节值得抄走:

  • 加索引 ≠ 老数据能用。 ADD INDEX 只对新写入的数据块生效,存量数据必须 MATERIALIZE INDEX 才补上。000080/000081/000082 三个迁移就是在补做前面几次忘了 materialize 的索引,迁移注释里逐条写明"materialize index from 000075"。
  • materialize 是异步 mutation。 迁移注释直接给了监控 SQL:SELECT * FROM system.mutations WHERE is_done = 0 AND table = 'spans'000082:3),并要求"等 000081 的 mutation 跑完再上 000082"——TB 级表上并发 mutation 会把集群压垮。
  • minmax 之所以对时间列有效,是因为注释里点破的前提:created_at / last_updated_at 与插入顺序高度相关(000088:3),所以每个块的 min/max 区间窄、裁剪率高。换成随机分布的列,minmax 就没用了。

5.4 过滤下推链路:从 URL 参数到 WHERE 片段

REST 层的 filters 是个 JSON 字符串。它要变成 SQL,走三个环节:

?filters=[{"field":"trace_id","operator":"=","value":"abc"}]


① FiltersFactory.newFilters() 反序列化 → 去重 → URL 解码 → 校验
api/filter/FiltersFactory.java:120 校验失败 → 400 BadRequest


② Field 枚举 → FieldType 每个可过滤字段声明自己的类型
api/filter/SpanField.java:9 trace_id 是 STRING_EXACT,name 是 STRING


③ FilterQueryBuilder.toAnalyticsDbFilters() (操作符 × 类型) → SQL 模板
domain/filter/FilterQueryBuilder.java:864 拼成 "(a AND b AND c)" 塞进模板

三环节各自的关键点:

  • **FieldTypeapi/filter/FieldType.java:12-36)**是这条链路的枢纽——18 种类型,每种决定值怎么校验、怎么解码、生成什么 SQL。
  • **校验表 FIELD_TYPE_VALIDATION_MAPFiltersFactory.java:40-116)**为每个类型挂一个断言。加了新类型忘了在这里登记,validateFieldType 直接 NPE。
  • 解码豁免(FiltersFactory.java:151-164:STRING / STRING_EXACT / STRING_LIST / ENUM 跳过 URL 解码,因为 JSON 反序列化时已经解过一次,再解会把值里的 +% 弄坏。
  • **模板表 ANALYTICS_DB_OPERATOR_MAPFilterQueryBuilder.java:177)**是二维的:Operator → FieldType → SQL 模板。查不到组合就返回 null,FiltersFactory.java:166-170 把它变成 400。

5.5 ID 字段必须用 STRING_EXACT——一条性能红线

同一个 = 操作符,两种字符串类型生成的 SQL 完全不同:

FieldType= 生成的 SQLcontains 生成的 SQL能吃主键索引?
STRINGlower(%1$s) = lower(:filter%2$d)ilike(%1$s, …)不能
STRING_EXACT%1$s = :filter%2$d%1$s LIKE …

依据:domain/filter/FilterQueryBuilder.java:231-232(EQUAL)、:147-148(CONTAINS)。

为什么 lower() 是灾难: ClickHouse 的主键索引建在列的原始值上。一旦写成 lower(trace_id) = lower(?),索引条件变成了对函数结果的判断,主键裁剪彻底失效——TB 级表全表扫描。

所以 api/filter/SpanField.java 里,ID:9)和 TRACE_ID:33)都声明为 FieldType.STRING_EXACT,而 NAME:10)、INPUT:15)、MODEL:20)这类人写的文本才用 STRING

判断准则一句话:UUID、外键、系统生成的标识符 → STRING_EXACT;人类可读的文本 → STRING 前者大小写敏感本来就是正确语义,后者用户希望忽略大小写。

5.6 顺带一个巧思:时间过滤变成主键范围扫描

用户按"最近一小时"过滤 span。start_time 不在主键里,怎么快?

答案:Opik 的 id 是 UUIDv7,前 48 位就是毫秒时间戳domain/IdGenerator.java:90Generators.timeBasedEpochGenerator(),且 validateVersion 强制 v7)。所以时间区间可以被翻译成 id 区间

// api/InstantToUUIDMapper.java:38 toLowerBound —— 造出该毫秒下字典序最小的 UUIDv7
long msb = (epochMilli << 16) | 0x7000L; // 时间戳 + 版本号 7,随机位全 0
long lsb = 0x8000000000000000L; // 变体位 10,其余 62 随机位全 0

配对的 toUpperBound:79)把随机位全置 1。于是查询条件变成 AND id >= :uuid_from_time AND id <= :uuid_to_timedomain/SpanDAO.java:1054-1057)——主键最后一列上的范围扫描,直接吃索引(同一处还并排加了一条 toMonday(id_at) 周一对齐条件,防跨周桶的边界串味)。绑定点在 api/resources/v1/priv/SpansResource.java:157


6. 请求侧概览(这一层这章只做地图)

6.1 端点分布

端点方法说明位置
/v1/private/spansPOST单条创建SpansResource.java:204
/v1/private/spans/batchPOST批量创建(SDK 主力):227
/v1/private/spans/{id}PATCH单条更新(end/output):268
/v1/private/spans/batchPATCH批量更新(按 id 集合):245
/v1/private/otel/v1/tracesPOSTOTLP 直接入口OpenTelemetryResource.java:35,43

TracesResource.java 是同构的一套。

6.2 OTLP 直连:不装 Opik SDK 也能上报

Opik 直接接受 OpenTelemetry 的原生协议,两种编码各有一个 JAX-RS MessageBodyReader

  • infrastructure/otel/OtelProtobufMessageBodyReader.java:18 —— @Consumes("application/x-protobuf"),用反射找 protobuf 类的 parseFrom(InputStream) 静态方法(:31-33)。
  • infrastructure/otel/OtelJsonMessageBodyReader.java —— 同一份 ExportTraceServiceRequest,走 JSON 编码。

两个端点最终汇到同一个 handleOtelTraceRequestOpenTelemetryResource.java:51),项目名从请求头取,落进 OpenTelemetryService.parseAndStoreSpans

意义: 任何已经接了 OTel 的服务,改个 endpoint 就能把数据打进 Opik,不需要碰 Opik SDK。

6.3 限流与 SDK 的 429 重试,是一对协议

后端侧是一个 Guice 方法拦截器:

  • @RateLimited 注解(infrastructure/ratelimit/RateLimited.java)声明桶名,支持 {workspaceId} / {apiKey} / {clientIp} 占位符。
  • 拦截器算令牌(Redis),超限就抛 ClientErrorException(…, SC_TOO_MANY_REQUESTS)infrastructure/ratelimit/RateLimitInterceptor.java:196-200),并在抛之前把剩余额度和重置时间写进响应头(setLimitHeaders:203)。
  • RateLimitResponseFilter.java:26-36 负责把这些头搬到最终响应上。

SDK 侧接住这个约定:读 429 响应头算出 retry_after,抛成 OpikCloudRequestsRateLimited 让上游退避重试(sdks/python/src/opik/message_processing/processors/online_message_processor.py:149-156)。同一个处理器对 409 直接静默吞掉:132-137),注释写明"重试机制有时会把同一请求发两次,第二次被后端拒绝不该让用户看到报错"——这正好接上第 4.3 节哨兵值产出的那个 409。

读端点的限流是分开计费的getSpansByProject 用独立桶且 shouldAffectWorkspaceLimit = falseSpansResource.java:115),查询不占用写入配额。

6.4 连接层

infrastructure/db/DatabaseAnalyticsModule.java 装配三样东西:

  • R2DBC ConnectionFactory:31-33)—— 响应式驱动,全部 DAO 的执行通道,外面套了 OTel 的 R2dbcTelemetry 埋点。
  • ClickHouse v2 Client:34)—— 少数需要新客户端能力的场景。
  • 一个只读客户端:43buildReadOnlyFreeFormSqlClient)—— 专供自由 SQL 功能,跑在 readonly=1 的受限用户下。

另一个可借鉴的小件:infrastructure/db/ZeroRowsRetryPolicy.java:43retryOnZeroRows。它针对的是 INSERT … SELECT 在复制表上"HTTP 层成功但写了 0 行"的业务级失败——驱动内建的重试只认传输层故障(类注释 :13-20 说明了这个区别)。目前的使用者是 domain/DatasetItemVersionDAO.java:2625,spans/traces 路径尚未接入。


7. 巧妙之处(可以带走的技术)

  1. 主键 = 排序键 = 去重键,三合一。ReplacingMergeTree(版本列) + ORDER BY (租户, 容器, …, 唯一 id),多租户裁剪和幂等写入一次性解决,零额外索引成本(000001_init_script.sql:24-25)。

  2. 把「合并」写进 INSERT 语句本身。 INSERT … SELECT new LEFT JOIN old + 每字段 multiIf,一条 SQL 覆盖 insert/update 两种情况,且不需要事务(domain/SpanDAO.java:174-352)。

  3. 两条相反优先级的写路径,让乱序收敛。 POST 路径旧值优先、PATCH-先到路径新值优先,任意到达序都得到同一结果(:169 vs :432)。

  4. 借类型系统当断言。 leftPad('', 40, '*') 写不进 FixedString(36),用溢出制造受控失败,再在服务层翻译成 409(:203 + domain/SpanService.java:316-327)。

  5. 用 epoch 时间戳当"半成品"标记。 PATCH 先到时写 start_time = 1970-01-01,后到的 POST 靠这个标记区分"补全"和"重复重试"(domain/SpanDAO.java:572 + SpanService.java:199)。

  6. 两段式分页:窄列定页、宽列回表。 扫描阶段丢掉 input/output,只为最终一页付宽列 IO(domain/SpanDAO.java:1043-1111)。

  7. UUIDv7 让时间过滤吃上主键索引。 时间边界翻译成 id 边界,非主键列的范围查询变成主键前缀扫描(api/InstantToUUIDMapper.java:38,77)。

  8. 零停机换表引擎:ATTACH PARTITION 硬链接搬家。 建影子表 → 挂分区 → 改名,TB 级表秒级切换(000017_change_tables_to_replicated.sql:196-201)。


8. 边界与局限(诚实版)

  • 冲突检测靠字符串匹配异常消息。 handleSpanDBErrordomain/SpanService.java:309-330)匹配的是 "_CAST(trace_id, FixedString(36))" 这类 ClickHouse 内部生成的表达式文本。ClickHouse 升级改了错误措辞,409 就会退化成 500。源码里三处 //TODO: refactor to implement proper conflict resolutionSpanDAO.java:357:431SpanService.java:250)说明作者知道这是权宜之计。

  • 写路径依赖 Redis 分布式锁串行化。 SpanService.java:184:217 都把写操作包在 lockService.executeWithLock(new LockService.Lock(id, SPAN_KEY), …) 里。锁丢了(Redis 抖动、锁超时)就可能出现两个并发写互相覆盖——multiIf 只能防语义冲突,防不了竞态。

  • 批量快路径没有字段级合并。 BULK_INSERT 直接写新行,同批内重复 id 只有最后一条的全部字段存活。应用层的 dedupSpansSpanService.java:394)是唯一防线。

  • 主键字段一旦写错就是新行。 因为去重键就是主键。这也是为什么必须靠哨兵值在写入时拦下冲突——事后无法修复。

  • MATERIALIZE INDEX 是重活。 迁移注释自己写明 spans 表约 19.5TB / 14.2 亿行(000082:3),materialize 是长时间后台 mutation,必须串行执行并盯 system.mutations

  • 本章不覆盖在线评估。 trace 落库后触发的规则引擎、Python 沙箱、prompt 自动优化,见 06-online-scoring-and-optimizer;把 dataset × task × metric 跑成 experiment 的部分见 04-evaluation-engine


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

主题文件符号名
初始表结构与表引擎apps/opik-backend/src/main/resources/liquibase/db-app-analytics/migrations/000001_init_script.sqlspans / traces / feedback_scores / dataset_items / experiments
切换复制表引擎.../migrations/000017_change_tables_to_replicated.sqlReplicatedReplacingMergeTree / ATTACH PARTITION
长度预计算列.../migrations/000024_add_output_input_pre_computed_columns_to_spans_and_traces.sqlinput_length / output_length
截断列.../migrations/000041_add_truncated_input_output_to_spans_and_traces_tables.sqltruncation_threshold / truncated_input
保结构瘦身列.../migrations/000054_add_slim_input_output_columns.sqlinput_slim / output_slim
skip index 物化.../migrations/000082_materialize_spans_skip_indexes.sqlidx_spans_source
时间列 skip index.../migrations/000088_add_created_at_last_updated_at_skip_indexes_to_spans_and_traces.sqlidx_spans_created_at / idx_traces_last_updated_at
「用 INSERT 做 UPDATE」主实现apps/opik-backend/src/main/java/com/comet/opik/domain/SpanDAO.javaINSERT / PARTIAL_INSERT / UPDATE
批量写快路径同上BULK_INSERT / batchInsert
两段式分页查询同上SELECT_BY_PROJECT_ID(CTE spans_deduped / page_ids / page_wide
trace 侧同构实现apps/opik-backend/src/main/java/com/comet/opik/domain/TraceDAO.javaINSERT / INSERT_UPDATE / BATCH_INSERT
写路径分流与冲突翻译apps/opik-backend/src/main/java/com/comet/opik/domain/SpanService.javainsertSpan / insertUpdate / handleSpanDBError / updateOrFail
过滤类型枚举apps/opik-backend/src/main/java/com/comet/opik/api/filter/FieldType.javaSTRING / STRING_EXACT / buildFilter
过滤校验与解码apps/opik-backend/src/main/java/com/comet/opik/api/filter/FiltersFactory.javaFIELD_TYPE_VALIDATION_MAP / toValidAndDecoded
过滤 SQL 生成apps/opik-backend/src/main/java/com/comet/opik/domain/filter/FilterQueryBuilder.javaANALYTICS_DB_OPERATOR_MAP / toAnalyticsDbFilters
span 可过滤字段声明apps/opik-backend/src/main/java/com/comet/opik/api/filter/SpanField.javaID / TRACE_ID(均为 STRING_EXACT
时间 → UUID 边界apps/opik-backend/src/main/java/com/comet/opik/api/InstantToUUIDMapper.javatoLowerBound / toUpperBound
UUIDv7 生成与校验apps/opik-backend/src/main/java/com/comet/opik/domain/IdGenerator.javagenerateId / validateVersion
REST 端点apps/opik-backend/src/main/java/com/comet/opik/api/resources/v1/priv/SpansResource.javacreateSpans / batchUpdate / getSpansByProject
OTLP 入口apps/opik-backend/src/main/java/com/comet/opik/api/resources/v1/priv/OpenTelemetryResource.javareceiveProtobufTraces / receiveJsonTraces
OTLP 反序列化apps/opik-backend/src/main/java/com/comet/opik/infrastructure/otel/OtelProtobufMessageBodyReader.javareadFrom
限流拦截与 429apps/opik-backend/src/main/java/com/comet/opik/infrastructure/ratelimit/RateLimitInterceptor.javaverifyRateLimit / setLimitHeaders
SDK 侧 429/409 处理sdks/python/src/opik/message_processing/processors/online_message_processor.pyOpikCloudRequestsRateLimited
存储连接装配apps/opik-backend/src/main/java/com/comet/opik/infrastructure/db/DatabaseAnalyticsModule.javagetConnectionFactory / buildReadOnlyFreeFormSqlClient
0 行写入重试apps/opik-backend/src/main/java/com/comet/opik/infrastructure/db/ZeroRowsRetryPolicy.javaretryOnZeroRows
MySQL 元数据 DAO 风格apps/opik-backend/src/main/java/com/comet/opik/domain/ProjectDAO.java@RegisterConstructorMapper / JDBI SQL 对象