Langfuse 实时聚合实战:ClickHouse 增量物化视图(Incremental MV)最佳实践
2026/9/10 20:46:37 网站建设 项目流程

Langfuse 实时聚合实战:ClickHouse 增量物化视图(Incremental MV)最佳实践

【免费下载链接】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 最佳实践规则query-mv-incremental(位于 .agents/skills/clickhouse-best-practices/rules/query-mv-incremental.md)展开,讲解在 Langfuse 这类高吞吐可观测性平台中,如何用增量物化视图把"每次查询全表聚合"的反模式改造成"插入时预聚合、查询时只读千行级结果"的正解。读完本文,你将掌握AggregatingMergeTree + Materialized View + State/Merge 函数三件套的完整落地姿势,并能在 Langfuse 真实迁移脚本(packages/shared/clickhouse/migrations/canonical/)中找到对应生产级范例作为参照。

背景:为什么 Langfuse 需要增量聚合

Langfuse 将 traces、observations、scores、events 等核心遥测数据存储在 ClickHouse 中,其中事件表(events)是典型的超大明细表,累积数据可达数十亿行。而 UI 侧的仪表盘、trace 列表页需要按小时、按事件类型、按项目做实时统计,若沿用传统数据库思维"每次页面加载都现算",一条GROUP BY就要扫描最近 7 天的全部明细,代价是从数十亿行中读数据,延迟与集群压力都不可接受。

规则文件给出的核心主张非常明确:

Impact: HIGH— Incremental MVs 在插入时自动把视图的查询应用到新数据块,结果写入目标表,部分结果随时间合并。读数千行而不是数十亿行,对集群开销极小。

这条规则属于该技能包 28 条规则中的query-mv-*(物化视图)类别,评级为 HIGH。物化视图的选型总览可见 .agents/skills/clickhouse-best-practices/SKILL.md:增量 MV(query-mv-incremental)用于实时聚合,Refreshable MV(query-mv-refreshable)用于复杂 JOIN 与批处理工作流。

反模式:每次查询全量聚合

规则文件首先给出了需要避免的写法。以事件流分析为例,若在每次仪表盘加载时直接对明细表做聚合:

-- Full aggregation on every dashboard load SELECT event_type, toStartOfHour(timestamp) as hour, count() as events, uniq(user_id) as unique_users FROM events WHERE timestamp >= now() - INTERVAL 7 DAY GROUP BY event_type, hour; -- Scans 7 days of data every time (billions of rows)

这条 SQL 的问题在于:聚合结果没有被保存,每次请求都要重新扫描 7 天明细(数十亿行)。在 Langfuse 的查询路径中,这类明细表(如 events)列式存储虽然能压缩数据,但高频、重复的全表聚合依然会持续消耗 CPU 与 IO,拖慢仪表盘首屏。

值得强调的是,Langfuse 技能包中有一条相关的 Langfuse 专属约束:events表被设计为无需FINAL,查询events时严禁使用FINAL关键字,因为该关键字会显著拖慢性能。因此对事件明细的查询优化思路应转向"预聚合落表"而非"靠 FINAL 现场合并"。

正解:增量 MV 三段式(预聚合方案)

规则文件给出了标准的增量物化视图三段式,这也是 ClickHouse 社区广泛采用的最佳实践。请完整复制以下结构:

-- 1. 创建聚合结果目标表 CREATE TABLE events_hourly ( event_type LowCardinality(String), hour DateTime, events AggregateFunction(count), unique_users AggregateFunction(uniq, UInt64) ) ENGINE = AggregatingMergeTree() ORDER BY (event_type, hour); -- 2. 创建物化视图,增量填充目标表 CREATE MATERIALIZED VIEW events_hourly_mv TO events_hourly AS SELECT event_type, toStartOfHour(timestamp) as hour, countState() as events, uniqState(user_id) as unique_users FROM events GROUP BY event_type, hour; -- 3. 查询预聚合数据 SELECT event_type, hour, countMerge(events) as events, uniqMerge(unique_users) as unique_users FROM events_hourly WHERE hour >= now() - INTERVAL 7 DAY GROUP BY event_type, hour; -- Reads thousands of rows instead of billions

