主题
第 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 把状态变更预聚合到 AggregatingMergeTree12.9 📌 与 MySQL / PG 的对比小框
| 维度 | ClickHouse | MySQL InnoDB | PostgreSQL |
|---|---|---|---|
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:
- 解析 WHERE,用稀疏主键索引找到所有可能命中的 Part。
- 对每个命中 Part 启动一个 Mutation 任务:
- 把 Part 内所有列读出来(列存格式,要 hardlink 或拷贝其他列);
- 按 UPDATE 表达式重写需要改的列;
- 生成一个新的 Part 目录(命名带新的 mutation 版本号);
- 把旧 Part 标记 inactive,由后台慢慢回收。
- 默认异步:
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 小) | 分钟到小时 |
| 语法 | 标准 SQL | CK 专属 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) → 写入直接拒绝
↓
集群事实上不可用避免方法:
- 合并多条 Mutation:
ALTER TABLE t UPDATE a=..., b=... WHERE x; DELETE WHERE y合成一条。 - 批量延迟:业务侧攒一批,间隔 N 分钟跑一次。
- 换引擎/思路:维度修改用 ReplacingMergeTree;状态翻转用 CollapsingMergeTree;GDPR 删除用轻量级 DELETE。
- 整批修改:用 staging 表 +
REPLACE PARTITION替代逐条 Mutation。 - 监控:
SELECT count() FROM system.mutations WHERE NOT is_done报警;system.replication_queue看副本延迟。
加分项:能给出 KILL MUTATION 的紧急救场命令;能解释"已写完的新 Part 不会回滚"。
易错点:以为发起 Mutation 后等几秒就完了 —— 实际异步,命令早返回,后台还在重写几小时。
Q4:在 ClickHouse 里如何实现"按主键更新一行"?有几种方案?
考察点:对 CK 数据模型替代方案的全面认识。
标准答案:
5 种方案(按推荐度排序):
- ReplacingMergeTree + 追加新版本:sql"修改" =
ENGINE = ReplacingMergeTree(version) ORDER BY user_idINSERT新版本;查询用argMax或SELECT FINAL拿最新。 - VersionedCollapsingMergeTree:带 sign + version,支持乱序 CDC。
- CollapsingMergeTree(+1/-1):状态翻转场景。
- 轻量级 DELETE + 重新 INSERT:先删旧的,再写新的;适合偶发场景。
- 重型
ALTER TABLE UPDATE:万不得已 —— 数据极少 / 离线时段 / 单次性。
配合:
- 字典 (
Dictionary) 缓存维度表,查询时dictGet拿最新值。 - 物化视图 (
AggregatingMergeTree + MV) 把"最新状态"持续聚合。
加分项:能解释 ReplacingMergeTree 的"最终一致"特性 —— 没合并完之前 SELECT 看到多版本,要用 FINAL 或 argMax 兜底;能补 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根据情况再发补救 Mutation 或 REPLACE PARTITION。
-- 看哪些 Part 已经被处理(mutation 版本号 > x) SELECT name, mutations FROM system.parts WHERE table='...' AND active;
加分项:能补"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 适用场景:
- 一次要修改的数据是完整一个或多个分区(按月、按天的整批修改)。
- 修改逻辑复杂,靠
UPDATE ... WHERE表达不清(需要 JOIN 维表 / 重新聚合)。 - 希望操作可回滚:staging 表保留原数据,REPLACE 后随时切回去。
- 希望操作对在线读不阻塞: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 Mutation | REPLACE 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()