Snowflake MigrationClickHouse Workshops

MergeTree 引擎

从 MergeTree 家族中选择引擎,并设计一个真正物有所值的 ORDER BY 键。

ClickHouse 的所有数据都存放在由某个 MergeTree 引擎变体支撑的表中。如果你是从 Snowflake 过来的,那么这里没有对应的概念,Snowflake 在内部处理所有存储决策。在 ClickHouse 中,选择正确的引擎是你的责任,而选错会产生静默的错误结果。

本指南涵盖你在 NYC Taxi 实验课中会用到的引擎,以及那些让每一位 Snowflake 迁移者踩坑的陷阱。


什么是 MergeTree?

MergeTree 是 ClickHouse 的主力存储引擎。数据被写入称为 part 的不可变列式文件。ClickHouse 会在后台周期性地合并 part,按引擎的规则对它们进行排序、压缩,以及可选的转换。

关键后果是:在合并发生之前,一次读取可能看到同一行的多个版本。 大多数引擎会透明地处理这一点,但有些引擎(尤其是 ReplacingMergeTree)要求你理解合并生命周期才能写出正确的查询。

创建 MergeTree 表时,你必须指定 ORDER BY。它决定了:

  1. 每个 part 内部数据的物理排序顺序
  2. 主索引(稀疏的、块级的,常驻内存)
  3. 对于会去重的引擎,哪些列构成去重所依据的“键”

这里没有独立的主键、聚簇索引或分布键概念。ORDER BY 同时扮演了这些全部角色。


MergeTree

何时使用: 表只有插入操作,或更新在外部处理。不需要去重。

CREATE TABLE default.some_events (
    event_id      String,
    occurred_at   DateTime64(3, 'UTC'),
    payload       String
)
ENGINE = MergeTree()
ORDER BY (occurred_at, event_id);

特性:

  • 插入会以新 part 的形式追加数据
  • 不去重,重复行会被保留
  • 合并会优化存储和压缩,但不改变逻辑内容
  • 查询会读取所有匹配 ORDER BY 前缀范围的 part

它出错的时候: 如果你把同一行插入两次(例如网络故障后重试),两行都会出现在查询结果中。对于真正只插入、不可能产生重复的管道,这是正确的。对于任何接收 CDC 更新或可重试加载的表,请使用 ReplacingMergeTree。


ReplacingMergeTree

何时使用: 行可能被更新(例如车费修正、状态变更)。你希望查询结果中每个键只有一行。

CREATE TABLE analytics.fact_trips (
    trip_id       String,
    pickup_at     DateTime64(3, 'UTC'),
    fare_amount   Float64,
    updated_at    DateTime64(3, 'UTC'),
    -- ...
)
ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id);

特性:

  • 在后台合并期间,具有相同 ORDER BY 键的行会被去重:只保留版本列取值最大的那一行
  • 版本列(这里是 updated_at)决定哪一行胜出,值越大 = 越新 = 被保留
  • 去重是异步的,在合并运行之前,旧版本和新版本会同时存在

关键陷阱:去重延迟

在两次合并之间,不带 FINAL 的查询会看到某一行的所有版本:

-- This may return multiple rows for the same trip_id
-- if the row has been updated since the last merge
SELECT * FROM analytics.fact_trips WHERE trip_id = 'abc123';

-- This returns exactly one row per trip_id, applying deduplication at query time
SELECT * FROM analytics.fact_trips FINAL WHERE trip_id = 'abc123';

FINAL 强制在读取时去重。它比不带 FINAL 的读取更慢,因为 ClickHouse 必须检查所有 part 中是否有重复键。在 NYC Taxi 实验课中,所有针对 fact_trips 的查询都使用 FINAL。

它出错的时候:

  • 在点查中省略 FINAL → 静默返回重复行;聚合结果偏大
  • 使用了错误的版本列(一个不会随更新递增的列)→ 旧值胜出
  • 对可变表使用 MergeTree 而非 RMT → 所有版本不断累积;行数无界增长
  • 期待同步去重 → ETL 作业在插入后立即读取,看到重复数据

RMT 与 dbt: delete_insert 增量策略会在插入前先删除来批数据键范围内的行,因此表里根本不会出现重复。为安全起见仍建议使用 FINAL,但在 dbt 策略正确的情况下它不那么关键。


AggregatingMergeTree

何时使用: 表中存放的是部分聚合状态,这些状态应在后台合并期间被合并,并在查询时被组合。

CREATE TABLE analytics.agg_hourly_revenue (
    hour_bucket   DateTime,
    borough       String,
    fare_sum      AggregateFunction(sum, Float64),
    trip_count    AggregateFunction(count, UInt64)
)
ENGINE = AggregatingMergeTree()
ORDER BY (hour_bucket, borough);

特性:

  • 具有相同 ORDER BY 键的行会使用聚合函数的合并逻辑被归并
  • 查询时使用 -Merge 后缀的组合器:sumMerge(fare_sum)、countMerge(trip_count)
  • 通常由一个 materialized view 供给,它把原始插入转换为部分状态

何时使用它: AggregatingMergeTree 适用于部分状态必须可组合的预聚合数据。在 NYC Taxi 实验课中,agg_hourly_zone_trips 每次运行都由 dbt 重建,它是一张全量替换的表,不是部分状态累加器。那里应该使用 ReplacingMergeTree。

它出错的时候: 查询时用 sum(fare_sum) 而不是 sumMerge(fare_sum),会把二进制聚合状态当成 Float64 处理,返回垃圾数字。这是一种静默的正确性错误。


CollapsingMergeTree

何时使用: 你需要通过插入一个“符号行”来删除或更新行(sign=1 表示插入,sign=-1 表示撤销)。较不常见,但对基于事件的 CDC 模式有用。

