☰
Langfuse ClickHouse 压测种子数据方案:用两段 `numbers()` 批量 SQL 生成 traces 与 observations
2026/9/29 20:47:24 网站建设 项目流程

Langfuse ClickHouse 压测种子数据方案:用两段numbers()批量 SQL 生成 traces 与 observations

【免费下载链接】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 开源仓库中一份用于 ClickHouse 压测的种子数据生成方案(原文见 clickhouse-load-seed-plan.md):用两条INSERT INTO ... SELECT FROM numbers(N)语句,分别向traces与observations表灌入可反复调用、可廉价参数化的真实感负载数据。读完本文,你将掌握这套方案的参数旋钮设计、trace↔observation 无 JOIN 的确定性关联技巧、基于 ReplacingMergeTree 的幂等重跑策略,以及配套的健全性校验查询,并理解它如何与仓库中已有的 seeder 实现(clickhouse-builder.ts)相互印证。

设计目标与约束

方案的出发点很明确:为 ClickHouse 负载测试生成"真实、可廉价参数化"的种子数据,最终形态是两张表各一条 SQL 语句,且不经过任何逐行的 Node.js 代码路径——数据完全由 ClickHouse 服务端生成,从而把写入压力集中在数据库本身。

核心目标与约束如下:

  • 每张表一条 INSERT(traces、observations),没有逐行 Node 代码路径;
  • 可配置旋钮通过 CTE 常量传入(详见下一节):
    • 数据量:每批总行数、每秒批次数、时间跨度(天);
    • 低基数桶:version、release、environment、name、level——从 3~5 个值的集合中选取;
    • 高基数字符串:user_id、session_id、trace_id、id——用池大小参数控制重复率;
    • 项目分布:静态的project_id数组,行按桶落入不同项目,单条 INSERT 即可覆盖多个项目;
    • input/output/metadata 长度分布:通过p95_bytes与max_bytes(上限 50 MiB)控制,正文由重复 ASCII 片段切片到目标长度;
    • 重型 metadata 比例:例如 5% 的行携带一个或两个值极大的 map 条目;
  • trace↔observation 关联:使用确定性 ID,使 observations 无需 JOIN 即可引用真实存在的 traces;
  • 不依赖任何 Postgres 前置状态——压测纯跑 ClickHouse;
  • 幂等性:以相同BATCH_NO重跑不会产生重复数据(id由(BATCH_NO, number)派生,重跑时 ID 碰撞,ReplacingMergeTree 按id去重)。

值得一提的是,这一"批量 SQL 直灌"的思路在仓库现有实现中已有对应物:seeder 的 bulk 生成器(clickhouse-builder.ts)正是通过xxHash32(toUInt64(number * 4 + seedSalt))这类加盐哈希列把每行变化与numbers()行号绑定,从而在确定性重跑与数据多样性之间取得平衡。

参数旋钮:CTE 常量区

方案的可调参数全部集中在WITH子句的常量区,改写即可调整压测形态,无需触碰 SELECT 主体。以traces为例,默认值如下:

常量默认值作用
row_count{{ROW_COUNT}}每批总行数
day_window{{DAY_WINDOW}}时间回退窗口(天),决定timestamp散布范围
batch_no{{BATCH_NO}}批次号,用于派生确定性 ID,也是幂等重跑的关键
versions['v1.0','v1.1','v2.0','v2.1','v3.0']低基数版本桶
releases['stable','canary','rc','nightly']低基数发布通道桶
envs['default','staging','prod','eu-prod']低基数环境桶
projects['proj-load-a'..'proj-load-d']参与压测的project_id静态数组
user_pool_size50000user_id池大小,控制高基数重复率
session_pool_size10000session_id池大小
name_pool_size2000trace 名称池大小
p95_bytes2048input/output 长度分布的 p95 分位点
max_bytes50 * 1024 * 1024input/output 长度上限(50 MiB)
large_pct5超过 p95 长度的"大 payload"行占比
heavy_meta_pct5携带超大 metadata 值的行占比
heavy_meta_bytes1024 * 1024超大 metadata 值的字节数

这些旋钮与 seeder 常量文件(clickhouse-seed-constants.ts)中沉淀的REALISTIC_*名称池互为补充:前者负责分布形状,后者提供贴近真实业务的名字(如ChatCompletion、TextSummarization、gpt-5.4-mini)。

最终 SQL(一):traces批量生成

