首页
/ Langfuse ClickHouse 写入最佳实践:INSERT 批量大小(10K–100K 行)的原理与落地

Langfuse ClickHouse 写入最佳实践:INSERT 批量大小(10K–100K 行)的原理与落地

2026-09-09 15:25:43作者:何将鹤

本篇基于 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 系引擎(MergeTreeReplacingMergeTreeCollapsingMergeTree 等)最关键的一条写入语义:

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.

从这条语义可以推断出完整的因果链:

  1. 一次 INSERT = 一个新 part:无论这条 INSERT 携带 1 行还是 10 万行,落盘后都是一个独立的 data part(目录化的列式数据块)。
  2. part 数量爆炸:逐行插入 10,000 条事件就产生 10,000 个 part;每个 part 100 行的小批量插入 100 万次同样产生 100 万个 part。
  3. 后台 merge 追不上:MergeTree 依赖后台 merge 进程将小 part 合并成大 part;part 过多时 merge 队列积压,磁盘上出现大量碎片化的小文件,读取时每次查询都要打开、扫描更多 part(稀疏索引的命中粒度也随之劣化)。
  4. 分区级 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),
  },
});

可以看到三个与规则直接相关的实现细节:

  1. 整批一次性提交values 是整批记录的数组,批次大小由上游调用方(ingestion 处理管线)决定,而不是按行循环调用 insert——这正是 insert-batch-size 所要求的形态;
  2. 使用 JSONEachRow 格式:与规则库 rules/insert-format-native.md 中“使用二进制高效格式”的精神一致;
  3. 查询归因标签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 规则的技术内核可以浓缩为三句话:

  1. 机制:ClickHouse 中每次 INSERT 产生一个 data part,小批量 INSERT 会让 part 数量以“行数/批大小”的速度膨胀,最终压垮后台 merge 并触发分区级 part 上限(>3000/分区时 INSERT 被阻塞)。
  2. 标准:同步写入目标 10,000–100,000 行/次、约 1 次 INSERT/秒;低于 1,000 行的写入应改走服务端 async insert 缓冲。
  3. 验证:用 system.partsactive = 1,按表聚合 count())持续监控 part 数量,作为批量化改造的验收指标。

对 Langfuse 这类高吞吐 LLM 可观测性平台而言,该规则由 clickhouse.ts 中的整批 insert 实现、worker 回补任务的显式 BATCH_SIZE 参数与规则库的评审清单共同保证,是写入侧性能与集群稳定性的第一道防线。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
924
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
599
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
394