ENGINE = CollapsingMergeTree(sign)

在合并期间,同一键上 sign=1 与 sign=-1 的成对行会相互抵消。NYC Taxi 实验课中未使用,对于本工作负载的插入重试模式,带版本列的 ReplacingMergeTree 更简单。


带 TTL 的 MergeTree

为任意 MergeTree 变体添加基于时间的数据过期:

CREATE TABLE default.trips_raw (
    trip_id    String,
    pickup_at  DateTime64(3, 'UTC'),
    _synced_at DateTime DEFAULT now(),
    -- ...
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id)
TTL toDate(pickup_at) + INTERVAL 2 YEAR;

TTL 在后台合并期间触发。过期的行会在其所在 part 合并时被移除。本实验课未配置 TTL,全部 4 年的数据都被保留。在生产环境中,TTL 对控制存储成本至关重要。


选择引擎:决策树

Does the table receive UPDATE or DELETE operations?
├── No (insert-only, e.g., event log, append-only stream)
│   └── MergeTree()
└── Yes
    ├── Do rows have a version/timestamp column that increases on update?
    │   ├── Yes → ReplacingMergeTree(version_col)
    │   └── No (full reload, e.g., dim tables rebuilt by dbt)
    │       └── MergeTree() — dbt atomic table swap (full rebuild) handles "upsert"
    └── Is the table a pre-aggregated accumulator with combinable states?
        └── AggregatingMergeTree()

对于 NYC Taxi 实验课:

表引擎原因
trips_rawReplacingMergeTree(_synced_at)迁移脚本重试和切换后的生产者重试都可能把同一个 trip_id 写入两次;_synced_at DEFAULT now() 确保较晚的写入胜出
fact_tripsReplacingMergeTree(updated_at)行程可能被修正;以 updated_at 作为版本
agg_hourly_zone_tripsReplacingMergeTree(updated_at)滚动重算 = upsert;以 updated_at 作为版本
dim_* 表MergeTree由 dbt 全量重载;没有部分更新
mv_hourly_revenue可刷新 MV按计划运行;每次替换整个结果

ORDER BY 设计

ORDER BY 是 ClickHouse 表中最重要的性能决策。它决定了:

  1. 主索引效率:过滤 ORDER BY 前缀列的查询可以跳过无关的数据块
  2. 压缩率:排好序的数据压缩效果更好(相近的值彼此相邻)
  3. 去重键(对 RMT/AMT 而言),只有当两行的 ORDER BY 列都相同时,它们才算重复

设计 ORDER BY 的规则:

  1. 把低基数列放在前面(例如 borough、payment_type):更多行共享同一个值,因此索引能跳过更多数据块
  2. 把高基数列放在最后(例如 trip_id、UUID):它们能收窄范围,但放在前面时压缩效果不佳
  3. 从实际的查询过滤条件推导列的选择,而不是照抄源端结构
  4. 对 RMT 表,最后一列应当是唯一行标识符(确保每个业务键只有一行)

反模式: 把源端主键照抄成 ORDER BY。如果 Snowflake 的 TRIPS_RAW 没有显式排序,照抄 Snowflake 的结构顺序(trip_id 在前)会让 ClickHouse 得到一个随机的 ORDER BY,任何分析型查询都无法跳过数据块。

fact_trips 的推导示例:

Q1–Q7 全部以某种形式过滤 pickup_at:

  • Q1:WHERE pickup_at >= ...
  • Q2:ORDER BY week, pickup_location_id
  • Q3:WHERE pickup_at >= CURRENT_DATE - 7
  • Q4:GROUP BY DATE_TRUNC('day', pickup_at)

因此 pickup_at 必须出现在 ORDER BY 中,并且应当靠前。用 toStartOfMonth(pickup_at) 作为第一列,会构造出一个粒度更粗的前缀,即使没有 PARTITION BY 子句也能实现分区级裁剪。trip_id 放在最后,用于保证 RMT 的唯一性。

结果:ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)


PARTITION BY

PARTITION BY 是可选的,并且与 ORDER BY 相互独立。它创建物理目录分区,每个分区都是一组独立的 part。

PARTITION BY toYYYYMM(pickup_at)

在以下情况使用 PARTITION BY:

  • 你需要高效地 DROP 整个时间范围(ALTER TABLE DROP PARTITION '202401')
  • 你希望 TTL 按月而不是按行生效
  • 表非常大(>1TB),且按分区的元数据有助于查询规划

不要用 PARTITION BY 来代替 ORDER BY。 一个常见错误是把 toYYYYMM(date) 放进 PARTITION BY 却把它从 ORDER BY 中省略,这会导致分区内部无法进行块级跳过。

对于 NYC Taxi 实验课,不需要 PARTITION BY,数据集为 5000 万行(压缩后约 8GB),完全在单分区的性能范围之内。


关键陷阱汇总

陷阱后果修复方法
对可变数据选错引擎重复行静默累积使用 ReplacingMergeTree + FINAL
RMT 查询缺少 FINAL合并延迟期间聚合结果偏大对 RMT 表的所有分析型查询都加上 FINAL
ORDER BY 照抄源端结构查询变慢;无法跳过数据块从实际的查询过滤条件推导 ORDER BY
高基数列排在 ORDER BY 首位索引选择性差低基数在前,高基数在后
用 sum() 而非 sumMerge() 查询 AggregateFunction 列静默产生垃圾数字对 AggregatingMergeTree 一律使用 -Merge 组合器
RMT 版本列不是单调递增旧版本随机胜出使用在更新时总被设为 now() 的时间戳

本页内容

ZH