这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 ClickHouse 的写入模型与 OLTP 数据库差异很大:它偏好大批量、稳定频率、可重试的追加写入。本章覆盖 INSERT、文件导入、外部表、批量写入策略和常见写入故障。
5.1 写入链路
INSERT / HTTP / Kafka
|
v
生成一个新 part
|
v
分区目录 + 每列文件 + 索引文件 + 校验文件
|
v
后台 merge 合并小 part
每一次小 INSERT 通常都会生成新的 part。part 数量过多后,CPU 会大量消耗在合并、定位和元数据管理上,最终可能触发 Too many parts。
推荐写入:
| 指标 | 参考 |
|---|---|
| 单批行数 | 1 万到 10 万以上 |
| 写入频率 | 每秒到每数秒一批 |
| 批次字段 | 尽量完整、类型稳定 |
| 并发 | 控制在写入端可管理范围 |
| 失败 | 按批次重试,必要时幂等 |
5.2 INSERT 语句
INSERT INTO analytics.events_local
VALUES
('2026-08-25', '2026-08-25 10:00:00', 'E001', 1001, 'view', 1, 0.00),
('2026-08-25', '2026-08-25 10:01:00', 'E002', 1002, 'pay', 1, 199.00);
指定列写入:
INSERT INTO analytics.events_local
(event_date, event_time, event_id, user_id, event_type, city_id, amount)
VALUES
('2026-08-25', '2026-08-25 10:02:00', 'E003', 1003, 'view', 2, 0);
从查询写入:
INSERT INTO analytics.events_backup
SELECT * FROM analytics.events_local
WHERE event_date = '2026-08-25';
注意:INSERT INTO ... SELECT 会占用源表读取和目标表写入资源,大范围数据迁移应分批执行。
5.3 常见输入格式
JSONEachRow:
{"event_date":"2026-08-25","event_time":"2026-08-25 10:00:00","event_id":"E001","user_id":1001,"event_type":"view","city_id":1,"amount":0}
{"event_date":"2026-08-25","event_time":"2026-08-25 10:01:00","event_id":"E002","user_id":1002,"event_type":"pay","city_id":1,"amount":199.00}
写入:
clickhouse-client --query "
INSERT INTO analytics.events_local FORMAT JSONEachRow
" < events.jsonl
CSV:
clickhouse-client --date_time_input_format=best_effort --query "
INSERT INTO analytics.events_local FORMAT CSVWithNames
" < events.csv
Parquet:
clickhouse-client --query "
INSERT INTO analytics.events_local FORMAT Parquet
" < events.parquet
5.4 HTTP 写入
curl -sS 'http://127.0.0.1:8123/' \
--data-binary 'INSERT INTO analytics.events_local FORMAT JSONEachRow
{"event_date":"2026-08-25","event_time":"2026-08-25 10:03:00","event_id":"E004","user_id":1004,"event_type":"view","city_id":3,"amount":0}
'
使用 query_id 便于追踪:
curl -sS 'http://127.0.0.1:8123/?query_id=load-20260825-0001' \
--data-binary @events.jsonl
HTTP 适合服务端导入;高吞吐应用写入通常优先 Native 协议或 Kafka。
5.5 从文件和对象存储读取
查询本地文件:
SELECT count()
FROM file('events.jsonl', JSONEachRow, 'event_date Date, event_time DateTime, event_id String');
查询 URL:
SELECT *
FROM url(
'https://example.com/events.jsonl',
JSONEachRow,
'event_date Date, user_id UInt64'
);
查询 S3:
SELECT count()
FROM s3(
'https://bucket.s3.amazonaws.com/path/events.parquet',
'Parquet',
'event_date Date, user_id UInt64'
);
访问外部资源要配置访问控制、网络策略、凭据和超时,不要让任意用户直接读取内部地址。
5.6 从 MySQL 导入
MySQL 引擎:
CREATE TABLE mysql_orders
ENGINE = MySQL
AS SELECT
order_id,
created_at,
status,
amount
FROM mysql('mysql-host:3306', 'shop', 'orders', 'user', 'password');
写入 ClickHouse:
INSERT INTO analytics.orders_local
SELECT
toDate(created_at) AS order_date,
order_id,
status,
amount,
toUInt64(unixTimestamp(created_at)) AS version
FROM mysql_orders
WHERE updated_at >= now() - INTERVAL 1 HOUR;
增量导入要解决:
- 时间基准;
- 时区;
- 更新时间回拨;
- 删除同步;
- 主键映射;
- 任务失败重放。
5.7 应用批量写入
简化伪代码:
queue = memory buffer
last_flush = now
on event:
queue.add(event)
if queue.size >= 50000 or now - last_flush >= 2s:
flush(queue)
flush(queue):
batch = queue.drain()
retry with backoff
on permanent failure:
write to local dead letter or send to Kafka retry topic
Go 示例核心结构:
type Writer struct {
conn clickhouse.Conn
batch *clickhouse.Batch
rows int
ctx context.Context
}
func (w *Writer) Add(e Event) error {
if err := w.batch.Append(
e.EventDate, e.EventTime, e.EventID,
e.UserID, e.EventType, e.CityID, e.Amount,
); err != nil {
return err
}
w.rows++
return nil
}
func (w *Writer) Flush() error {
if w.rows == 0 {
return nil
}
return w.batch.Send()
}
5.8 写入幂等
ClickHouse 没有像 OLTP 主键那样提供全局实时唯一约束。常用方案:
| 方案 | 做法 | 适合 |
|---|---|---|
| 批次 ID | 同一重试批次复用 block ID | 应用重试 |
| ReplacingMergeTree | 排序键 + 版本 | 订单、用户状态 |
| 事件表 | append + 查询去重 | 行为事件 |
| 外部幂等 | Kafka transactional id 或 ETL 水位 | 复杂链路 |
ReplacingMergeTree 只能降低合并后的物理重复,不能把查询实时变成唯一表。
5.9 删除和更新
轻量删除:
DELETE FROM analytics.events_local
WHERE event_date = '2026-08-01';
变更:
ALTER TABLE analytics.events_local
UPDATE amount = amount * 2
WHERE event_date = '2026-08-01' AND city_id = 1;
这类操作会生成变更 part,代价较高。分析表更常用的做法是:
- 修正上游后重写指定分区;
- 用替换表切换;
- 追加修正记录;
- 在查询层过滤无效标记。
5.10 写入问题排查
查看异常 part:
SELECT
database,
table,
partition,
name,
active,
rows,
bytes_on_disk
FROM system.parts
WHERE active
ORDER BY rows DESC
LIMIT 50;
查看合并:
SELECT database, table, elapsed, num_parts, total_size
FROM system.merges
ORDER BY elapsed DESC;
常见问题:
| 现象 | 可能原因 | 处理 |
|---|---|---|
| Too many parts | 小写入太多、合并慢 | 批量写入、降并发、扩容 |
| Disk space insufficient | 磁盘满 | 清理 TTL、扩容、停止写入 |
| Merge with where part is lost | 副本缺 part | 从副本恢复 |
| Quota exceeded | 配额限制 | 调整用户配额 |
| Timeout | 写入批次过大或队列拥塞 | 拆批、重试、限流 |
本章小结
写入 ClickHouse 的关键是把小写入合并成大批量,并让失败可重试、数据可幂等。文件、HTTP、MySQL、S3 和 Kafka 都是输入通道,但都要关注格式、类型、时间分区和资源隔离。删除与更新应被视为特殊操作,而不是日常高频操作。
思考题
- 为什么小 INSERT 会引发 Too many parts?
- JSONEachRow 和 Parquet 分别适合什么场景?
- 应用写入端如何设计缓冲和重试?
- ReplacingMergeTree 如何辅助幂等?
- 大范围 UPDATE 为什么不适合 ClickHouse?