的创建与分批回填脚本实战)
Civitai Event EngineClickHouse 指标日聚合物化视图entityMetricDaily的创建与分批回填脚本实战【免费下载链接】civitaiA repository of models, textual inversions, and more项目地址: https://gitcode.com/GitHub_Trending/ci/civitai导读在 Civitai 的 event-engine 服务中海量实体指标事件entityMetricEvents需要被实时聚合为按实体 × 指标 × 天的日聚合数据供图片 Feed、SSR 等热路径快速读取。本篇文章围绕 clickhouse-metric-rollup-view-script.md 这份规划文档完整讲解其提出的核心方案创建 SummingMergeTree 支撑表backing table与物化视图Materialized View在物化视图建立瞬间记录时间戳再以该时间戳为边界、从可配置的起始日期默认2022-11-01起按月脚本实现中为按周分批回填历史数据。读完本文你将掌握这套MV 创建时间戳 分片回填的可重复执行模式并理解仓库中 setup-clickhouse-rollup-view.ts 脚本的完整实现与消费端 metrics.ts 的读取方式。一、问题背景事件流与日聚合表的缺口event-engine 是 Civitai 仓库中的一个 Kafka 消费者服务负责监听 Postgres 与 ClickHouse 的数据库事件并处理指标、信号、CDC 副作用等见 package.json 的描述。其指标链路大致是业务侧把每次指标变化作为事件写入entityMetricEvents包含entityType、entityId、metricType、metricValue、createdAt聚合层需要把事件流按天折叠成一行(entityType, entityId, metricType, day, total)读取端如图片 Feed / SSR再基于日聚合表查询某类实体的累计指标。问题在于物化视图Materialized View从创建那一刻起才会开始增量消费新写入的事件创建之前的历史事件它看不到。如果直接基于全表SELECT ... GROUP BY做一次性聚合在数据量巨大时又会长时间占用 ClickHouse 资源。于是规划文档提出了一个标准解法——先建支撑表与物化视图记录视图创建时间matViewStart再按时间片月分批回填createdAt matViewStart的历史数据。这样增量由 MV 保证实时性存量由回填保证完整性且回填压力被摊薄到多个小批次上。二、核心 SQL 方案支撑表 物化视图 创建时间戳 按月回填规划文档给出了完整的可执行 SQL 骨架下面逐段拆解保留原文语义并补充注释。2.1 支撑表Backing Table-- Backing table DROP TABLE IF EXISTS entityMetricDailyAgg; CREATE TABLE entityMetricDailyAgg ( entityType LowCardinality(String), entityId Int32, metricType LowCardinality(String), day Date, total Int64 ) ENGINE SummingMergeTree ORDER BY (entityType, entityId, metricType, day) SETTINGS index_granularity 8192;要点SummingMergeTree引擎写入时不会即时合并而是在后台合并过程中把排序键相同的行按数值列求和。也就是说同一天同一实体同一指标的多条记录最终会折叠成一行total为和的记录查询时必须带GROUP BY才能拿到正确的汇总结果读取端 metrics.ts 正是这么做的见下文第五节。LowCardinality(String)entityType与metricType取值集合有限Image/Article/User 等实体Like/Heart/View 等指标用 LowCardinality 压缩字典编码显著降低存储与扫描开销。ORDER BY (entityType, entityId, metricType, day)排序键即聚合粒度保证同一聚合键的行物理相邻使合并与范围查询都高效。index_granularity 8192ClickHouse 默认的稀疏索引粒度每 8192 行一个索引标记这里显式声明以便与其他表保持一致。2.2 物化视图Materialized View-- Mat view DROP VIEW IF EXISTS entityMetricDaily; CREATE MATERIALIZED VIEW entityMetricDaily TO entityMetricDailyAgg AS SELECT entityType, entityId, metricType, toDate(createdAt) AS day, sum(metricValue) AS total FROM entityMetricEvents GROUP BY entityType, entityId, metricType, day;要点TO entityMetricDailyAgg物化视图的目标表就是支撑表。每当entityMetricEvents有新数据插入MV 会执行SELECT转换并把结果写入entityMetricDailyAgg实现增量实时聚合。toDate(createdAt) AS day把事件的时间戳截断到天粒度这是日聚合的维度切分点。sum(metricValue) AS total同一天内同一聚合键的多条事件在此求和。2.3 记录视图创建时间matViewStart-- Get time of mat view creation SELECT now(); -- matViewStart在 MV 创建之后立刻执行SELECT now()把当前时间记为matViewStart。它的意义是MV 只会捕获它创建之后写入entityMetricEvents的新事件因此回填只需处理createdAt matViewStart的历史数据两条链路增量 存量以这个时间戳为界不会重叠也不会遗漏。2.4 按月分批回填Backfill-- Backfill INSERT INTO entityMetricDailyAgg SELECT entityType, entityId, metricType, toDate(createdAt) AS day, sum(metricValue) AS total FROM entityMetricEvents -- WHERE createdAt matViewStart -- Process in batches of months using: -- AND createdAt 2025-01-01 AND createdAt 2025-02-01 GROUP BY entityType, entityId, metricType, day;要点回填的SELECT与 MV 的SELECT逐字段一致保证历史聚合与实时聚合口径完全相同注释中的两处WHERE揭示了分批策略外层用createdAt matViewStart圈定历史范围内层每次用createdAt YYYY-MM-01 AND createdAt YYYY-MM-01下一个月 1 号按月切片执行由于写入目标仍是SummingMergeTree表分批写入的多个月份数据最终会在后台合并时自然求和无需担心批次边界重复同键同天同指标的行在合并时相加等价于一次全量聚合的结果。三、从规划到实现setup-clickhouse-rollup-view.ts 脚本规划文档停留在 SQL 骨架层面而仓库里已经把它落成了可重复执行的 TypeScript 脚本 setup-clickhouse-rollup-view.ts并通过npm run setup:clickhouse-rollup见 package.json即可运行。下面结合源码逐步骤说明。3.1 环境变量与连接const clickhouseUrl process.env.CLICKHOUSE_URL; if (!clickhouseUrl) { throw new Error(CLICKHOUSE_URL is not defined); } const backfillStartDate process.env.BACKFILL_START_DATE || 2022-11-01; const client createClient({ url: clickhouseUrl });CLICKHOUSE_URL必填ClickHouse HTTP 接口地址BACKFILL_START_DATE回填起始日期默认2022-11-01与规划文档要求的默认值完全一致可用环境变量覆盖。3.2 步骤一创建支撑表const createBackingTableQuery CREATE TABLE IF NOT EXISTS entityMetricDailyAgg ( entityType LowCardinality(String), entityId Int32, metricType LowCardinality(String), day Date, total Int64 ) ENGINE SummingMergeTree ORDER BY (entityType, entityId, metricType, day) SETTINGS index_granularity 8192 ;与规划文档 SQL 唯一的差异是把CREATE TABLE改成了CREATE TABLE IF NOT EXISTS使脚本可安全重跑幂等。3.3 步骤二创建物化视图const createMatViewQuery CREATE MATERIALIZED VIEW IF NOT EXISTS entityMetricDaily TO entityMetricDailyAgg AS SELECT entityType, entityId, metricType, toDate(createdAt) AS day, sum(metricValue) AS total FROM entityMetricEvents GROUP BY entityType, entityId, metricType, day ;同样使用IF NOT EXISTS保证重复执行不会报错。3.4 步骤三记录视图创建时间const matViewStartResult await client.query({ query: SELECT now() as matViewStart, format: JSONEachRow, }); const matViewStartData await matViewStartResult.json() as Array{matViewStart: string}; const matViewStart new Date(matViewStartData[0].matViewStart);脚本把规划文档中手工记下SELECT now()结果的动作自动化从查询结果解析出matViewStart字符串并转成Date对象作为后续回填循环的终止边界。3.5 步骤四分批回填规划文档要求按月、一次一个月month by month源码实现选择了按周切片并特意从包含起始日期的那个周一开始对齐周边界const startDate new Date(backfillStartDate); let currentWeek new Date(startDate); currentWeek.setDate(currentWeek.getDate() - currentWeek.getDay() 1); // 回退到本周周一 currentWeek.setHours(0, 0, 0, 0);随后先遍历一遍计算总周数用于进度展示再进入主循环while (currentWeek matViewStart) { const nextWeek new Date(currentWeek); nextWeek.setDate(nextWeek.getDate() 7); const backfillQuery INSERT INTO entityMetricDailyAgg SELECT entityType, entityId, metricType, toDate(createdAt) AS day, sum(metricValue) AS total FROM entityMetricEvents WHERE createdAt ${weekStr} AND createdAt ${nextWeekStr} GROUP BY entityType, entityId, metricType, day ; await client.exec({ query: backfillQuery }); // ... currentWeek nextWeek; }关键点窗口右开每个批次用createdAt 周一起 AND createdAt 下周一批次之间严格无缝、不重叠保证既不丢数据也不重复计数时间边界与 MV 的衔接循环条件currentWeek matViewStart确保回填只覆盖 MV 建立之前的历史窗口MV 建立之后的数据由实时增量负责进度与耗时脚本在每批执行前后打印[n/total] Processing week A to B、✓ Processed N rows in Xs便于长时间运行时的观察与监控源码 setup-clickhouse-rollup-view.ts。这里需要如实说明一个规划与实现的差异文档建议按月切片仓库脚本实际按周切片。从实现角度看周窗口更短、单批更轻、失败重跑粒度更细且对 SummingMergeTree 而言分片粒度不影响最终正确性——两种粒度在语义上等价你可以按数据量与 ClickHouse 负载自行调整循环步长。3.6 步骤五验证聚合结果回填完成后脚本会执行两类验证查询setup-clickhouse-rollup-view.ts整体统计count(*)总行数、min(day)/max(day)日期范围、uniq(entityType)/uniq(metricType)去重类型数、sum(total)总指标值按实体类型抽样GROUP BY entityType后输出每种实体类型的行数与总值 Top 10。这相当于把规划文档中回填后如何确认数据正确的隐含问题显式化给运维一个快速 sanity check。3.7 辅助能力--drop-existing 与干净重建脚本还导出了一个dropExistingRollupTables()函数并支持命令行参数--drop-existingsetup-clickhouse-rollup-view.tsawait client.exec({ query: DROP VIEW IF EXISTS entityMetricDaily }); await client.exec({ query: DROP TABLE IF EXISTS entityMetricDailyAgg });node scripts/setup-clickhouse-rollup-view.ts --drop-existing会先删除旧 MV 与旧表再重建适合需要完全重置聚合数据例如 schema 变更或数据修复的场景。四、更进一步的仓库佐证current-totals 点查表与历史回填实战这套MV 支撑表 回填模式在 event-engine 仓库里被反复使用进一步印证了规划文档的通用性。4.1 演进entityMetricCurrentTotals_v2 点查表在 clickhouse-current-totals-table.md 中记录了同一模式的进阶形态由于读取端直接查聚合视图entityMetricDailyAgg占了 ClickHouse 约 86% 的流量约 28 q/s团队新增了entityMetricCurrentTotals_v2点查表ORDER BY (entityType, entityId, metricType)配合REFRESH EVERY 1 HOUR的可刷新物化视图每小时原子重算。它的回填同样沿用了INSERT INTO ... SELECT模式SQL 与可运行脚本分别是 entity-metric-current-totals.sql 和 setup-clickhouse-current-totals.ts后者明言mirrors setup-clickhouse-rollup-view.ts。这段文档还记录了一个与回填高度相关的正确性陷阱历史表是SharedReplacingMergeTree按sealedAt去重的直接sum(total)会重复计数必须argMax(total, sealedAt) GROUP BY day取每天最后一次 sealed 的值再求和该语义已用entityId % 9973 0的确定性样本验证21,388 个 key 0 处不一致。这提醒我们在动手写回填 SQL 前必须先弄清源表的去重/合并语义而不是盲目GROUP BY求和。4.2 复杂回填案例backfill-user-comment-count.mjs另一个高价值的参考实现是 backfill-user-comment-count.mjs——它演示了事件流上线时间点不一致时的回填边界处理不同 surface 的前向事件起始时间不同用单一边界回填会双算或漏算脚本因此按 surface 区分firstEvent与--new-cutover两个边界并做 dry-run--execute才真正写入、--no-dirty/--dirty-only分阶段演练、dirty 标记分批限速每批 2500、间隔 60s等精细控制。这与规划文档回填要分批、要控制对生产库的冲击的思想一脉相承可作为复杂回填任务的模板。五、消费端如何读取MetricService 与 entityMetricDailyAgg回填与 MV 的最终目标是服务读取端。在 metrics.ts 中MetricService的fetchFromClickhouse()方法正是对entityMetricDailyAgg做点查聚合metrics.tsSELECT entityId, metricType, sum(total) AS value FROM entityMetricDailyAgg WHERE entityType entityType AND entityId IN (batch) AND metricType IN (metricTypes) GROUP BY entityId, metricType HAVING value 0注意几个与本文主题呼应的细节查询带GROUP BY entityId, metricType——再次印证 SummingMergeTree 表读取必须显式聚合按FETCH_BATCH_SIZE 1000分批查询见 metrics.ts避免 IN 子句过大上层用 Redis 做 24 小时缓存含 10% 概率的 TTL 滑动、notFound负缓存 5 分钟、setNx分布式锁防缓存击穿见 metrics.tsClickHouse 只在缓存 miss 时被访问——这也是为什么日聚合表的查询要尽量廉价。fetchTimeframes()则利用sumIf按 Day / Week / Month / Year / AllTime 五个时间窗一次算出metrics.ts同样是建立在每天一行的预聚合基础之上的。六、运行方式与注意事项6.1 如何运行# 在 apps/event-engine 目录下 npm run setup:clickhouse-rollup # 用默认起始日期 2022-11-01 CLICKHOUSE_URLhttp://localhost:8123 npm run setup:clickhouse-rollup BACKFILL_START_DATE2023-01-01 CLICKHOUSE_URL... npm run setup:clickhouse-rollup # 自定义起始日期 npm run setup:clickhouse-rollup -- --drop-existing # 先删除旧 MV 与表再重建对应 npm script 见 package.json底层是tsx scripts/setup-clickhouse-rollup-view.ts。6.2 实操注意事项entityMetricEvents表必须先存在MV 的FROM entityMetricEvents依赖源表脚本本身不创建源表回填期间 MV 持续写入同一张支撑表这是设计内行为——增量与存量回填并行写入靠 SummingMergeTree 的合并语义收敛无需暂停 MV长耗时任务建议在低峰执行历史数据量大时回填会持续扫描entityMetricEvents对 ClickHouse 有一定负载仓库在 current-totals 的运行手册中明确标注回填是 HUMAN-GATED 操作不要在饱和的 ClickHouse 上自动执行见 clickhouse-current-totals-table.md这一原则同样适用于本脚本验证不可省略回填完成后用脚本自带的统计与抽样查询核对日期范围、类型数与总量必要时抽样对比源表。结语从 clickhouse-metric-rollup-view-script.md 的规划草稿到 setup-clickhouse-rollup-view.ts 的可运行实现再到 metrics.ts 的消费端我们看到了 event-engine 指标链路中一条完整且可复用的实践SummingMergeTree 支撑表承接物化视图的增量写入SELECT now()时间戳划定增量与存量的边界按周或按月切片 右开区间保证回填不重不漏最后用统计查询验证结果。这套模式同样支撑了后续的entityMetricCurrentTotals_v2点查表与复杂的用户评论数回填是理解整个指标系统数据管线的最佳入口。【免费下载链接】civitaiA repository of models, textual inversions, and more项目地址: https://gitcode.com/GitHub_Trending/ci/civitai创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考