Langfuse ClickHouse 写入最佳实践:INSERT 批量大小(10K–100K 行)的原理与落地
本篇基于 Langfuse 仓库内置的 ClickHouse 最佳实践规则库(.agents/skills/clickhouse-best-practices/)中影响级别为 CRITICAL 的 insert-batch-size 规则展开,讲解“为什么每条 INSERT 都会产生一个新数据 part”这一底层机制、10K–100K 行/次 INSERT 的批量标准及其背后的 part/merge 权衡,并给出可复制的 Python 示例、批量大小决策表和 part 数量监控 SQL,最后结合 Langfuse 仓库自身的 ClickHouse 写入代码展示该规则在真实 LLM 可观测性平台中的落地方式。读完本篇,你可以为自己的 ClickHouse 写入链路做批量大小审计,并用 system.parts 监控查询验证批量化改造的效果。
1. 规则定位:insert-batch-size 在规则库中的地位
该规则文档位于 rules/insert-batch-size.md,其 YAML frontmatter 声明了规则的元信息:
| 字段 | 值 |
|---|---|
| title | Batch Inserts Appropriately (10K-100K rows) |
| impact | CRITICAL |
| impactDescription | Each INSERT creates a part; single-row inserts overwhelm merge process |
| tags | insert, batching, parts, performance |
在技能入口 SKILL.md 的“Rule Categories by Priority”优先级表中,Insert Batching(insert-batch- 前缀)被列为 CRITICAL 级别,并规定了 Insert Strategy 类问题的审查流程:任何涉及数据摄入(data ingestion)的方案,必须先读取 rules/insert-batch-size.md,再依次检查 mutation 规避、async insert 等相邻规则。其审查清单中第一项就是:
[ ] Batch size 10K-100K rows per INSERT
也就是说,在 Langfuse 仓库的约定中,“每次 INSERT 承载 1 万到 10 万行”不是建议值,而是写入链路评审的硬性检查点。
2. 底层机制:为什么小批量 INSERT 会压垮集群
规则文档给出的核心解释只有一句话,但它揭示了 ClickHouse MergeTree 系引擎(MergeTree、ReplacingMergeTree、CollapsingMergeTree 等)最关键的一条写入语义:
Each INSERT creates a new data part. Single-row or small-batch inserts create thousands of tiny parts, overwhelming the merge process and causing cluster instability.
从这条语义可以推断出完整的因果链:
- 一次 INSERT = 一个新 part:无论这条 INSERT 携带 1 行还是 10 万行,落盘后都是一个独立的 data part(目录化的列式数据块)。
- part 数量爆炸:逐行插入 10,000 条事件就产生 10,000 个 part;每个 part 100 行的小批量插入 100 万次同样产生 100 万个 part。
- 后台 merge 追不上:MergeTree 依赖后台 merge 进程将小 part 合并成大 part;part 过多时 merge 队列积压,磁盘上出现大量碎片化的小文件,读取时每次查询都要打开、扫描更多 part(稀疏索引的命中粒度也随之劣化)。
- 分区级 part 上限被触发:ClickHouse 对每个分区内的活跃 part 数量有上限(
max_parts_in_total),超过后 INSERT 会被直接拒绝——这正是规则文档中“>3000 per partition blocks inserts”所指的写入失败场景,也是集群“instability”(不稳定)的直接来源。
因此批量化的本质是用更少的 part 承载同样的行数:10 万行数据按 100 行/批要 1000 个 part,按 1 万行/批只需 10 个 part,part 数量相差两个数量级,merge 压力与查询扫描的 part 数同步下降。
3. 反模式与正确写法(完整示例)
规则文档给出了两组对比示例,下面完整保留并补充注释。
3.1 反模式:逐行或极小批量插入
# Single-row inserts - creates 10,000 parts!
for event in events:
client.execute("INSERT INTO events VALUES", [event])
# Tiny batches - still too many parts
for batch in chunks(events, 100): # 100 rows per INSERT
client.execute("INSERT INTO events VALUES", batch)
- 第一段:10,000 次逐行 INSERT 直接产生 10,000 个 part,是最典型的错误。
- 第二段:把批次提升到 100 行看似有改进,但 100 行仍属于“tiny batch”——1,000 万行数据仍会产生 10 万个 part,远离目标区间。
3.2 正确写法:1 万~10 万行/次
# Ideal batch size: 10,000-100,000 rows
BATCH_SIZE = 10_000
for batch in chunks(events, BATCH_SIZE):
client.execute("INSERT INTO events VALUES", batch)
chunks(events, BATCH_SIZE) 将事件流按 1 万行切块,每次 client.execute 提交一整批。这个写法同样适用于 Node.js 的 @clickhouse/client:调用 client.insert({ table, values, format }) 时把整批记录放入 values 一次提交,而不是在循环里逐条调用。
3.3 推荐批量大小阈值
规则文档给出的决策表如下,可直接用于写入链路评审:
| 阈值 | 取值 | 说明 |
|---|---|---|
| 最小批量(Minimum) | 1,000 行/次 | 低于此值的同步 INSERT 已被视为不健康,应改用 async insert(见第 5 节) |
| 理想区间(Ideal range) | 10,000–100,000 行/次 | 覆盖大多数 LLM tracing/eval 场景的吞吐 |
| 同步写入速率(Insert rate, sync) | 约 1 次 INSERT/秒 | 同步 INSERT 的合理节奏;更高频的小写入应交给服务端 async 缓冲 |
三个阈值互为印证:以 1 次/秒的速率插入 1 万行/批,每天产生的 part 数为 86,400 个,分散在多分区后可被后台 merge 平稳消化;而同样流量若逐行插入,part 数放大 1 万倍。
4. 验证手段:用 system.parts 监控 part 数量
规则文档提供了标准化的验证 SQL,用于观察每张表的活跃 part 数与总行数:
-- Monitor part count (>3000 per partition blocks inserts)
SELECT table, count() as parts, sum(rows) as total_rows
FROM system.parts
WHERE active AND database = 'default'
GROUP BY table
ORDER BY parts DESC;
使用要点:
active = 1过滤掉正在被合并而标记为active = 0的 part,反映真实占用;database条件需替换为实际数据库名(Langfuse 的部署中表通常位于独立的 ClickHouse 数据库,而不是示例中的default);- 关注按表聚合后的
parts排序:某张高频写入表的 part 数长期居高,即是批量过小或分区粒度过细的信号; - 该查询适合作为批量化改造前后的对照指标:改造成功意味着同样的行数下
parts显著下降。
5. 关联规则:客户端无法批量时改用 async insert
insert-batch-size 不是孤立规则。同一规则库中 rules/insert-async-small-batches.md(HIGH 级)给出了客户端实在无法攒批时的兜底方案:打开服务端 async insert,让 ClickHouse 自己把小批缓冲成大 part:
# Enable async_insert with safe defaults
client.execute("SET async_insert = 1")
client.execute("SET wait_for_async_insert = 1") # Confirms durability
for batch in chunks(events, 100):
client.execute("INSERT INTO events VALUES", batch)
# Server buffers and creates larger parts automatically
-- Configure server-side for specific users
ALTER USER my_app_user SETTINGS
async_insert = 1,
wait_for_async_insert = 1,
async_insert_max_data_size = 10000000, -- Flush at 10MB
async_insert_busy_timeout_ms = 1000; -- Flush after 1s
两条规则的边界是明确的:首选客户端批量(10K–100K 行/次);只有当业务语义决定客户端不能攒批(如实时性要求极高的单条事件流)时,才启用服务端 async insert,且 flush 条件取“先到者”——缓冲达到 async_insert_max_data_size、超过 async_insert_busy_timeout_ms 或累计 INSERT 查询数达到上限。返回模式上 wait_for_async_insert = 1(等待 flush 确认持久化)是推荐配置;= 0 是 fire-and-forget,存在数据丢失风险,仅在接受丢数的前提下使用。
6. Langfuse 仓库中的实际落地证据
规则库只是“纸面约束”,Langfuse 仓库源码中有几处可以直接印证该规则的落地方式。
6.1 upsertClickhouse:一次 INSERT 携带整批记录
核心写入路径在 clickhouse.ts。该函数接收 records: T[](整批记录数组),映射出全部行后只发起一次 clickhouseClient().insert(...) 调用:
const res = await clickhouseClient().insert({
table: opts.table,
values: opts.records.map((record) => ({
...record,
event_ts: convertDateToClickhouseDateTime(new Date()),
})),
format: "JSONEachRow",
clickhouse_settings: {
log_comment: buildClickHouseLogComment(opts.tags),
},
});
可以看到三个与规则直接相关的实现细节:
- 整批一次性提交:
values是整批记录的数组,批次大小由上游调用方(ingestion 处理管线)决定,而不是按行循环调用 insert——这正是insert-batch-size所要求的形态; - 使用
JSONEachRow格式:与规则库 rules/insert-format-native.md 中“使用二进制高效格式”的精神一致; - 查询归因标签:
log_comment通过buildClickHouseLogComment注入 surface/route/projectId 等标签,便于事后在system.query_log中按来源(trpc、publicapi、worker、mcp 等)审计写入行为——批量化改造前后对比时,可以用它区分不同入口的 INSERT 特征。
6.2 回补任务的显式批次大小
Worker 的后台迁移任务把批次大小做成显式参数,例如 backfillValidToForDatasetItems.ts 定义了 DEFAULT_BATCH_SIZE = 1000,并通过 args.batchSize 允许运行时覆盖。从源码结构看,这类回补任务单次处理的行数受任务边界约束,取 1,000 行这个“最小阈值”量级是权衡了回补吞吐与单 part 大小后的选择;这也印证了第 3.3 节决策表中“1,000 行是下限而非理想值”的边界语义。
6.3 评审检查清单的闭环
回到 SKILL.md 的 Insert Strategy Reviews 清单,insert-batch-size 与四条相邻规则共同构成写入链路的完整评审面:
- [ ] Batch size 10K-100K rows per INSERT(本规则)
- [ ] 不对频繁变更使用
ALTER TABLE UPDATE(mutation 规避) - [ ] 更新语义使用 ReplacingMergeTree / CollapsingMergeTree
- [ ] 高频小批次场景启用 async inserts(第 5 节的兜底方案)
一条符合 Langfuse 仓库评审标准的写入链路应当是:上游先攒批(事件按项目/类型聚合到批),批量落入 MergeTree 系表(一次 INSERT 一个 part,part 数量可控),确无法攒批的入口再叠加 async insert 兜底,最后用第 4 节的 system.parts 查询周期性验证 part 水位。
7. 小结
insert-batch-size 规则的技术内核可以浓缩为三句话:
- 机制:ClickHouse 中每次 INSERT 产生一个 data part,小批量 INSERT 会让 part 数量以“行数/批大小”的速度膨胀,最终压垮后台 merge 并触发分区级 part 上限(>3000/分区时 INSERT 被阻塞)。
- 标准:同步写入目标 10,000–100,000 行/次、约 1 次 INSERT/秒;低于 1,000 行的写入应改走服务端 async insert 缓冲。
- 验证:用
system.parts(active = 1,按表聚合count())持续监控 part 数量,作为批量化改造的验收指标。
对 Langfuse 这类高吞吐 LLM 可观测性平台而言,该规则由 clickhouse.ts 中的整批 insert 实现、worker 回补任务的显式 BATCH_SIZE 参数与规则库的评审清单共同保证,是写入侧性能与集群稳定性的第一道防线。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00