Langfuse 的 ClickHouse 写入调优:用 async_insert 应对高频小批量写入
【免费下载链接】langfuse🪢 Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. 🍊YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse
Langfuse 是一个开源的 AI 工程平台(LLM 可观测性、评估、指标、Prompt 管理、数据集),其核心事件数据存储在 ClickHouse 中。本文基于 Langfuse 仓库内 ClickHouse 最佳实践规则集(clickhouse-best-practicesskill)中的insert-async-small-batches规则展开,讲清"当客户端无法自行攒批时,如何用服务端异步插入(async insert)自动合并成更大的数据 part",并结合 Langfuse 自身的 ClickHouse 客户端实现 展示这条规则在真实生产系统中的落地方式。读完本文,你能掌握 async insert 的适用场景、关键参数与返回模式取舍,并知道在 Langfuse 中应通过哪些环境变量调优吞吐。
问题背景:为什么高频小批量写入会压垮 ClickHouse
Langfuse 的 ClickHouse 规则集将插入策略分为几类规则,其中两条直接决定了本文的上下文:
insert-batch-size(影响级别 CRITICAL):每条INSERT都会生成一个新的数据 part,单行或极小批量写入会制造数千个微小 part,压垮 merge 进程并引发集群不稳定。规则给出的建议批量区间是10,000–100,000 行/次 INSERT(最低不低于 1,000 行,同步插入速率约 1 次/秒)。insert-async-small-batches(影响级别 HIGH):本文主角——当客户端侧攒批(client-side batching)不可行时,用服务端异步插入做缓冲,自动合成更大的 part。
也就是说,两者是"首选"与"兜底"的关系:客户端能攒批就攒批(见 insert-batch-size.md);攒不了(例如上游事件天然零散、写入方众多、每个写入方量都很小)时,才启用 async insert 让 ClickHouse 服务端替你攒批。规则集 SKILL.md 把这两条归入 "Insert Strategy" 类别,并明确要求在 Insert Strategy Review 时逐条检查:批量大小是否落在 10K–100K、高频小批量场景是否启用了 async insert。
规则核心:async_insert 的用法
规则原文(insert-async-small-batches.md)给出的反模式是"小批量 + 同步直插":
# Small batches without async_insert - creates too many parts for batch in chunks(events, 100): client.execute("INSERT INTO events VALUES", batch)每 100 行一次 INSERT,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规则给出的两个关键示例参数及其含义:
| 参数 | 示例值 | 含义 |
|---|---|---|
async_insert_max_data_size | 10,000,000 | 异步缓冲区累计到 10MB 即触发一次真实写入(生成 part) |
async_insert_busy_timeout_ms | 1000 | 距上次数据进入缓冲区 1 秒后强制刷盘 |
刷盘时机:三个条件"先到先触发"
规则明确了缓冲区何时真正落成一个数据 part,以下任一条件先满足即触发:
- 缓冲区数据量达到
async_insert_max_data_size; - 距上次有数据进入缓冲区的时间超过
async_insert_busy_timeout_ms; - 累积的插入查询数量达到上限。
调参的思路由此而来:max_data_size控制 part 的"重量下限"(越大 part 越肥、merge 压力越小),busy_timeout_ms控制数据的可见性延迟上界。两者需要按业务对"写入延迟 vs part 数量"的容忍度权衡。
返回模式:wait_for_async_insert 决定你能否感知错误
| 设置 | 行为 | 适用场景 |
|---|---|---|
wait_for_async_insert=1 | 等待缓冲区真正刷盘,确认数据持久化 | 推荐 |
wait_for_async_insert=0 | Fire-and-forget,客户端对服务端错误无感知 | 危险——仅在你明确接受数据丢失时使用 |
规则将wait_for_async_insert=1标记为推荐项,将=0明确标注为 Risky:关掉等待后,即使服务端刷盘失败,客户端也会以为插入成功,错误被彻底吞掉。对 Langfuse 这类承载可观测性数据的系统而言,"丢写入"本身就是最严重的故障类型,因此这条取舍在生产系统中几乎没有争议。
Langfuse 源码中的落地:客户端全局开启 async insert
规则讲"应该怎么做",Langfuse 的源码展示了"生产系统实际怎么做"。在 client.ts 的ClickHouseClientManager.getClient()中,所有客户端(读写、只读、事件只读三类服务)在创建时都硬性注入了异步插入设置(约 L233–L234):
clickhouse_settings: { // Overwrite async insert settings to tune throughput ...(env.CLICKHOUSE_ASYNC_INSERT_MAX_DATA_SIZE ? { async_insert_max_data_size: env.CLICKHOUSE_ASYNC_INSERT_MAX_DATA_SIZE } : {}), ...(env.CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MS ? { async_insert_busy_timeout_ms: env.CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MS } : {}), ...(env.CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MIN_MS ? { async_insert_busy_timeout_min_ms: env.CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MIN_MS } : {}), // ... async_insert: 1, wait_for_async_insert: 1, // if disabled, we won't get errors from clickhouse },这段代码有三个值得注意的事实:
async_insert = 1与wait_for_async_insert = 1被硬编码为所有客户端的默认行为,与规则中"推荐等待刷盘以确认持久化"一致。注释// if disabled, we won't get errors from clickhouse恰好印证了规则表格里wait_for_async_insert=0被标为 Risky 的原因——关闭等待会丢失服务端错误信号。- 两个刷盘参数通过环境变量可选覆盖(
CLICKHOUSE_ASYNC_INSERT_MAX_DATA_SIZE、CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MS、CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MIN_MS),代码注释写明目的:"Overwrite async insert settings to tune throughput"。对应的环境声明在 env.ts(约 L133–L140),其中CLICKHOUSE_ASYNC_INSERT_MAX_DATA_SIZE是字符串型可选项,CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MS是整数可选项,CLICKHOUSE_ASYNC_INSERT_BUSY_TIMEOUT_MIN_MS额外要求最小值 50ms。注释 "Optional to allow for server-setting fallbacks" 说明设计上允许把这三个参数交给数据库服务端(如上文ALTER USER)兜底,客户端不传时不覆盖服务端设置。 - Langfuse 自身还叠加了一层客户端攒批:worker 侧的 ClickhouseWriter 按表把事件先放入内存队列,依据
LANGFUSE_INGESTION_CLICKHOUSE_WRITE_BATCH_SIZE(批次行数)和LANGFUSE_INGESTION_CLICKHOUSE_WRITE_INTERVAL_MS(最大间隔)两种条件触发 flush——这与insert-batch-size规则"客户端优先攒批"的思想完全对应。也就是说 Langfuse 是"双保险":应用层先攒到大批量,同时底层客户端对仍然零散的 INSERT 用 async insert 兜底合成更大 part。SKILL.md 的 Langfuse 专属规则也提到 ClickhouseWriter 的插入在查询归属(query attribution)中使用projectId = "MULTI_PROJECT",可据此在system.query_log中定位这些写入。
如何验证 async insert 是否生效:监控 part 数
无论采用哪种策略,最终都要回到同一个可验证指标——活动 part 数量。insert-batch-size规则给出的验证 SQL 同样适用于 async insert 调优后的效果确认:
-- 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;开启 async insert 前,高频小批量写入会让parts持续膨胀(单分区超过 3000 个活动 part 时会阻塞新插入);开启并调好async_insert_max_data_size/async_insert_busy_timeout_ms后,同等写入量下 part 数量应显著下降。这条 SQL 是判断"async insert 是否真正减少了 part 产出"的直接依据。
规则在 Langfuse 规则集中的位置与适用边界
结合 SKILL.md 的规则优先级表,insert-async-前缀(Async Inserts)的影响级别为 HIGH,仅次于 CRITICAL 级的批量大小(insert-batch-)与 mutation 规避(insert-mutation-)规则。由此可以推断出该规则集的决策顺序:
- 先检查客户端能否把批量做到 10K–100K 行(CRITICAL,首选方案);
- 无法做到时,再对高频小批量写入启用 async insert(HIGH,服务端兜底);
- 无论哪种方案,都用
system.parts监控 part 数验证效果。
适用前提与限制也需要说明:async insert 是 ClickHouse 服务端能力,具体参数名与行为以你所运行 ClickHouse 版本的官方"Selecting an Insert Strategy"最佳实践文档为准(规则文件末尾即引用了该文档);本文代码证据来自当前 Langfuse 仓库,其中CLICKHOUSE_*环境变量的解析位于 packages/shared/src/env.ts,客户端注入位于 packages/shared/src/server/clickhouse/client.ts。
小结
insert-async-small-batches规则的核心结论可以浓缩为三点:
- 适用场景:客户端侧攒批不现实的高频小批量写入,用服务端缓冲自动合成大 part,避免 part 爆炸;
- 安全配置:
async_insert = 1+wait_for_async_insert = 1,通过async_insert_max_data_size(按数据量)与async_insert_busy_timeout_ms(按时间)控制刷盘时机,先到先触发; - 生产印证:Langfuse 在 client.ts 中对全部 ClickHouse 客户端硬编码了这两个设置,并保留三个
CLICKHOUSE_ASYNC_INSERT_*环境变量用于吞吐调优,同时在上层 ClickhouseWriter 保留客户端攒批作为第一道防线。
【免费下载链接】langfuse🪢 Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. 🍊YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考