三段式逐一拆解:

  1. 目标表events_hourly:使用AggregatingMergeTree引擎,聚合列的类型是AggregateFunction(count)AggregateFunction(uniq, UInt64)。后台合并时,相同ORDER BY (event_type, hour)键的行会把聚合中间态合并起来,从而把"一天内同小时、同类型的多条聚合中间态"收敛成一行。event_type使用LowCardinality(String),符合技能包中schema-types-lowcardinality(低基数字符串用 LowCardinality)的配套建议。
  2. 物化视图events_hourly_mvTO子句把结果定向写入目标表;countState()uniqState()这类带-State后缀的函数输出的是聚合中间态而非最终数值,供后续合并。
  3. 查询语句:用countMerge(events)uniqMerge(unique_users)把中间态"展开"为最终数值。由于数据已被压缩到小时粒度,同样的 7 天窗口只读数千行。

规则文件总结的要点如下,必须严格遵循:

  • MV 内用-State函数,查询内用-Merge函数——这是中间态与终态的正确配对方式;
  • 增量语义——MV 只处理创建之后新插入的数据块,已有历史数据不会自动纳入,需要单独回填(backfill);
  • 插入时开销极小——预聚合只发生在数据写入路径上,对集群的额外负担可以忽略。

State/Merge 机制与 SimpleAggregateFunction 的差异

深入源码可以发现,Langfuse 实际生产表并未一律使用AggregateFunction+-State,而是大量使用SimpleAggregateFunction。以 packages/shared/clickhouse/migrations/canonical/0023_traces_aggregating_merge_trees.up.sql 中的traces_all_amt为例:

CREATE TABLE traces_all_amt {CLICKHOUSE_CLUSTER_CLAUSE} ( `project_id` String, `id` String, `timestamp` SimpleAggregateFunction(min, DateTime64(3)), `end_time` SimpleAggregateFunction(max, Nullable(DateTime64(3))), `name` SimpleAggregateFunction(anyLast, Nullable(String)), `metadata` SimpleAggregateFunction(maxMap, Map(String, String)), `tags` SimpleAggregateFunction(groupUniqArrayArray, Array(String)), `cost_details` SimpleAggregateFunction(sumMap, Map(String, Decimal(38, 12))), `bookmarked` AggregateFunction(argMax, Nullable(Bool), DateTime64(3)), `input` AggregateFunction(argMax, String, DateTime64(3)) CODEC (ZSTD(3)) ) Engine = AggregatingMergeTree() ORDER BY (project_id, id);

对比规则文件示例,可以提炼出两条配套规律:

  1. SimpleAggregateFunction(f, T):适用于minmaxanyLastsumMapgroupUniqArrayArray这类合并结果与数据无关(可交换、幂等性要求低)的函数。它在 MV 里可以直接写min(...)anyLast(...)等普通函数(见同文件 MV 中min(tn.start_time) as start_timemax(coalesce(tn.end_time, tn.start_time)) as end_time),查询端也不需要-Merge,直接SELECT列即可,ClickHouse 在合并时自动按类型语义收敛。
  2. AggregateFunction(f, T[, X]):适用于argMaxuniqcount等需要保留中间态的复杂聚合。这类列在 MV 中必须显式使用-State变体(同文件中argMaxState(tn.input, ...) as input),查询端则需argMaxMerge(...)展开。

traces_all_amt中还有一个值得注意的细节:bookmarkedpublicinputoutput这类"最后一次写入胜出"的字段,用argMaxState(col, event_ts)以事件时间戳为版本,确保合并时取的是时间上最新的值而非任意值——这正是-State中间态函数表达力所在。而规则文件示例中uniqState(user_id)同理,保证uniq的去重统计在多次合并后依然精确。

增量语义:已有数据不会自动回填

规则文件明确强调:"Incremental - existing data not automatically included (backfill separately)"。增量 MV 从创建时刻起,只消费之后插入到源表的数据块;创建前已存在的数十亿历史行不会自动出现在目标表中。

这意味着在 Langfuse 这类有存量数据的系统里落地新 MV 时,必须配套回填流程。仓库中恰好有对应的回填实现可参考:

  • worker/src/backgroundMigrations/backfillEventsFullFromObservations.ts:从 observations 表回填 events 明细;
  • worker/src/backgroundMigrations/backfillEventsFullFromDatasetRunItems.ts:从 dataset run items 回填。

这些后台迁移(background migrations)由 worker/src/backgroundMigrations/backgroundMigrationManager.ts 统一调度,与 MV 的"增量填充"形成"历史全量 + 实时增量"的完整数据闭环。实操时建议:先建 MV 保证新数据开始累积,再跑回填任务补历史,最后在查询侧把两部分结果按聚合键合并。