INSERT INTO traces WITH toUInt64({{ROW_COUNT}}) AS row_count, toUInt32({{DAY_WINDOW}}) AS day_window, toUInt64({{BATCH_NO}}) AS batch_no, ['v1.0','v1.1','v2.0','v2.1','v3.0'] AS versions, ['stable','canary','rc','nightly'] AS releases, ['default','staging','prod','eu-prod'] AS envs, ['proj-load-a','proj-load-b','proj-load-c','proj-load-d'] AS projects, toUInt64(50000) AS user_pool_size, toUInt64(10000) AS session_pool_size, toUInt64(2000) AS name_pool_size, toUInt64(2048) AS p95_bytes, toUInt64(50 * 1024 * 1024) AS max_bytes, toUInt8(5) AS large_pct, toUInt8(5) AS heavy_meta_pct, toUInt64(1024 * 1024) AS heavy_meta_bytes SELECT concat('trace-', toString(batch_no), '-', toString(number)) AS id, toDateTime64(now() - randUniform(0, day_window * 86400), 3) AS timestamp, concat('trace-name-', toString(rand() % name_pool_size)) AS name, if(rand() % 100 < 70, concat('user-', toString(cityHash64(rand64()) % user_pool_size)), NULL) AS user_id, multiIf( (rand() % 100) < heavy_meta_pct, map('big_payload', randomPrintableASCII(heavy_meta_bytes), 'kind', 'oversized'), map('env', envs[1 + (rand() % length(envs))], 'tenant', concat('t-', toString(rand() % 200)), 'region', arrayElement(['us-east','eu-west','ap-south'], 1 + (rand() % 3))) ) AS metadata, releases[1 + (rand() % length(releases))] AS release, versions[1 + (rand() % length(versions))] AS version, projects[1 + (number % length(projects))] AS project_id, envs[1 + (rand() % length(envs))] AS environment, rand() % 10 < 8 AS public, rand() % 10 < 1 AS bookmarked, if(rand() % 10 < 3, ['production','ai-agent'], []) AS tags, randomPrintableASCII( multiIf((rand() % 100) < (100 - large_pct), toUInt64(64) + (rand64() % (p95_bytes - 64)), p95_bytes + (rand64() % (max_bytes - p95_bytes))) ) AS input, randomPrintableASCII( multiIf((rand() % 100) < (100 - large_pct), toUInt64(64) + (rand64() % (p95_bytes - 64)), p95_bytes + (rand64() % (max_bytes - p95_bytes))) ) AS output, if(rand() % 100 < 60, concat('sess-', toString(cityHash64(rand64()) % session_pool_size)), NULL) AS session_id, now() AS created_at, now() AS updated_at, now() AS event_ts, toUInt8(0) AS is_deleted FROM numbers(row_count) SETTINGS max_block_size = 256;

逐字段设计要点

  • id(幂等基石):concat('trace-', toString(batch_no), '-', toString(number))。同一BATCH_NO重跑时行号number相同,ID 必然碰撞;配合traces表的ReplacingMergeTree(event_ts, is_deleted)引擎(见 0001_traces.up.sql),重跑数据会被去重而不会翻倍。
  • timestamp散布:now() - randUniform(0, day_window * 86400),毫秒精度DateTime64(3),将批次均匀回退到day_window天窗口内。
  • project_id确定性分桶:projects[1 + (number % length(projects))]——用行号而非随机数取模,保证同一批内项目分布确定、可复现,这是与user_id等随机字段的关键区别。
  • user_id/session_id高基数:cityHash64(rand64()) % pool_size把随机值稳定映射到指定池大小,从而可以"拨动"重复率:池越小重复越多,越接近真实的多用户共享场景。
  • metadata混合:multiIf让heavy_meta_pct(默认 5%)的行携带 1 MiB 的big_payload超大 map 值,其余行携带包含env/tenant/region的常规 map,用于测试 map 键/值上的 bloom filter 索引(表定义中的idx_res_metadata_key/idx_res_metadata_value)在极端负载下的表现。
  • 长度分布:multiIf以100 - large_pct的概率落入[64, p95_bytes)区间,否则落入[p95_bytes, max_bytes)区间,从而构造出"p95 以内为主、偶发超大"的拖尾分布。

最终 SQL(二):observations批量生成

