ClickHouseNotes

第 05 章:数据写入与导入

zjc 于 2026-01-05 发布

这是《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;

增量导入要解决:

  1. 时间基准;
  2. 时区;
  3. 更新时间回拨;
  4. 删除同步;
  5. 主键映射;
  6. 任务失败重放。

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,代价较高。分析表更常用的做法是:

  1. 修正上游后重写指定分区;
  2. 用替换表切换;
  3. 追加修正记录;
  4. 在查询层过滤无效标记。

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 都是输入通道,但都要关注格式、类型、时间分区和资源隔离。删除与更新应被视为特殊操作,而不是日常高频操作。

思考题

  1. 为什么小 INSERT 会引发 Too many parts?
  2. JSONEachRow 和 Parquet 分别适合什么场景?
  3. 应用写入端如何设计缓冲和重试?
  4. ReplacingMergeTree 如何辅助幂等?
  5. 大范围 UPDATE 为什么不适合 ClickHouse?