Skip to content

第 12 章 Mutation:UPDATE / DELETE 的真相

学习目标:彻底搞懂 ClickHouse 的 UPDATE / DELETE 为什么"看起来像 MySQL,实际完全不是";能区分重型 Mutation轻量级 DELETE;知道 system.mutations 怎么用、什么时候 KILL MUTATION、为什么 ReplacingMergeTree / CollapsingMergeTree / REPLACE PARTITION 才是 CK 里"修改"数据的正确姿势。


12.0 一句话总览

「在 ClickHouse 里,UPDATE 一行 = 重新印一本书;DELETE 一行 = 把一整本书撕了重装。所以 CK 不是『拿来频繁改的』,它是『拿来频繁查的』。」

 MySQL/PG 的 UPDATE                    ClickHouse 的 UPDATE
 ─────────────────                    ─────────────────────
                                                          
   找到 row → 改 in-place                找到匹配的 Part(s)
   写 redo / undo                        ↓
   commit                                把整个 Part 复制一份, 改字段, 写新 Part
   ✓ 毫秒级                              旧 Part 标记 inactive
                                         (后台再回收文件)
                                        ✗ 几秒到几小时, IO 巨大

12.1 ClickHouse 的 UPDATE / DELETE 不是「就地修改」

12.1.1 你以为它是这样:

sql
UPDATE users SET nickname = 'Alice2' WHERE user_id = 1;

像 MySQL 一样,找到那一行,把 nickname 字段写入新值,搞定。

12.1.2 实际它是这样:

sql
ALTER TABLE users UPDATE nickname = 'Alice2' WHERE user_id = 1;

注意:CK 没有标准的 UPDATE/DELETE 语句(22.8+ 引入了"轻量级 DELETE"是例外,下面讲)。常规修改必须写成 ALTER TABLE ... UPDATE/DELETE,俗称 Mutation(变更)。它做的事是:

   Step 1: 解析 WHERE,找到所有可能匹配的 Part。
           (用稀疏主键索引快速排除大部分 Part)
                          
   Step 2: 对每个匹配 Part,启动一个 Mutation 任务:
           - 读出这个 Part 的所有列
           - 按 UPDATE 表达式重写需要改的列
           - 写一份**新的 Part** (例如 part_0_2_1 → part_0_2_1_5)
           - 把旧 Part 标记为 inactive
                          
   Step 3: 后台回收旧 Part 的物理文件 (默认 8 分钟后)
                          
   全程异步,命令立即返回 → 但数据不会立刻"看起来"变好。

12.1.3 生活类比:「修改一本书 = 重新印一本」

   UPDATE 一行 = 你想改一本厚书的某一句话

   出版社不能用涂改液 (列存压缩了,不能局部覆盖)

   只能:把整本书拆掉、按你的修改重新打字、重新印刷、装订成一本新的。

   书架上"这本书"还在,但实际是新印的那本,旧的进了仓库。

UPDATE 一列 != 改一个字节,而是整个 Part 重写一遍。一个 50GB 的 Part,UPDATE 一行也要重写 50GB(虽然只动一列,列存其他列也要 hardlink 过来)。


12.2 跟踪 Mutation:system.mutations

sql
SELECT
    database, table,
    mutation_id,                         -- 唯一标识
    command,                             -- ALTER UPDATE 的原始语句
    create_time,
    is_done,                             -- 是否完成
    parts_to_do,                         -- 还有多少 Part 没处理
    latest_failed_part,
    latest_fail_reason
FROM system.mutations
WHERE database = 'learn_ck'
  AND NOT is_done
ORDER BY create_time DESC;

每发起一个 ALTER TABLE ... UPDATE/DELETE,CK 就在 system.mutations 里登记一条记录。is_done = 0 表示"还在跑",parts_to_do = 0 表示完成。