INSERT INTO observations WITH toUInt64({{ROW_COUNT}}) AS row_count, -- e.g. 5x trace rows toUInt32({{DAY_WINDOW}}) AS day_window, toUInt64({{BATCH_NO}}) AS batch_no, toUInt64({{TRACES_IN_BATCH}}) AS traces_in_batch, toUInt8({{OBS_PER_TRACE}}) AS obs_per_trace, -- typical fanout ['default','staging','prod','eu-prod'] AS envs, ['DEFAULT','DEBUG','WARNING','ERROR'] AS levels, ['GENERATION','SPAN','EVENT','AGENT','TOOL'] AS obs_types, ['gpt-4o','gpt-4o-mini','claude-sonnet-4','claude-haiku-4','gemini-2.0'] AS models, ['v1.0','v1.1','v2.0'] AS versions, ['proj-load-a','proj-load-b','proj-load-c','proj-load-d'] AS projects, toUInt64(2048) AS p95_bytes, toUInt64(50 * 1024 * 1024) AS max_bytes, toUInt8(5) AS large_pct, toUInt8(5) AS heavy_meta_pct, toUInt64(1024 * 1024) AS heavy_meta_bytes, toUInt64(2000) AS name_pool_size SELECT concat('obs-', toString(batch_no), '-', toString(number)) AS id, -- trace_id derived from the same (batch_no, number/obs_per_trace) scheme concat('trace-', toString(batch_no), '-', toString(intDiv(number, obs_per_trace) % traces_in_batch)) AS trace_id, -- project_id must match the parent trace projects[1 + ((intDiv(number, obs_per_trace) % traces_in_batch) % length(projects))] AS project_id, envs[1 + (rand() % length(envs))] AS environment, obs_types[1 + (rand() % length(obs_types))] AS type, if(number % obs_per_trace = 0, NULL, concat('obs-', toString(batch_no), '-', toString(number - 1))) AS parent_observation_id, toDateTime64(now() - randUniform(0, day_window * 86400), 3) AS start_time, addMilliseconds(start_time, toInt64(randUniform(50, 5000))) AS end_time, concat('obs-name-', toString(rand() % name_pool_size)) AS name, multiIf( (rand() % 100) < heavy_meta_pct, map('big_payload', randomPrintableASCII(heavy_meta_bytes), 'kind', 'oversized'), map('step', toString(number % obs_per_trace), 'env', envs[1 + (rand() % length(envs))]) ) AS metadata, levels[1 + (rand() % length(levels))] AS level, if(rand() % 100 < 5, 'failed downstream call', NULL) AS status_message, versions[1 + (rand() % length(versions))] AS version, -- input/output: gate on type, same length distribution as traces if(type IN ('GENERATION','EMBEDDING','TOOL'), randomPrintableASCII( multiIf((rand() % 100) < (100 - large_pct), toUInt64(64) + (rand64() % (p95_bytes - 64)), p95_bytes + (rand64() % (max_bytes - p95_bytes)))), NULL) AS input, if(type IN ('GENERATION','EMBEDDING','TOOL','RETRIEVER','EVALUATOR'), randomPrintableASCII( multiIf((rand() % 100) < (100 - large_pct), toUInt64(64) + (rand64() % (p95_bytes - 64)), p95_bytes + (rand64() % (max_bytes - p95_bytes)))), NULL) AS output, if(type = 'GENERATION', models[1 + (rand() % length(models))], NULL) AS provided_model_name, if(type = 'GENERATION', concat('model-', toString(rand() % 50)), NULL) AS internal_model_id, if(type = 'GENERATION', '{"temperature":0.7,"max_tokens":2000}', NULL) AS model_parameters, if(type = 'GENERATION', map('input', toUInt64(randUniform(20, 4000)), 'output', toUInt64(randUniform(10, 2000)), 'total', toUInt64(randUniform(30, 6000))), map()) AS provided_usage_details, if(type = 'GENERATION', map('input', toUInt64(randUniform(20, 4000)), 'output', toUInt64(randUniform(10, 2000)), 'total', toUInt64(randUniform(30, 6000))), map()) AS usage_details, if(type = 'GENERATION', map('input', toDecimal64(randUniform(0.00001, 0.005), 12), 'output', toDecimal64(randUniform(0.00001, 0.01), 12), 'total', toDecimal64(randUniform(0.00002, 0.015), 12)), map()) AS provided_cost_details, if(type = 'GENERATION', map('input', toDecimal64(randUniform(0.00001, 0.005), 12), 'output', toDecimal64(randUniform(0.00001, 0.01), 12), 'total', toDecimal64(randUniform(0.00002, 0.015), 12)), map()) AS cost_details, if(type = 'GENERATION', toDecimal64(randUniform(0.00002, 0.015), 12), NULL) AS total_cost, if(type = 'GENERATION', addMilliseconds(start_time, toInt64(randUniform(50, 500))), NULL) AS completion_start_time, NULL AS prompt_id, NULL AS prompt_name, NULL AS prompt_version, start_time AS created_at, start_time AS updated_at, start_time AS event_ts, toUInt8(0) AS is_deleted, '' AS usage_pricing_tier_id, '' AS usage_pricing_tier_name, map() AS tool_definitions, [] AS tool_calls, [] AS tool_call_names FROM numbers(row_count) SETTINGS max_block_size = 256;

