前置:/clickhouse-data-ingestion/(数据导入与格式)、/clickhouse-merge-tree-principle/(Part 落盘与合并原理)、/clickhouse-kafka-engine/(Kafka 表引擎)。
目录
- 1. 写入全链路:从 INSERT 到 Part 落盘
- 2. 批量 INSERT:一次 1 行 vs 一次十万行
- 3. async_insert:异步缓冲与合并写入
- 4. 吞吐瓶颈诊断:query_log 与 ProfileEvents
- 5. 分区过多与 parts 膨胀的代价
- 6. 写入幂等:deduplicate 与 ReplacingMergeTree
- 7. Kafka 消费写入优化
- 8. 分区策略与写入并发
- 9. 吞吐测试方法论与基准
- 10. 速查表与一句话记忆
- 延伸阅读
1. 写入全链路:从 INSERT 到 Part 落盘
理解写入瓶颈,先看懂一条数据从客户端到磁盘经历了什么。ClickHouse 的写入不是逐行落盘,而是攒成块后一次刷盘。
写入链路:
□ 客户端发 INSERT → 按列组装内存 Block
□ 攒够 min_insert_block_size_rows(约 100 万行)
□ Block 压缩编码 → 写列文件 + 索引文件
□ 落盘为不可变 Part,后台异步 merge
为什么单条 INSERT 慢:
□ 每条 1 行 → Block 永远攒不满
□ 每 INSERT 建 Part、写索引、fsync,海量小 part
结论:写入单元 = 批量 Block,不是单行
-- 观察最近写入操作
SELECT event_time, query, written_rows, query_duration_ms
FROM system.query_log
WHERE query_kind = 'Insert'
ORDER BY event_time DESC LIMIT 20;
-- 当前 parts 总数:小 part 多的信号
SELECT count() AS parts, sum(rows) AS total_rows
FROM system.parts WHERE active = 1 AND table = 'events';
工程要点:ClickHouse 的写入是按块(Block)落盘的,最小落盘块默认约 100 万行;单条 INSERT 慢是因为攒不满块、Part 过多、网络往返多——优化写入的第一原则是把「次数」降下来,把「单次行数」提上去。
2. 批量 INSERT:一次 1 行 vs 一次十万行
批量 INSERT 是最直接、最可靠的提速手段,不需要任何服务器端魔法。
批量的收益:
□ 单次 Block 越大 → 压缩率越高
□ 写入次数少 → 建 part 开销分摊
□ 网络往返少 → RTT 不再是瓶颈
□ Part 数量少 → 后续 merge 压力小
经验值:
□ 单批 10 万~100 万行(或 100MB~1GB)
□ 用一条 SQL 带多行 VALUES
□ 避免逐行 insertOne
□ 批次太大内存上升 → 按行数/字节拆批
-- 正例:一条 SQL 批量插入
INSERT INTO events (event_time, user_id, event_type, value)
VALUES
('2026-09-30 10:00:00', 1, 'view', 10),
('2026-09-30 10:00:01', 2, 'click', 20),
...; -- 一次带十万行
-- 更优:从文件批量导入
INSERT INTO events SELECT * FROM file('events.csv', 'CSVWithNames');
工程要点:批量 INSERT 把写入次数降低几个数量级——单批 10 万~100 万行、一条 SQL、避免逐行插;服务器端 zero-config,却同时改善压缩率、Part 数量与网络开销,是写入调优的第一优先级动作。
3. async_insert:异步缓冲与合并写入
如果业务只能逐条发送(如 SDK 埋点、Logstash 转发),可以用 async_insert 在服务器端把多条小 INSERT 攒成大块。
async_insert 原理:
□ 开启后 INSERT 先进内存缓冲
□ 服务器按超时/大小把缓冲合并成块落盘
□ 对外表现为「每条成功」,实则批量落盘
关键设置:
□ async_insert = 1:开启异步缓冲
□ wait_for_async_insert = 1:等待真正落盘
□ async_insert_max_data_size:缓冲上限(1MB)
□ async_insert_busy_timeout_ms:攒批时间窗(1s)
权衡:吞吐高 RTT 低;wait=0 崩溃时缓冲数据可能丢失
SET async_insert = 1;
SET wait_for_async_insert = 1;
-- 多条小 INSERT 会被服务器攒批落盘
INSERT INTO events (event_time, user_id, event_type, value) VALUES ('2026-09-30 10:00:00', 1, 'view', 10);
INSERT INTO events (event_time, user_id, event_type, value) VALUES ('2026-09-30 10:00:00', 1, 'click', 20);
-- 观察异步缓冲产生的落盘效果
SELECT query, written_rows, query_duration_ms, memory_usage
FROM system.query_log
WHERE query_kind = 'Insert' AND event_time > now() - INTERVAL 10 MINUTE;
工程要点:async_insert 把服务器端的多次小 INSERT 攒成一个块再落盘,让逐条写入也享受批量红利;但要用 wait_for_async_insert=1 换取可查性,同时接受缓冲数据在崩溃时可能丢失——它优化的是「频率高」而非「总量大」的场景。
4. 吞吐瓶颈诊断:query_log 与 ProfileEvents
写入慢先别改配置,先量化瓶颈在哪一环。system.query_log 记录了每次 INSERT 的完整画像。
诊断指标:
□ written_rows / written_bytes:单次写入量
□ query_duration_ms:写入耗时
□ memory_usage:写入内存峰值
看什么:
□ 单批行数小 → 批太小,合并次数多
□ 耗时与行数不成比例 → 网络/磁盘 fsync
□ parts 增长快 → merge 追不上
工具:system.processes / system.merges / system.part_log
-- 最近 1 小时每次 INSERT 的画像
SELECT query_duration_ms, written_rows, written_bytes,
ProfileEvents['InsertedRows'] AS inserted, memory_usage
FROM system.query_log
WHERE query_kind = 'Insert' AND event_time > now() - INTERVAL 1 HOUR
ORDER BY query_duration_ms DESC LIMIT 20;
-- 看是否有 merge 长期追不上
SELECT table, is_mutation, elapsed, rows_written
FROM system.merges ORDER BY elapsed DESC LIMIT 10;
工程要点:写入瓶颈诊断看 system.query_log 的单次写入量×次数×耗时三个维度,配合 system.merges 看 merge 是否追得上;当「parts 增速 > merge 消化速度」时,真正的问题不在写入本身,而在批次太小或分区过细。
5. 分区过多与 parts 膨胀的代价
写入吞吐经常被分区粒度拖垮:分区越细,同样数据产生的 part 越多,merge 越忙,查询越碎。
parts 膨胀的传导链:
□ 每批数据进入所属分区 → 产生一个 part
□ 分区多 → part 基数大 → merge 排队
□ part 多 → 查询要读更多文件头
□ 极端:分区数 > 查询线程数 → 并行浪费
分区过多场景:按小时/天分区、每分区只写几十行
控制手段:
□ 分区粒度对齐查询与 TTL 粒度(天/月)
□ 控制单分区 part 数
□ OPTIMIZE FINAL 主动合并不该频繁做
-- 看各分区 part 数与行数分布
SELECT partition, count() AS parts, sum(rows) AS rows,
sum(bytes_on_disk) AS bytes
FROM system.parts
WHERE active = 1 AND table = 'events'
GROUP BY partition ORDER BY parts DESC LIMIT 20;
-- 主动合并(应急手段,非日常)
OPTIMIZE TABLE events FINAL;
工程要点:分区过细与 parts 膨胀是写入吞吐的隐形杀手——每个分区每次写入都产生一个 part,分区越多 part 越多,merge 后台永远追不上,查询也被拖慢;解法是分区粒度对齐真实查询与 TTL 粒度,让单分区拥有足够大的 part。
6. 写入幂等:deduplicate 与 ReplacingMergeTree
写入链路会有重试(网络超时、Kafka 至少一次),没有幂等机制就会产生重复数据。MergeTree 家族提供了两种兜底思路。
去重方案:
□ ReplicatedMergeTree 自带 block 去重
→ 相同 part 名只落盘一次
□ ReplacingMergeTree:按排序键去重
→ 相同 ORDER BY 键只留最新(按版本/时间)
→ 合并时生效,查询用 FINAL
适用边界:
□ block 去重防「同一批数据重复插入」
□ Replacing 去重防「同 key 多次更新累积」
注意:去重靠排序键/版本,不是任意列
-- ReplacingMergeTree:以 version 取最新
CREATE TABLE events_dedup (
event_time DateTime, user_id UInt64,
event_type String, version UInt32
) ENGINE = ReplacingMergeTree(version)
PARTITION BY toYYYYMM(event_time)
ORDER BY (user_id, event_type);
-- 重复写入同一 key:合并后只留 version 最大者
INSERT INTO events_dedup VALUES ('2026-09-30 10:00:00', 1, 'view', 1);
INSERT INTO events_dedup VALUES ('2026-09-30 10:00:00', 1, 'view', 2);
-- 强制看最新态
SELECT * FROM events_dedup FINAL;
工程要点:写入幂等分两层——block 级去重(ReplicatedMergeTree 对同批 Part 只落盘一次)与 行级去重(ReplacingMergeTree 按 ORDER BY 键 + 版本列保留最新);重试型管道至少要二选一,生产上推荐「复制表 + 版本列」双保险。
7. Kafka 消费写入优化
Kafka → ClickHouse 是最常见的实时写入管道,它的吞吐瓶颈往往不在 ClickHouse,而在消费批的大小与频率。
Kafka 写入优化要点:
□ 增大单批拉取:每批攒几万~几十万行再写
□ 用 INSERT ... SELECT 从 kafka 引擎表装载
□ 并行消费:分区数 = 消费并行度
Materialized View 管道:
□ Kafka 引擎表 + 物化视图 → 目标 MergeTree
□ 视图「消费即写」,频率由 kafka_max_block_size 控制
注意:消费速度 > 写入速度 → 堆积,看 system.kafka_consumers lag
-- Kafka 引擎表:原始层
CREATE TABLE kafka_events (
event_time DateTime, user_id UInt64, event_type String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'events',
kafka_group_name = 'ch_events',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
-- 物化视图:批量转存到 MergeTree
CREATE MATERIALIZED VIEW mv_events TO events AS
SELECT * FROM kafka_events;
-- 检查消费进度与堆积
SELECT topic, consumer, partition, last_pulled_offset
FROM system.kafka_consumers WHERE topic = 'events';
工程要点:Kafka→ClickHouse 管道的吞吐取决于单批消费规模——用 Kafka 引擎表 + 物化视图批量转存、把 kafka_num_consumers 对齐分区数、攒批写入目标表;消费 lag 是第一个要盯的指标,堆积通常源于「批太小 + 写入次数太多」。
8. 分区策略与写入并发
写入吞吐还和分区键的写入分布有关:如果所有数据都写进同一个分区,单个 part 的 merge 与写入会互相竞争。
并发与分区的相互作用:
□ 并发 INSERT 到同一分区 → 更多 part
□ 均匀分区键 → 写入分散,merge 也分散
□ 极端热分区(当天日期)→ 热点写入
策略:
□ 时间分区天然「写最新分区」→ 热点不可避免
□ 避免高基数随机键分区 → 每批都进新分区
□ 大批装载用 max_insert_threads 并行
□ 写入前按分区键排序 → part 更整
-- 控制并行写入线程
SET max_insert_threads = 4;
-- 大批装载:按分区键排序后再并行写(减少 part 碎片)
-- 例如把输入按 event_time 排序,再分段 INSERT
SELECT partition, count() AS parts, sum(rows) AS rows
FROM system.part_log
WHERE table = 'events' AND event_time > now() - INTERVAL 1 DAY
GROUP BY partition;
工程要点:分区策略决定写入的并发格局——时间分区带来「写热点分区」但历史分区稳定,随机键分区则让每批都产生新 part;写入前按分区键排序、用 max_insert_threads 并行装载,能让 part 更整、merge 更顺。
9. 吞吐测试方法论与基准
调优要有可复现的度量。写一套标准测试,对比改动前后的写入吞吐与 parts 状态。
测试方法论:
□ 固定数据量与机器,一次只改一个变量
□ 指标:行/秒、写入耗时;辅指标:parts 数、merge 耗时
测试矩阵:
□ 基准:逐行 INSERT(最差基线)
□ 批量 1k / 10k / 100k / 1M 行每批
□ async_insert 开关、分区粒度(天 vs 月)、并行线程数
判定标准:吞吐提升同时 parts 数不爆炸
-- 生成测试数据并批量装载
INSERT INTO events
SELECT now() - rand() % 86400, rand() % 1000000,
['view', 'click', 'buy'][rand() % 3 + 1], rand() % 1000
FROM numbers(100000000);
-- 用 query_log 量化本次装载(行/秒)
SELECT count() AS insert_count, sum(written_rows) AS total_rows,
sum(written_rows) / (sum(query_duration_ms) / 1000) AS rows_per_s
FROM system.query_log
WHERE query_kind = 'Insert' AND query ILIKE '%numbers%'
AND event_time > now() - INTERVAL 1 HOUR;
-- 装载后 parts 健康度
SELECT count() AS parts, sum(rows) AS rows, max(rows) AS max_part_rows
FROM system.parts WHERE active = 1 AND table = 'events';
工程要点:写入调优要一次只改一个变量并记录「吞吐 + parts 数 + merge 耗时」三组数字——用 numbers() 造数、query_log 量化行/秒、装载后检查 parts 健康度,才能证明「快」不是错觉、也没把读性能拖垮。
10. 速查表与一句话记忆
把全文压成可对照的清单。
写入优化清单:
□ 批次:单批 10 万~100 万行,一条 SQL
□ 异步:逐条场景开 async_insert,wait=1
□ 分区:粒度对齐查询与 TTL,避免过细
□ 并发:max_insert_threads + 按分区键排序
□ 幂等:Replicated 去重 + Replacing 版本列
□ Kafka:大 batch 消费 + 物化视图转存
□ 监控:query_log 画像 + system.merges + lag
一句记忆:写入快 = 次数少、批次大、part 整
-- 写入优化的最小体检包
SELECT query_duration_ms, written_rows, memory_usage
FROM system.query_log WHERE query_kind = 'Insert'
ORDER BY event_time DESC LIMIT 20;
SELECT count(), sum(rows) FROM system.parts WHERE active = 1 AND table = 'events';
工程要点:写入吞吐的终极公式是**「次数少、批次大、part 整」**——批量 INSERT 解决次数、async_insert 解决逐条、分区策略与排序解决 part 碎片、幂等机制兜底重试;监控三件套(query_log、system.merges、parts 数)随时验证方向正确。
延伸阅读
- /clickhouse-data-ingestion/ — 数据导入格式与大批量装载方式
- /clickhouse-merge-tree-principle/ — Part 生命周期与合并的底层原理
- /clickhouse-kafka-engine/ — Kafka 表引擎与实时写入管道
- /clickhouse-replicated-tables-disaster-recovery/ — 复制表的写入去重与数据一致性
- /clickhouse-schema-modeling-best-practices/ — 分区键与排序键建模实践
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。