📌 Mutation 命令是异步的

  • 命令本身立即返回(默认 mutations_sync = 0
  • 真正的重写在后台进行
  • 想等完成再返回,设 SET mutations_sync = 1(等本节点完成)或 2(等所有副本完成)

12.3 KILL MUTATION:取消未完成的 Mutation

误发了一条扫全表的 UPDATE?赶紧杀:

sql
KILL MUTATION
WHERE database = 'learn_ck'
  AND table = 'big_table'
  AND mutation_id = 'mutation_42.txt';

注意:

  • 只能杀未完成的 Mutation。
  • 已经重写完的 Part 不会回滚(CK 没有 Undo 概念,不像 PG 有版本元组)。
  • 杀完之后,已经写好的新 Part 仍然是有效数据;只是没改的 Part 保持原样。结果:可能造成"半新半旧"的数据态,需要手动用一条新 Mutation 修补。

12.4 轻量级 DELETE(22.8+)

CK 22.8 引入了"看起来像 MySQL 的 DELETE":

sql
DELETE FROM events_log WHERE event_date < '2026-01-01';

不是重型 Mutation,而是用「墓碑列」实现:

12.4.1 内部机制

原表:每个 Part 有一个隐藏的 _row_exists 列 (UInt8)
       
DELETE FROM t WHERE x = 1;

   不重写整个 Part,只往 _row_exists 列写 0
   (列存里改一列比改全表轻得多)
   
SELECT * FROM t;

   引擎自动加 WHERE _row_exists = 1 过滤掉墓碑行
   (用户看不见已删除的行)

12.4.2 优势 vs 局限

维度轻量级 DELETE重型 Mutation DELETE
语法DELETE FROM t WHERE ...ALTER TABLE t DELETE WHERE ...
实现_row_exists 列 (墓碑)重写整个 Part
速度⚡ 几秒(小 IO)🐢 几分钟到几小时 (整 Part 重写)
物理空间❌ 不回收(墓碑还在)✅ 立即回收
真正回收后续 Merge 时一并清理立即
对查询影响轻微(多一个 WHERE 列过滤)

📌 不能滥用:连续几百条轻量级 DELETE 会让墓碑列爆炸,每次查询都要带 _row_exists=1 过滤,性能慢慢退化。轻量级 DELETE 适合"每天一两次的批量删",不适合"每秒来一条"。

12.4.3 轻量级 UPDATE(24.x+,实验性)

sql
SET allow_experimental_lightweight_update = 1;
UPDATE users SET nickname = 'Alice2' WHERE user_id = 1;

机制类似:用「补丁列」记录新值,查询时合并最新版本。截至 24.x 仍是实验性,生产慎用。


12.5 REPLACE PARTITION / EXCHANGE PARTITION:批量"修改"的最佳姿势

如果你要"修改某个月的所有数据",正确姿势是先在 staging 表里把新版本算好,再整块替换

sql
-- 1. 用 staging 表算新版本
CREATE TABLE staging.events_2026_04 AS prod.events_log;
INSERT INTO staging.events_2026_04
SELECT
    event_date, user_id,
    if(event_type = 'CLICK_OLD', 'CLICK', event_type) AS event_type,    -- 修改逻辑
    payload
FROM prod.events_log
WHERE toYYYYMM(event_date) = 202604;

-- 2. 一行命令整块替换 (秒级,仅元数据切换)
ALTER TABLE prod.events_log REPLACE PARTITION '202604' FROM staging.events_2026_04;

-- 3. 清理 staging
DROP TABLE staging.events_2026_04;

优势

  • 几秒完成(仅元数据),不阻塞读。
  • 失败可回滚 (staging 表保留)。
  • 不产生 Mutation 任务,不影响后台 Merge。

EXCHANGE PARTITION:双向原子交换两个表的同名分区(A↔B),适合灰度切换。


12.6 不可逆操作警示

12.6.1 Mutation 期间会发生什么?

  你执行: ALTER TABLE big DELETE WHERE event_date < '2026-01-01';
  
  ┌─────────────────────────────────────────────────────────────┐
  │  并发的副作用                                                  │
  ├─────────────────────────────────────────────────────────────┤
  │  ✗ 占用大量磁盘 IO (重写所有命中 Part, 临时空间 1.5x)         │
  │  ✗ 占用 CPU (压缩 / 解压 / 写入)                              │
  │  ✗ 阻塞同表的后台 Merge (Merge 要等 Mutation 完)              │
  │  ✗ 阻塞 ALTER 类 DDL (TTL 修改、ADD COLUMN 都要排队)          │
  │  ✗ 副本表会同步重做(带宽 / 各副本独立 IO 都 double)           │
  │  ✗ 查询不受影响(旧 Part 还在用),但偶发的 read 内存上升       │
  └─────────────────────────────────────────────────────────────┘

12.6.2 高并发 Mutation = 灾难

  发起 Mutation #1: 重写 30 个 Part
  发起 Mutation #2: 又重写同样 30 个 Part 的不同字段
  发起 Mutation #3: 再次重写...
  
  → CK 会按顺序串行处理,但每个 Mutation 都要把 Part 完整重写一遍。
  → 30 个 Part × 3 个 Mutation = 90 次重写
  → 后台 Merge 全停 → Part 数继续涨 → 写入开始报 Too many parts → 表挂

🔥 生产铁律单表同时挂着的 Mutation 数控制在 1~2 个。需要批量改?合并成一条 SQL:

sql
-- ❌ 烂写法
ALTER TABLE t UPDATE col_a = ... WHERE ...;
ALTER TABLE t UPDATE col_b = ... WHERE ...;
ALTER TABLE t DELETE        WHERE ...;

-- ✅ 合并写法
ALTER TABLE t
    UPDATE col_a = ..., col_b = ... WHERE cond1,
    DELETE WHERE cond2;

12.7 替代方案:CK 里"软更新 / 软删除"的正确姿势

ClickHouse 的设计哲学是「数据是事实,事实不应该被修改」。如果你"必须"频繁改数据,先想想能不能换个思路。

12.7.1 用 ReplacingMergeTree 替代 UPDATE

sql
CREATE TABLE users_dim
(
    user_id    UInt64,
    nickname   String,
    region     LowCardinality(String),
    updated_at DateTime
)
ENGINE = ReplacingMergeTree(updated_at)        -- 按 updated_at 取最新
ORDER BY user_id;

-- "修改" = 追加新版本
INSERT INTO users_dim VALUES (1, 'Alice2',  'CN', now());
INSERT INTO users_dim VALUES (1, 'Alice3',  'CN', now());

-- 查询时:
SELECT * FROM users_dim FINAL WHERE user_id = 1;
-- 或者
SELECT argMax(nickname, updated_at), argMax(region, updated_at)
FROM users_dim WHERE user_id = 1 GROUP BY user_id;

ReplacingMergeTree 在后台 Merge 时自动按主键去重,留下 updated_at 最大的那个版本。坏处是:

  • 「最终一致」—— 没合并完之前 SELECT 会看到多个版本。
  • SELECT FINAL 在线触发去重,性能较低。

12.7.2 用 CollapsingMergeTree 实现「软删除」

sql
CREATE TABLE orders_log
(
    order_id  UInt64,
    status    LowCardinality(String),
    amount    UInt64,
    sign      Int8                     -- 关键:+1 表示新增,-1 表示作废
)
ENGINE = CollapsingMergeTree(sign)
ORDER BY order_id;

INSERT INTO orders_log VALUES (1, 'CREATED', 100,  1);
INSERT INTO orders_log VALUES (1, 'CREATED', 100, -1);   -- 作废
INSERT INTO orders_log VALUES (1, 'PAID',    100,  1);   -- 新版本

-- 查询时
SELECT order_id, sumIf(amount, sign = 1) - sumIf(amount, sign = -1) AS amt
FROM orders_log GROUP BY order_id;

后台 Merge 时引擎会抵消同主键的 +1/-1 行,最终只留有效版本。适合订单状态这种"频繁翻转 + 最终需要净额"的场景。

12.7.3 VersionedCollapsingMergeTree:乱序写入也能折对

CDC(变更数据捕获)场景,多个客户端可能乱序写入新旧版本。VersionedCollapsingMergeTree 加一个 version 列保证按版本号折叠,乱序也无所谓。

12.7.4 用 TTL DELETE 做生命周期清理

sql
ALTER TABLE events_log MODIFY TTL event_date + INTERVAL 90 DAY DELETE;

不需要发起 DELETE 命令,引擎自动按时间清理(详见第 11 章)。

12.7.5 用 REPLACE PARTITION 做批量 "修改"

如本章 12.5。

12.7.6 选型总结

业务诉求推荐方案
维度表 (用户画像),需要"取最新"ReplacingMergeTree + argMax 查询
订单状态,需要"净状态"CollapsingMergeTree
CDC 同步,乱序到达VersionedCollapsingMergeTree
"删除" 90 天前的数据TTL DELETE
修改某个月所有数据REPLACE PARTITION
GDPR/合规 真正删除某用户轻量级 DELETE → 后台 Merge 清理
单条数据精确修改真没办法,上 Mutation,但控制频率

12.8 真实案例

12.8.1 用户隐私字段的"合规删除"(GDPR)

sql
-- 用户申请删除:把 user_id = 12345 的所有数据物理消除
DELETE FROM events_log WHERE user_id = 12345;       -- 轻量级 DELETE
DELETE FROM clicks_log WHERE user_id = 12345;
DELETE FROM users_dim  WHERE user_id = 12345;

-- 30 天后强制 Merge 把墓碑彻底清掉,满足 GDPR 30 天物理删除要求
OPTIMIZE TABLE events_log FINAL;

12.8.2 订单状态的频繁变更

错误做法 (CK 反模式):
   每次状态变更 → ALTER TABLE orders UPDATE status = 'PAID' WHERE order_id = ...;
   假如每秒 100 次状态变更 → 每秒触发 100 个 Mutation → 集群崩溃

正确做法:
   1. 写入只追加 (CollapsingMergeTree + sign)
   2. 状态聚合靠查询时 GROUP BY
   3. 实时大屏靠 MV 把状态变更预聚合到 AggregatingMergeTree

12.9 📌 与 MySQL / PG 的对比小框

维度ClickHouseMySQL InnoDBPostgreSQL
UPDATE 语法ALTER TABLE t UPDATE ... WHERE
(无 UPDATE 标准语句)
标准 UPDATE标准 UPDATE
修改方式重写整个 Part就地改 + Undo Log新插入元组 + 标记旧版死
修改速度秒~小时(重 IO)毫秒毫秒
事务❌ 无✅ ACID✅ ACID
DELETE 是否物理回收重型: 立即;轻量: 等 Merge等 purge 线程等 VACUUM
高并发改❌ 灾难✅ 设计目标✅ 设计目标
轻量删22.8+ DELETE FROM(墓碑)
软更新ReplacingMT / CollapsingMT 模式直接 UPDATE直接 UPDATE
适合场景写多 + 查多 + 几乎不改OLTP 频繁改OLTP 频繁改

金句

MySQL 把每行当成可变状态PG 把每行当成多版本快照ClickHouse 把每行当成不可变事实。如果你需要"修改",去 MySQL/PG;ClickHouse 是给你聚合分析的。」


12.10 本章小结

┌──────────────────────────────────────────────────────────────┐
│                     第 12 章核心要点                            │
├──────────────────────────────────────────────────────────────┤
│                                                                │
│  ① ClickHouse 没有传统 UPDATE/DELETE,只有:                    │
│     ALTER TABLE ... UPDATE / DELETE → 重型 Mutation             │
│     (异步 + 整 Part 重写)                                       │
│                                                                │
│  ② 22.8+ 引入轻量级 DELETE = 墓碑列实现,比 Mutation 快很多      │
│     但仍然不可滥用 (墓碑会让 SELECT 慢慢退化)                   │
│                                                                │
│  ③ system.mutations 跟踪 Mutation 进度                         │
│     KILL MUTATION 取消未完成任务(已写部分不回滚)               │
│                                                                │
│  ④ Mutation 期间高 IO/CPU + 阻塞 Merge → 频繁 Mutation 必崩     │
│                                                                │
│  ⑤ 替代方案矩阵:                                                │
│     - 维度表"修改" → ReplacingMergeTree + argMax                │
│     - 状态翻转    → CollapsingMergeTree (+1/-1)                │
│     - 乱序 CDC    → VersionedCollapsingMergeTree                │
│     - 生命周期    → TTL DELETE                                  │
│     - 整批修改    → REPLACE PARTITION                           │
│                                                                │
│  ⑥ MySQL/PG 把数据当可变状态;CK 把数据当不可变事实。            │
│     "需要频繁改" 是 CK 反模式,应换思路或换数据库。              │
│                                                                │
└──────────────────────────────────────────────────────────────┘

12.11 面试高频题

Q1:ClickHouse 的 UPDATE 为什么这么慢?底层做了什么?

考察点:是否真懂 Mutation 的实现机制。

标准答案

ClickHouse 的 ALTER TABLE ... UPDATE 不是就地修改,而是 Mutation

  1. 解析 WHERE,用稀疏主键索引找到所有可能命中的 Part。
  2. 对每个命中 Part 启动一个 Mutation 任务:
    • 把 Part 内所有列读出来(列存格式,要 hardlink 或拷贝其他列);
    • 按 UPDATE 表达式重写需要改的列;
    • 生成一个新的 Part 目录(命名带新的 mutation 版本号);
    • 把旧 Part 标记 inactive,由后台慢慢回收。
  3. 默认异步ALTER 命令立即返回,可在 system.mutations 跟踪进度。

慢的根本原因:

  • 即使只改一行、一列,也要把整 Part 复制重写(列存压缩文件不能局部覆盖);
  • 50GB 的 Part,UPDATE 一行也要写 50GB(其他列 hardlink 也要做元数据 + 校验);
  • 高并发 Mutation 排队 + 阻塞后台 Merge → 表挂。

加分项:能说出 mutations_sync 控制同步等待,KILL MUTATION 可取消未完成任务但不回滚已写 Part;能解释为什么 ReplacingMergeTree / CollapsingMergeTree 是更好的"修改"姿势。

易错点:以为可以用并发 UPDATE 来加速 —— 反而会让 Part 数爆炸+互相阻塞。


Q2:轻量级 DELETE(DELETE FROM t WHERE)和重型 ALTER TABLE t DELETE 有什么区别?

考察点:对 22.8+ 新机制的理解。

标准答案

维度DELETE FROM(轻量级)ALTER TABLE DELETE(Mutation)
实现写隐藏列 _row_exists = 0重写整个匹配 Part
速度秒级(IO 小)分钟到小时
语法标准 SQLCK 专属 ALTER
物理回收后续 Merge 时一并清理立即(旧 Part 退役)
查询影响引擎自动加 WHERE _row_exists=1 过滤
适用场景偶发 / 合规删除 / GDPR计划性大规模数据清理

轻量级 DELETE 的本质是把"墓碑"写进数据,查询透明跳过 —— 类似 PG 的 MVCC 删除。但 CK 没有 VACUUM,回收靠 Merge,所以滥用会导致墓碑堆积,慢慢拖慢查询

加分项:能补 lightweight_deletes_sync 控制同步语义;能说"轻量级 DELETE 不能跨副本一致"在 Replicated 场景下的特殊行为;能提"24.x 引入轻量级 UPDATE,机制类似补丁列"。

易错点:认为轻量级 DELETE 和 MySQL 的 DELETE 等价 —— 它仍然有"墓碑积累"问题,需要后台 Merge 才能彻底清理。


Q3:高并发发起 Mutation 会怎样?怎么避免?

考察点:生产事故经验。

标准答案

症状链

   每秒 100 条 ALTER UPDATE

   每条都要重写匹配 Part 的所有列

   后台 Merge 通道被占满 (与 Mutation 共享同一个池)

   Part 数因为 INSERT 继续涨, 但 Merge 跟不上

   超过 parts_to_throw_insert (默认 3000) → 写入直接拒绝

   集群事实上不可用

避免方法

  1. 合并多条 MutationALTER TABLE t UPDATE a=..., b=... WHERE x; DELETE WHERE y 合成一条。
  2. 批量延迟:业务侧攒一批,间隔 N 分钟跑一次。
  3. 换引擎/思路:维度修改用 ReplacingMergeTree;状态翻转用 CollapsingMergeTree;GDPR 删除用轻量级 DELETE。
  4. 整批修改:用 staging 表 + REPLACE PARTITION 替代逐条 Mutation。
  5. 监控SELECT count() FROM system.mutations WHERE NOT is_done 报警;system.replication_queue 看副本延迟。

加分项:能给出 KILL MUTATION 的紧急救场命令;能解释"已写完的新 Part 不会回滚"。

易错点:以为发起 Mutation 后等几秒就完了 —— 实际异步,命令早返回,后台还在重写几小时。


Q4:在 ClickHouse 里如何实现"按主键更新一行"?有几种方案?

考察点:对 CK 数据模型替代方案的全面认识。

标准答案

5 种方案(按推荐度排序):

  1. ReplacingMergeTree + 追加新版本
    sql
    ENGINE = ReplacingMergeTree(version)
    ORDER BY user_id
    "修改" = INSERT 新版本;查询用 argMaxSELECT FINAL 拿最新。
  2. VersionedCollapsingMergeTree:带 sign + version,支持乱序 CDC。
  3. CollapsingMergeTree(+1/-1):状态翻转场景。
  4. 轻量级 DELETE + 重新 INSERT:先删旧的,再写新的;适合偶发场景。
  5. 重型 ALTER TABLE UPDATE:万不得已 —— 数据极少 / 离线时段 / 单次性。

配合

  • 字典 (Dictionary) 缓存维度表,查询时 dictGet 拿最新值。
  • 物化视图 (AggregatingMergeTree + MV) 把"最新状态"持续聚合。

加分项:能解释 ReplacingMergeTree 的"最终一致"特性 —— 没合并完之前 SELECT 看到多版本,要用 FINALargMax 兜底;能补 OPTIMIZE TABLE FINAL 强制立即合并的代价。

易错点:直接照搬 MySQL 的 INSERT ... ON DUPLICATE KEY UPDATE 思维 —— CK 没有 UPSERT,硬模拟会写出各种性能陷阱。


Q5:KILL MUTATION 杀掉一个 Mutation 后,数据状态是什么样的?

考察点:对 Mutation 原子性的理解。

标准答案

  • KILL MUTATION 只能杀未完成的 Mutation 任务。
  • 已经被该 Mutation 重写的 Part:新 Part 仍然是有效数据,CK 不会回滚它们(CK 没有 Undo / 事务概念)。
  • 还没开始或正在重写中的 Part:保持原样,UPDATE 不生效。
  • 结果:数据可能呈现"半新半旧"的状态,需要后续手动用一条新 Mutation 修补或人工对账。

最佳实践

  • 大表上发起 Mutation 前先在 staging 表试跑或者先用 WHERE 1=0 LIMIT 1 估算耗时。
  • mutations_sync = 1 同步等待(小数据量场景),避免发现误发已经跑掉。
  • 紧急情况下 KILL 后立刻:
    sql
    -- 看哪些 Part 已经被处理(mutation 版本号 > x)
    SELECT name, mutations FROM system.parts WHERE table='...' AND active;
    根据情况再发补救 Mutation 或 REPLACE PARTITION。

加分项:能补"KILL 后需要确认 system.mutations.is_done = 1, latest_failed_part 不为空";提到 query_log 里能看到 Mutation 的发起者和参数。

易错点:以为 KILL MUTATION 像 MySQL 的 ROLLBACK 一样回到 Mutation 之前 —— 完全错。


Q6:什么时候应该选 REPLACE PARTITION 而不是 ALTER TABLE UPDATE

考察点:对批量修改的最佳实践。

标准答案

REPLACE PARTITION 适用场景

  1. 一次要修改的数据是完整一个或多个分区(按月、按天的整批修改)。
  2. 修改逻辑复杂,靠 UPDATE ... WHERE 表达不清(需要 JOIN 维表 / 重新聚合)。
  3. 希望操作可回滚:staging 表保留原数据,REPLACE 后随时切回去。
  4. 希望操作对在线读不阻塞:REPLACE PARTITION 仅做元数据切换(秒级),UPDATE 是后台几小时重写。

操作步骤

sql
-- 1. 在 staging 表里按期望状态写入
CREATE TABLE staging.t_2026_04 AS prod.t;
INSERT INTO staging.t_2026_04 SELECT ... FROM prod.t WHERE toYYYYMM(d)=202604;

-- 2. 整块替换 (毫秒级)
ALTER TABLE prod.t REPLACE PARTITION '202604' FROM staging.t_2026_04;

优势对比

维度UPDATE MutationREPLACE PARTITION
耗时几分钟到几小时几秒(仅元数据)
IO重写所有命中 Part仅切元数据 + hardlink
是否可回滚不可可(保留 staging)
阻塞 Merge
跨字段重算困难任意 SQL 都行

加分项:能补 EXCHANGE PARTITION 用于双向原子互换、MOVE PARTITION TO TABLE 用于跨表搬迁;能讲 Replicated 场景下 REPLACE PARTITION 的副本同步。

易错点:忘了 staging 表必须 Schema 一致 / 同 PARTITION BY 表达式。


📌 下一章预告:第 13 章我们讲 副本与分布式 —— ReplicatedMergeTree + ClickHouse Keeper 怎么协调多副本,Distributed 引擎怎么把单表查询 fan-out 到多个分片再 fan-in。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
第 12 章 · Mutation 实操对比脚本

演示 4 件事:
  1. 灌入 ~100 万行原始数据,观察初始 Part 数
  2. 重型 Mutation:   ALTER TABLE ... UPDATE   →  耗时 + Part 重写
  3. 轻量级 DELETE:   DELETE FROM ...           →  耗时 + 墓碑列
  4. REPLACE PARTITION: staging 表整块替换       →  毫秒级
  + 5. ReplacingMergeTree "软更新" 演示

前置:
    1. clickhouse-client --multiquery < ../init.sql
    2. pip install clickhouse-connect

用法:
    python3 mutation_play.py                  # 全部跑一遍
    python3 mutation_play.py --rows 500000    # 调整数据量
    python3 mutation_play.py --skip-seed      # 不重灌, 复用现有数据
"""
from __future__ import annotations

import argparse
import random
import time
from datetime import date, datetime, timedelta

import clickhouse_connect


def banner(title: str) -> None:
    print("\n" + "=" * 72)
    print(f"  {title}")
    print("=" * 72)


def show_parts(client, table: str) -> dict:
    sql = f"""
        SELECT
            count()                                AS n_parts,
            sum(rows)                              AS rows,
            formatReadableSize(sum(bytes_on_disk)) AS size,
            max(mutations_id)                      AS max_mut
        FROM (
            SELECT name, rows, bytes_on_disk,
                   length(extractAll(name, '_')) - 3 AS mutations_id
            FROM system.parts
            WHERE database='learn_ck' AND table='{table}' AND active
        )
    """
    r = client.query(sql).result_rows
    if not r or r[0][0] == 0:
        return {"n_parts": 0, "rows": 0, "size": "(none)"}
    n, rows, size, mx = r[0]
    return {"n_parts": n, "rows": rows, "size": size, "max_mut": mx}


def wait_mutation(client, table: str, mut_id: str = None, timeout: int = 120) -> dict:
    """阻塞等 Mutation 完成"""
    start = time.time()
    while time.time() - start < timeout:
        cond = f"AND mutation_id='{mut_id}'" if mut_id else ""
        sql = f"""
            SELECT mutation_id, is_done, parts_to_do, latest_fail_reason
            FROM system.mutations
            WHERE database='learn_ck' AND table='{table}' {cond}
            ORDER BY create_time DESC LIMIT 1
        """
        r = client.query(sql).result_rows
        if not r:
            time.sleep(0.5)
            continue
        m_id, done, todo, fail = r[0]
        if done:
            return {"mutation_id": m_id, "done": True, "fail": fail}
        time.sleep(0.5)
    return {"done": False, "timeout": True}


def seed(client, total: int) -> None:
    cols = ["event_date", "event_time", "user_id", "event_type", "payload", "score"]
    today = date.today()
    batch_size = 50_000
    sent = 0
    while sent < total:
        n = min(batch_size, total - sent)
        rows = []
        for _ in range(n):
            d = today - timedelta(days=random.randint(0, 60))
            rows.append((
                d,
                datetime(d.year, d.month, d.day, random.randint(0,23), random.randint(0,59)),
                random.randint(1, 100_000),
                random.choice(["click", "view", "buy", "share"]),
                "P" * random.randint(50, 200),
                random.randint(0, 1000),
            ))
        client.insert("learn_ck.mut_events", rows, column_names=cols)
        sent += n


def demo_mutation(client) -> None:
    banner("Demo 2: 重型 ALTER TABLE UPDATE  → 整 Part 重写")
    before = show_parts(client, "mut_events")
    print(f"  Mutation 前: parts={before['n_parts']}, rows={before['rows']}, size={before['size']}")

    sql = "ALTER TABLE learn_ck.mut_events UPDATE score = score + 1 WHERE event_type = 'click'"
    print(f"\n  执行: {sql}")
    t0 = time.time()
    client.command(sql)
    print(f"  命令立即返回, 耗时 {1000*(time.time()-t0):.1f}ms (异步, 实际还在跑)")

    print("\n  跟踪 system.mutations:")
    res = wait_mutation(client, "mut_events")
    cost = time.time() - t0
    print(f"  Mutation 完成: id={res.get('mutation_id')}  总耗时={cost:.2f}s  fail={res.get('fail')}")

    after = show_parts(client, "mut_events")
    print(f"\n  Mutation 后: parts={after['n_parts']}, rows={after['rows']}, size={after['size']}")
    print(f"  注意 Part 名字带新版本号 (mutation_id), 旧 Part 已 inactive")


def demo_lightweight_delete(client) -> None:
    banner("Demo 3: 轻量级 DELETE FROM (墓碑列) ")
    before = show_parts(client, "mut_events")
    print(f"  DELETE 前: parts={before['n_parts']}, rows={before['rows']}")

    sql = "DELETE FROM learn_ck.mut_events WHERE event_type = 'view' AND score < 100"
    print(f"\n  执行: {sql}")
    t0 = time.time()
    try:
        client.command(sql)
        cost = time.time() - t0
        print(f"  耗时 {1000*cost:.1f}ms (轻量级 DELETE, 写墓碑列, 不重写 Part)")
    except Exception as e:
        print(f"  ⚠ 服务端可能未开启 lightweight delete 或版本 < 22.8: {e}")
        return

    # 看可见行数 (引擎自动加 _row_exists=1 过滤)
    visible = client.query("SELECT count() FROM learn_ck.mut_events").result_rows[0][0]
    after = show_parts(client, "mut_events")
    print(f"\n  DELETE 后:")
    print(f"    SELECT count() (可见行) = {visible}")
    print(f"    system.parts 物理行数 = {after['rows']}  (墓碑还在, 等 Merge 清)")
    print(f"    parts={after['n_parts']}, size={after['size']}")
    print(f"    差值 = {after['rows'] - visible}  (这些是墓碑行)")


def demo_replace_partition(client) -> None:
    banner("Demo 4: REPLACE PARTITION (毫秒级整块替换)")

    # 找一个有数据的分区
    parts = client.query(
        "SELECT DISTINCT partition FROM system.parts "
        "WHERE database='learn_ck' AND table='mut_events' AND active "
        "ORDER BY partition LIMIT 1"
    ).result_rows
    if not parts:
        print("  没数据, 跳过")
        return
    target = parts[0][0]
    print(f"  目标分区: {target}")

    # 1. 清 staging
    client.command("TRUNCATE TABLE learn_ck.mut_events_staging")

    # 2. 在 staging 里准备"修改后"的数据
    sql_stage = f"""
        INSERT INTO learn_ck.mut_events_staging
        SELECT event_date, event_time, user_id,
               'rewritten' AS event_type,
               concat('NEW-', payload) AS payload,
               score * 10 AS score
        FROM learn_ck.mut_events
        WHERE toYYYYMM(event_date) = {target}
    """
    print(f"\n  Step 1: 在 staging 表里准备新版本数据...")
    t0 = time.time()
    client.command(sql_stage)
    print(f"    耗时 {1000*(time.time()-t0):.0f}ms (这部分可以离线慢慢跑)")

    # 3. REPLACE
    print(f"\n  Step 2: ALTER TABLE ... REPLACE PARTITION '{target}' FROM mut_events_staging")
    t0 = time.time()
    client.command(
        f"ALTER TABLE learn_ck.mut_events REPLACE PARTITION '{target}' "
        f"FROM learn_ck.mut_events_staging"
    )
    print(f"    耗时 {1000*(time.time()-t0):.0f}ms  ★ 仅元数据切换!")

    # 验证
    sample = client.query(
        f"SELECT event_type, payload FROM learn_ck.mut_events "
        f"WHERE toYYYYMM(event_date) = {target} LIMIT 1"
    ).result_rows
    print(f"\n  验证: {sample}")


def demo_replacing_mt(client) -> None:
    banner("Demo 5: ReplacingMergeTree 软更新 - INSERT 替代 UPDATE")
    client.command("TRUNCATE TABLE learn_ck.mut_users_dim")

    # 三次"修改"同一个 user_id, 全用 INSERT
    rows = [
        (1, "Alice",   "CN",    datetime.now() - timedelta(days=2)),
        (1, "Alice2",  "CN-SH", datetime.now() - timedelta(days=1)),
        (1, "Alice3",  "CN-BJ", datetime.now()),
    ]
    client.insert("learn_ck.mut_users_dim", rows,
                  column_names=["user_id","nickname","region","updated_at"])

    print("  原始 INSERT 三次, 没 Merge 时:")
    for r in client.query("SELECT * FROM learn_ck.mut_users_dim WHERE user_id=1").result_rows:
        print(f"    {r}")

    print("\n  SELECT FINAL 在线去重:")
    for r in client.query("SELECT * FROM learn_ck.mut_users_dim FINAL WHERE user_id=1").result_rows:
        print(f"    {r}")

    print("\n  推荐: 用 argMax 在线 (不需要 FINAL):")
    sql = """
        SELECT user_id,
               argMax(nickname, updated_at) AS nickname,
               argMax(region,   updated_at) AS region,
               max(updated_at) AS updated_at
        FROM learn_ck.mut_users_dim
        WHERE user_id = 1
        GROUP BY user_id
    """
    for r in client.query(sql).result_rows:
        print(f"    {r}")

    print("\n  强制 OPTIMIZE FINAL 后:")
    client.command("OPTIMIZE TABLE learn_ck.mut_users_dim FINAL")
    for r in client.query("SELECT * FROM learn_ck.mut_users_dim WHERE user_id=1").result_rows:
        print(f"    {r}  ← 物理上只剩最新版本")


def main() -> None:
    parser = argparse.ArgumentParser()
    parser.add_argument("--host", default="127.0.0.1")
    parser.add_argument("--port", type=int, default=8123)
    parser.add_argument("--user", default="default")
    parser.add_argument("--password", default="")
    parser.add_argument("--rows", type=int, default=300_000)
    parser.add_argument("--skip-seed", action="store_true")
    parser.add_argument("--skip-mutation", action="store_true",
                        help="跳过 demo 2 (重型 Mutation 较慢)")
    args = parser.parse_args()

    client = clickhouse_connect.get_client(
        host=args.host, port=args.port,
        username=args.user, password=args.password,
    )

    if not args.skip_seed:
        banner("Demo 1: 灌入数据")
        client.command("TRUNCATE TABLE learn_ck.mut_events")
        t0 = time.time()
        seed(client, args.rows)
        cost = time.time() - t0
        s = show_parts(client, "mut_events")
        print(f"  灌入 {args.rows} 行  cost={cost:.1f}s  parts={s['n_parts']} size={s['size']}")

    if not args.skip_mutation:
        demo_mutation(client)
    demo_lightweight_delete(client)
    demo_replace_partition(client)
    demo_replacing_mt(client)

    banner("总结")
    print("""
  ┌─────────────────────────────────────────────────────────────┐
  │  操作                          典型耗时    适用场景           │
  ├─────────────────────────────────────────────────────────────┤
  │  ALTER TABLE UPDATE (Mutation) 几分钟~小时  万不得已          │
  │  DELETE FROM (轻量级)           几秒        合规/低频删       │
  │  REPLACE PARTITION              几毫秒      整批修改最佳      │
  │  ReplacingMergeTree + INSERT    毫秒        维度表"修改"      │
  │  CollapsingMergeTree + sign     毫秒        状态翻转          │
  │  TTL DELETE                     0 (自动)    生命周期管理      │
  └─────────────────────────────────────────────────────────────┘
""")


if __name__ == "__main__":
    main()

mutation_play.py ↗