与 traces 的关联设计(无 JOIN 的确定性引用)

这是本方案最精巧的部分:observations 的trace_id完全由算术推导,不查询任何表。

  • trace_id = concat('trace-', batch_no, '-', intDiv(number, obs_per_trace) % traces_in_batch):假设每批有traces_in_batch条 trace(ID 为trace-<batch>-0到trace-<batch>-(traces_in_batch-1)),那么第number行 observation 归属于第intDiv(number, obs_per_trace)条 trace,再对traces_in_batch取模即可保证落在本批 trace ID 空间内。
  • project_id同样用(intDiv(number, obs_per_trace) % traces_in_batch)映射回父 trace 所在的项目,保证"observation 的 project_id 永远与父 trace 一致"——这正是孤儿校验(见下文 sanity-check)通过的前提。
  • parent_observation_id:每个 fanout 组的首行(number % obs_per_trace = 0)为根(NULL),其余行指向本组前一行(obs-<batch>-(number-1)),从而构造出链状/树状层级。这比仓库中 bulk builder 的父子方案(父节点指向number - tracesCount,见 clickhouse-builder.ts)更紧凑,二者思路一致:用行号算术替代逐行查询。

类型门控的字段

observations表(见 0002_observations.up.sql)中大量字段只对GENERATION有意义,因此用if(type = 'GENERATION', ..., NULL)门控:

  • 模型三件套:provided_model_name(从模型池选取)、internal_model_id(随机 50 个)、model_parameters(固定 JSON);
  • 用量与成本:provided_usage_details/usage_details(Map(LowCardinality(String), UInt64),token 数)、provided_cost_details/cost_details与total_cost(Decimal64(12)精度,模拟每 token 成本);
  • TTFT:completion_start_time = addMilliseconds(start_time, randUniform(50, 500)),落在 generation 时长内部,保证首 token 延迟语义合理(seeder 的数据完整性契约同样要求 TTFT 落在 duration 内,见 seeder/README.md 的 Data integrity guarantees 一节);
  • prompt 字段:方案中显式置 NULL——绝不伪造 prompt ID,这与 seeder 契约"generations 要么关联真实 Postgres prompt,要么携带 NULL"完全一致;
  • input/output按type白名单门控:只有GENERATION/EMBEDDING/TOOL有 input,RETRIEVER/EVALUATOR另有 output,其余类型为 NULL。

幂等重跑:为什么重跑不会爆表

方案要求"以相同BATCH_NO重跑不产生重复数据",其机制依赖两张表的表引擎与排序键:

  • traces与observations均为ReplacingMergeTree(event_ts, is_deleted)(见 0001_traces.up.sql 与 0002_observations.up.sql);
  • ORDER BY元组是去重键:traces为(project_id, toDate(timestamp), id),observations为(project_id, type, toDate(start_time), id);
  • 由于id由(batch_no, number)确定性派生,重跑产生的行与旧行拥有完全相同的主键元组,merge 时按event_ts保留新版本,旧版本被替换——数据量不会翻倍。

这里有一个从源码可印证的硬性规则:凡落入 ORDER BY 键的值,绝不能来自顺序随机流或墙钟,否则重跑会静默重复。仓库 seeder 的 README 明确记录了这条"来之不易"的纪律(ClickHouse determinism rules 一节):时间锚点用 TS 侧计算的 UTC 午夜(utcDayStartMs()),行级变化来自加盐的xxHash32(number),并且哈希输入要包toUInt64——xxHash32 哈希的是二进制表示,类型收窄取模会悄悄改变同一值的哈希结果。本方案中的traces把id直接用number派生,规避了这个问题;SETTINGS max_block_size = 256则限制单块规模,便于观察写入节奏。