Langfuse 仓库中的生产级应用实例

实例一:project_environments(项目环境集合)

packages/shared/clickhouse/migrations/canonical/0009_add_project_environments.up.sql 是规则文件方案最简洁的仓库内印证——用聚合表保存"每个项目出现过哪些环境":

CREATE TABLE project_environments {CLICKHOUSE_CLUSTER_CLAUSE} ( `project_id` String, `environments` SimpleAggregateFunction(groupUniqArrayArray, Array(String)) ) ENGINE = {CLICKHOUSE_REPLICATION_PREFIX}AggregatingMergeTree ORDER BY (project_id); CREATE MATERIALIZED VIEW project_environments_traces_mv {CLICKHOUSE_CLUSTER_CLAUSE} TO project_environments AS SELECT project_id, groupUniqArray(environment) AS environments FROM traces GROUP BY project_id;

注意这里同时存在三条 MV(project_environments_traces_mvproject_environments_observations_mvproject_environments_scores_mv),分别从 traces、observations、scores 三个源表聚合写入同一个目标表。这展示了增量 MV 的一个进阶用法:多源合并。三个数据流各自增量写入AggregatingMergeTree,后台合并时按project_id收敛,最终一行包含全部来源的环境集合。查询端直接SELECT environments FROM project_environments WHERE project_id = ...即可,无需任何-Merge

实例二:traces 的 AMT 三级分层(TTL 与版本化字段)

同文件0023_traces_aggregating_merge_trees.up.sql展示了更完整的工业级形态。Langfuse 为 traces 建了traces_nullEngine = Null()的触发器表)、traces_all_amt(全量)、traces_7d_amt(TTL 7 天)、traces_30d_amt(TTL 30 天)三张聚合表,每张配一个 MV:

-- Null 引擎触发器表:只触发 MV,不落明细,节省存储 CREATE TABLE traces_null {CLICKHOUSE_CLUSTER_CLAUSE} ( ... ) Engine = Null(); -- 7 天 TTL 的 AMT 目标表 CREATE TABLE traces_7d_amt {CLICKHOUSE_CLUSTER_CLAUSE} ( ... ) Engine = AggregatingMergeTree() ORDER BY (project_id, id) TTL toDate(start_time) + INTERVAL 7 DAY; CREATE MATERIALIZED VIEW IF NOT EXISTS traces_7d_amt_mv {CLICKHOUSE_CLUSTER_CLAUSE} TO traces_7d_amt AS SELECT tn.project_id as project_id, tn.id as id, min(tn.start_time) as start_time, anyLast(tn.name) as name, argMaxState(tn.input, if(tn.input <> '', tn.event_ts, toDateTime64(0, 3))) as input, ... FROM traces_null tn GROUP BY project_id, id;

这里可以提炼出三个增量 MV 的设计模式:

  1. Null 引擎触发器表:真实数据写入traces_null,MV 从 Null 表消费并转发聚合结果,明细本身不落地,避免重复存储;
  2. TTL 分层:同一份数据按保留窗口建多张聚合表(7 天 / 30 天 / 全量),配合TTL toDate(start_time) + INTERVAL 7 DAY自动淘汰旧分区,让"热查询只访问小表";
  3. 版本化聚合argMaxState(col, event_ts)保证合并后保留最新值,避免乱序写入导致旧值覆盖新值。

实例三:events_core(行级投影 MV)

并非所有 MV 都是聚合型。0041_create_events_core_mv.up.sql中的events_core_mv行级投影MV——把events_full的字段裁剪、leftUTF8(input, 200)截断大字段后写入更轻量的events_core

CREATE MATERIALIZED VIEW IF NOT EXISTS events_core_mv {CLICKHOUSE_CLUSTER_CLAUSE} TO events_core AS SELECT project_id, trace_id, span_id, ..., leftUTF8(input, 200) as input, leftUTF8(output, 200) as output, ... FROM events_full;

这印证了TO目标表 MV 的通用性:无论目标是聚合表还是投影表,"插入时增量消费、写入独立目标表"的机制一致。而 Langfuse 对events的查询统一走 packages/shared/src/server/queries/clickhouse-sql/event-query-builder.ts 查询构建器(技能包中的 Langfuse 专属规则明确要求:对events表的查询必须先尝试该构建器,不要手写 SQL),预聚合/投影落表后,上层查询读的是聚合结果或精简列,天然贴合"读数千行而非数十亿行"的目标。

