这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 写入优化的核心是减少 part 数量、控制合并压力、避免热点分片和保证失败可重试。
26.1 写入模型
小 INSERT
-> 多个新 part
-> 后台频繁合并
-> CPU / IO / 元数据压力
优化后:
写入端缓冲
-> 大 block
-> 少量 part
-> 合并压力可控
推荐:
| 指标 | 参考 |
|---|---|
| 单批行数 | 1 万到 10 万以上 |
| 写入频率 | 每秒到每数秒 |
| 单批数据 | 根据行宽压测确定 |
| 并发 | 固定有限并发 |
| 失败 | 指数退避重试 |
| 死信 | 保存不可恢复批次 |
26.2 合并批次
应用端伪代码:
type Buffer struct {
rows []Event
maxRows int
maxInterval time.Duration
}
func (b *Buffer) Add(e Event, flush func([]Event) error) error {
b.rows = append(b.rows, e)
if len(b.rows) >= b.maxRows {
return flush(b.rows)
}
return nil
}
定时器触发:
ticker := time.NewTicker(2 * time.Second)
for range ticker.C {
if len(buffer.rows) > 0 {
flush(buffer.Take())
}
}
关闭服务时必须执行 final flush。
26.3 控制并发
并发写入不是越高越好:
- 每批都会生成 part;
- 合并线程有限;
- 高并发会放大内存和网络;
- 分布式写入增加转发压力;
- 分区热点会集中到单节点。
建议写入端使用固定 worker 池,并按表、分区、分片维度限流。
26.4 分区控制
避免一次写入覆盖过多分区:
1000 行
覆盖 86400 个每分钟分区
这会生成大量小 part。分区粒度应与数据量匹配,历史迟到数据可单独低并发补写。
检查分区 part:
SELECT table, partition, count() AS parts, sum(rows) AS rows
FROM system.parts
WHERE active
GROUP BY table, partition
ORDER BY parts DESC;
26.5 异步插入
ClickHouse 提供异步插入设置,适合无法立刻改造的客户端:
INSERT INTO analytics.events_local
SETTINGS async_insert = 1, wait_for_async_insert = 1
FORMAT JSONEachRow
...
它可以让服务端聚合部分小插入,但会增加延迟,并有故障语义需要评估。更推荐在写入端或消息层完成聚合。
26.6 批量文件导入
clickhouse-client --query "
INSERT INTO analytics.events_local FORMAT CSVWithNames
" < events.csv
大文件可拆分:
split -l 1000000 events.csv part_
for f in part_*; do
clickhouse-client --query "INSERT INTO analytics.events_local FORMAT CSVWithNames" < "$f"
done
补历史数据时控制并发,避免影响实时查询和合并。
26.7 分布式写入
两种方式:
| 方式 | 特点 |
|---|---|
| 写 Distributed 表 | 简单统一,可能异步转发 |
| 应用直连分片 | 更可控,需要维护拓扑 |
如果写 Distributed 表,要监控:
SELECT database, table, error_count, data_files, bytes_to_send
FROM system.distribution_queue;
队列积压时不要盲目重启,先检查目标分片、网络和磁盘。
26.8 写入 Schema 优化
- 减少超宽列;
- 高频字段使用窄类型;
- 低基数字符串使用 LowCardinality;
- 复杂 JSON 先在写入端展开;
- 避免过多 Nullable;
- 合理使用 Codec。
示例:
CREATE TABLE analytics.events_compact
(
event_date Date,
event_type LowCardinality(String),
platform LowCardinality(String),
city_id UInt32,
user_id UInt64,
amount Decimal64(2)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, city_id, user_id);
26.9 监控指标
| 指标 | 含义 |
|---|---|
| insert rows / bytes | 写入吞吐 |
| failed queries | 写入失败 |
| parts per partition | 合并压力 |
| merges count / elapsed | 合并负载 |
| disk IO | 磁盘瓶颈 |
| CPU iowait | IO 等待 |
| distribution queue | 分布式积压 |
| too many parts | 关键告警 |
查看合并:
SELECT database, table, elapsed, num_parts, total_size
FROM system.merges;
26.10 排查 Too many parts
SELECT table, partition, count() AS parts
FROM system.parts
WHERE active
GROUP BY table, partition
ORDER BY parts DESC
LIMIT 20;
处理:
- 降低写入频率;
- 增大批次;
- 减少分区数;
- 降低补数并发;
- 检查磁盘和合并线程;
- 暂停低优先级写入;
- 扩容。
本章小结
写入优化优先发生在写入端:微批、固定并发、稳定分区和失败重试。ClickHouse 侧配合合理 schema、分区和监控。目标不是追求单次最快,而是让 part 生成速度长期低于合并消化速度。
思考题
- 为什么小 INSERT 会拖慢查询?
- 应用缓冲和异步插入如何取舍?
- 分布式写入队列积压如何处理?
- 什么情况下需要降低补数并发?
- 如何判断集群写入容量已到瓶颈?