Sanity-check:每批落库后的三把尺子

文档为每批数据落地后提供了三条校验 SQL,全部基于FINAL(ReplacingMergeTree 去重语义下的最终视图):

-- length distribution per project SELECT project_id, quantile(0.50)(length(input)) p50, quantile(0.95)(length(input)) p95, quantile(0.99)(length(input)) p99, max(length(input)) mx FROM traces FINAL WHERE timestamp > now() - INTERVAL 1 DAY GROUP BY project_id; -- heavy-metadata fraction (should be ~5 %) SELECT countIf(arrayExists(v -> length(v) > 100000, mapValues(metadata))) / count() AS heavy_frac FROM traces FINAL WHERE timestamp > now() - INTERVAL 1 DAY; -- orphan check: every observation should resolve to a trace SELECT count() FROM observations o LEFT ANY JOIN traces t ON o.trace_id = t.id AND o.project_id = t.project_id WHERE o.start_time > now() - INTERVAL 1 DAY AND t.id = '';

三条查询分别验证:

  1. 长度分布:按项目统计 input 长度的 p50/p95/p99 与最大值,确认拖尾分布符合p95_bytes/max_bytes设定,且各项目数据形状一致;
  2. 重型 metadata 占比:统计metadatamap 中任一值超过 100000 字节的行占比,应接近heavy_meta_pct(约 5%);
  3. 孤儿检查:LEFT ANY JOIN找trace_id无法解析到任何 trace 的 observation——应为 0。这是对"确定性关联"设计的最终验收:只要trace_id/project_id的算术推导正确,孤儿数为 0。

与仓库 seeder 生态的呼应

这份 plan 是 seeder 工作区的一部分。虽然它描述的纯 CH 直灌方案刻意不依赖 Postgres 状态,但仓库已有的 bulk 路径提供了可对照的工程化落地:

  • 执行方式:seeder orchestrator 通过clickhouseClient().command({ query, clickhouse_settings: { wait_end_of_query: 1 } })执行批量 SQL(seeder-orchestrator.ts),并可用logStatistics()以bar()可视化各project_id行数分布;本 plan 的两条 SQL 同样可用该入口跑通。
  • 入口脚本:seed-clickhouse.ts 与 load-seed-clickhouse.ts 是两条可执行入口(后者支持项目数 / 天数 / 最大观测数三个 CLI 参数),它们都先通过 Prisma upsert 组织与 API Key,再调用prepareClickhouse()灌数据——本 plan 的 SQL 是"零前置依赖"的简化变体,适合纯写入压测。
  • 确定性纪律一脉相承:bulk builder 用xxHash32(toUInt64(number * 4 + seedSalt))系列构造确定性列(clickhouse-builder.ts),plan 用(batch_no, number)派生 ID,两者都遵守"ORDER BY 键不得来自随机流"的铁律,保证重跑幂等、uniqExact读回可精确断言。
  • 可组合演进:seeder README 的 roadmap 提到未来可加 "Ingestion API writer",用真实 ingestion 批量写入替代直灌 SQL;本 plan 的价值正在于先以最小代价验证 ClickHouse 侧的写入吞吐与形状,为后续真实链路压测提供基线。

使用前提与注意事项

  • 运行环境:SQL 依赖 ClickHouse 的numbers()表函数、randomPrintableASCII、randUniform、cityHash64、multiIf等函数,需在 ClickHouse 服务端直接执行(clickhouse-client或clickhouseClient().command),且目标库已按迁移脚本建好traces/observations表;
  • 占位符替换:{{ROW_COUNT}}、{{DAY_WINDOW}}、{{BATCH_NO}}、{{TRACES_IN_BATCH}}、{{OBS_PER_TRACE}}为模板占位符,执行前需替换为具体数值,且应保证traces_in_batch与 traces 批次实际行数一致,否则孤儿校验会失败;
  • 重跑语义:仅当BATCH_NO保持不变时幂等;换用新批次号意味着全新的一批 ID,可用于叠加更多数据制造持续写入压力;
  • 与 v4 的关系:本 plan 只覆盖 v3 的traces/observations两张表,不涉及 v4 的events_full(seeder 中由--v4标志与event-mirror.ts负责),如需 v4 压测需另行扩展。

【免费下载链接】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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询