MV 运维:修改与演进(drop-recreate 禁区)

增量 MV 上线后仍会面临查询逻辑调整。技能包 SKILL.md 针对 Langfuse 给出了两条硬性约束,直接适用于规则文件场景:

  1. 严禁对正在接收实时写入的源表执行"先 DROP 再 CREATE"的 MV 重建——DROPCREATE之间的每一行新插入数据都会静默、永久地丢失,不再进入目标表。正确做法是ALTER TABLE <mv> {CLICKHOUSE_CLUSTER_CLAUSE} MODIFY QUERY <select>,在不中断写入的前提下替换转换逻辑。若新查询增加了列,需对目标表执行ALTER TABLE ... ADD COLUMN IF NOT EXISTS ...,且这些目标表 ALTER 必须携带{CLICKHOUSE_CLUSTERED_ONLY: SETTINGS alter_sync = 2}模板片段,确保任何副本不会在新列就绪前应用新 MV 查询。MODIFY QUERY仅对TO目标表型 MV 可行——而 Langfuse 全部 MV 均使用TO
  2. 不要在迁移中使用CREATE OR REPLACE VIEW/CREATE OR REPLACE TABLE/EXCHANGE TABLES——其原子替换依赖renameat2,在 NFS 挂载(如 AWS EFS)的自托管部署上会失败导致启动中止。普通视图请在同一迁移文件内拆成DROP VIEW IF EXISTS ...+CREATE VIEW ...两条语句。

另外,Langfuse 的迁移模板使用{CLICKHOUSE_CLUSTER_CLAUSE}占位符在集群/单机两种部署模式下渲染 DDL,复制上方任何 SQL 到自己的集群时,需按自身拓扑决定是否追加ON CLUSTER子句。

选型:增量 MV 还是 Refreshable MV?

技能包把物化视图规则拆成两条(query-mv-incremental与 query-mv-refreshable.md),二者的适用边界可以这样区分:

维度增量 MV(本规则)Refreshable MV
触发方式源表插入数据块时实时触发按计划周期性执行
目标引擎常配AggregatingMergeTree(部分结果合并)MergeTree等,全量重写或追加
延迟秒级以内(近乎实时)取决于刷新周期(如每 5 分钟)
典型场景实时聚合计数、去重统计、实时看板复杂多表 JOIN 反范式化、Top N 缓存、批处理 DAG
数据新鲜度始终最新允许轻微滞后换取亚毫秒查询

规则文件示例中的events_hourly(小时级实时聚合)属于前者;而"订单 JOIN 客户 JOIN 商品"这类复杂关联查询(见query-mv-refreshableREFRESH EVERY 5 MINUTE示例)属于后者。Langfuse 中 trace 的 AMT 聚合链路是纯增量 MV 的典型代表,因为遥测写入是持续高频流,实时聚合价值最高。选型时遵循技能包原则:实时聚合选增量 MV,复杂 JOIN 与批处理选 Refreshable MV;同时谨记 Refreshable MV 的警告——查询执行耗时应远小于刷新间隔,不要每 10 秒刷新一个要跑 10 秒以上的查询。

总结

规则query-mv-incremental的结论可以浓缩为一句话:把聚合从"查询时"前移到"插入时",让AggregatingMergeTree在后台完成部分结果的合并。落地清单如下:

  1. 明细表之上建AggregatingMergeTree目标表,聚合列用AggregateFunction(...)(复杂聚合)或SimpleAggregateFunction(...)(可交换简单聚合);
  2. TO目标表的增量 MV,MV 内对AggregateFunction列使用-State变体函数,查询时用对应-Merge变体展开;
  3. 牢记增量语义:历史数据需独立回填(参考仓库的 backgroundMigrations 实现);
  4. 上线后修改 MV 用ALTER TABLE ... MODIFY QUERY,永远不要 DROP 正在接收写入的 MV;
  5. 与 Refreshable MV 按"实时 vs 批处理"分流,各司其职。

Langfuse 在packages/shared/clickhouse/migrations/canonical/0023_traces_aggregating_merge_trees.up.sql0009_add_project_environments.up.sql0041_create_events_core_mv.up.sql中的真实迁移脚本,就是这条规则最完整的生产级注脚——你可以直接打开这些文件,对照本文逐段验证State/Merge配对、SimpleAggregateFunction简化、TTL 分层与多源合并的每一种写法。

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

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

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

立即咨询