Learn
ClickHouse/10-data-ingestion

数据导入与导出

数据进不来,前面九章都是空谈。这一章覆盖从 CSV 文件到 Kafka 流的全部入库方式,以及导出。

1. INSERT 的黄金法则

第 4 章讲过 merge,这里把写入原则固化下来:

规则说明
每批 1 万~100 万行太小产生大量 part,太大占内存
每秒不超过 1~2 次 INSERT给后台 merge 留出时间
一次 INSERT 尽量只写少数分区跨分区会同时产生多个 part
攒不了批就开 async_insert让服务端帮你攒
-- 服务端攒批(应用端无法改造时的救命方案)
INSERT INTO events SETTINGS
    async_insert = 1,
    wait_for_async_insert = 1,           -- 等真正落盘再返回(推荐)
    async_insert_max_data_size = 10000000,   -- 攒够 10 MB 就刷
    async_insert_busy_timeout_ms = 1000      -- 或攒够 1 秒就刷
VALUES ('2024-07-01 10:00:00', 1, 'view', '/a', 'CN', 'web', 100);
wait_for_async_insert = 1(默认,安全)
  客户端 ──INSERT──▶ 服务端缓冲区 ──攒够刷盘──▶ part
        ◀──── 返回成功 ────────────────────────┘
  延迟高(最多等 1 秒),但返回成功就一定落盘了
 
wait_for_async_insert = 0(快但有风险)
  客户端 ──INSERT──▶ 服务端缓冲区
        ◀─立即返回─┘
  延迟低,但服务端崩溃时缓冲区数据丢失
⚠️Too many parts 是最常见的生产故障
Code: 252. DB::Exception: Too many parts (300).
Merges are processing significantly slower than inserts.

根因永远是「小批量高频写入」。应急处理:

-- 1. 先临时抬高阈值争取时间(治标)
ALTER TABLE events MODIFY SETTING parts_to_throw_insert = 3000;
 
-- 2. 手动触发合并
OPTIMIZE TABLE events PARTITION '202407';
 
-- 3. 检查是不是分区键太细导致的
SELECT count(DISTINCT partition) FROM system.parts WHERE table = 'events';

治本必须改写入方式:攒批或开 async_insert。

2. 数据格式

ClickHouse 支持 70 多种输入输出格式。常用的这几种要熟。

格式输入输出特点
CSV / CSVWithNames是是通用,WithNames 带表头
TSV / TSVWithNames是是制表符分隔,比 CSV 更不易出错
JSONEachRow是是每行一个 JSON 对象,流式友好
Parquet是是列存文件格式,压缩好,数据湖标配
Native是是ClickHouse 内部格式,最快
Values是是就是 INSERT VALUES 的括号语法
Pretty / PrettyCompact否是终端表格显示
Vertical否是竖排,宽表调试用

2.1 从 CSV 导入

准备文件 events.csv:

event_time,user_id,event_type,page,country,device,duration_ms
2024-07-01 10:00:00,1001,view,/home,CN,web,1200
2024-07-01 10:00:05,1002,click,/product/1,US,ios,340
2024-07-01 10:00:11,1001,purchase,/checkout,CN,web,5600
clickhouse-client --query "INSERT INTO demo.events FORMAT CSVWithNames" < events.csv

有表头就用 CSVWithNames(按列名匹配,列顺序可以和表不一致),没表头用 CSV(严格按位置匹配)。

常用的容错设置:

clickhouse-client --query "
INSERT INTO demo.events
SETTINGS
    input_format_allow_errors_num = 100,        -- 允许 100 行解析失败
    input_format_allow_errors_ratio = 0.01,     -- 或允许 1% 失败
    input_format_csv_skip_first_lines = 1,      -- 跳过首行
    date_time_input_format = 'best_effort'      -- 宽松解析时间格式
FORMAT CSV" < events.csv

date_time_input_format = 'best_effort' 特别有用,它能识别 ISO8601、RFC1123 等多种时间写法。

2.2 从 JSON 导入

{"event_time":"2024-07-01 10:00:00","user_id":1001,"event_type":"view","page":"/home","country":"CN","device":"web","duration_ms":1200}
{"event_time":"2024-07-01 10:00:05","user_id":1002,"event_type":"click","page":"/p/1","country":"US","device":"ios","duration_ms":340}
clickhouse-client --query "INSERT INTO demo.events FORMAT JSONEachRow" < events.jsonl

字段缺失时的处理:

SET input_format_skip_unknown_fields = 1;   -- 忽略表里没有的 JSON 字段
SET input_format_null_as_default = 1;       -- JSON 里的 null 用列默认值填充

2.3 Parquet

# 导入
clickhouse-client --query "INSERT INTO demo.events FORMAT Parquet" < events.parquet
 
# 导出
clickhouse-client --query "SELECT * FROM demo.events WHERE country='CN' FORMAT Parquet" > cn.parquet

Parquet 是数据湖交换的事实标准,和 Spark、Pandas、DuckDB 都能互通。

2.4 Native 格式:表间搬运最快

# 从一个集群导到另一个集群,比 CSV 快 5~10 倍
clickhouse-client --host src -q "SELECT * FROM demo.events FORMAT Native" \
  | clickhouse-client --host dst -q "INSERT INTO demo.events FORMAT Native"

Native 是二进制列存格式,不需要文本解析,也不丢精度。

💡压缩传输

大文件传输记得加压缩:

clickhouse-client -q "SELECT * FROM events FORMAT Native" | zstd > events.native.zst
zstd -d -c events.native.zst | clickhouse-client -q "INSERT INTO events FORMAT Native"

或者直接让 clickhouse-client 压缩:--compression=1。

3. 表函数:把外部数据当表查

表函数让你不建表就能读外部数据,是 ETL 的利器。

3.1 file

-- 直接查本地文件(文件要在 user_files 目录下)
SELECT country, count()
FROM file('events.csv', 'CSVWithNames',
          'event_time DateTime, user_id UInt64, event_type String,
           page String, country String, device String, duration_ms UInt32')
GROUP BY country;
 
-- 让 ClickHouse 自动推断 schema(24.x 支持得很好)
SELECT * FROM file('events.parquet') LIMIT 5;
 
-- 通配符批量读
SELECT count() FROM file('logs/2024-07-*.csv', CSVWithNames);

导入就是一句 INSERT ... SELECT:

INSERT INTO demo.events
SELECT * FROM file('events.csv', CSVWithNames);

文件必须放在 /var/lib/clickhouse/user_files/ 下(由 user_files_path 配置)。

3.2 url

SELECT * FROM url(
    'https://example.com/data/events.csv',
    'CSVWithNames',
    'event_time DateTime, user_id UInt64, country String'
) LIMIT 10;

3.3 s3

-- 读单个文件
SELECT count() FROM s3(
    'https://mybucket.s3.amazonaws.com/events/2024-07-01.parquet',
    'ACCESS_KEY', 'SECRET_KEY', 'Parquet'
);
 
-- 通配符批量读整个目录
INSERT INTO demo.events
SELECT * FROM s3(
    'https://mybucket.s3.amazonaws.com/events/2024-07-*.parquet',
    'ACCESS_KEY', 'SECRET_KEY', 'Parquet'
);
 
-- 导出到 S3
INSERT INTO FUNCTION s3(
    'https://mybucket.s3.amazonaws.com/export/cn_events.parquet',
    'ACCESS_KEY', 'SECRET_KEY', 'Parquet'
)
SELECT * FROM demo.events WHERE country = 'CN';

3.4 mysql / postgresql

直接从 MySQL 拉数据,这是 MySQL 到 ClickHouse 迁移最常用的方式:

-- 一次性查询
SELECT * FROM mysql(
    'mysql-host:3306', 'shop_db', 'orders',
    'ch_reader', 'password'
) LIMIT 10;
 
-- 全量导入
INSERT INTO demo.orders
SELECT
    order_id, user_id, order_time,
    toDecimal64(amount, 2), status, country, now()
FROM mysql('mysql-host:3306', 'shop_db', 'orders', 'ch_reader', 'password');
 
-- 增量导入(按时间切片,避免一次拉太多)
INSERT INTO demo.orders
SELECT order_id, user_id, order_time, toDecimal64(amount, 2), status, country, now()
FROM mysql('mysql-host:3306', 'shop_db', 'orders', 'ch_reader', 'password',
           SETTINGS connection_pool_size = 4)
WHERE order_time >= '2024-07-01' AND order_time < '2024-07-02';
⚠️mysql 表函数的 WHERE 下推

WHERE 条件会被下推到 MySQL 执行(对简单条件而言),但 GROUP BY、JOIN 不会。

如果写 SELECT country, count() FROM mysql(...) GROUP BY country,ClickHouse 会先把整张表拉过来再聚合。千万行以上的表要按时间/ID 分片拉取,否则会把 MySQL 拖垮。

也可以建持久的映射表:

CREATE TABLE mysql_orders
(
    order_id UInt64, user_id UInt64, order_time DateTime,
    amount Decimal(18,2), status String
)
ENGINE = MySQL('mysql-host:3306', 'shop_db', 'orders', 'ch_reader', 'password');
 
SELECT count() FROM mysql_orders;   -- 每次查询都实时打到 MySQL

3.5 numbers 与 generateRandom

造测试数据的两把刀:

-- numbers:生成 0 到 N-1
SELECT number FROM numbers(5);
 
-- generateRandom:按 schema 生成随机数据
SELECT * FROM generateRandom(
    'event_time DateTime, user_id UInt32, country String',
    1,      -- random_seed
    10,     -- max_string_length
    2       -- max_array_length
) LIMIT 3;

4. Kafka 引擎:流式接入

生产环境最常见的接入方式。架构是三件套:

Kafka Topic
     │
     ▼
┌─────────────────┐
│ Kafka 引擎表     │  ← 消费者,不存数据,读一次就消失
│ events_queue    │
└────────┬────────┘
         │ 物化视图触发
         ▼
┌─────────────────┐
│ MergeTree 表     │  ← 真正存数据
│ events          │
└─────────────────┘

4.1 完整配置

-- 1. Kafka 消费表
CREATE TABLE demo.events_queue
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  String,
    page        String,
    country     String,
    device      String,
    duration_ms UInt32
)
ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka1:9092,kafka2:9092',
    kafka_topic_list = 'web_events',
    kafka_group_name = 'clickhouse_consumer',
    kafka_format = 'JSONEachRow',
    kafka_num_consumers = 4,                     -- 并发消费者数,不超过分区数
    kafka_max_block_size = 100000,               -- 攒够 10 万行才写一批
    kafka_poll_max_batch_size = 10000,
    kafka_flush_interval_ms = 5000,              -- 或攒够 5 秒
    kafka_skip_broken_messages = 10,             -- 容忍 10 条坏消息
    kafka_handle_error_mode = 'stream';          -- 坏消息进虚拟列而不是丢弃
 
-- 2. 目标存储表
CREATE TABLE demo.events
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    page        String,
    country     LowCardinality(String),
    device      LowCardinality(String),
    duration_ms UInt32
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (country, event_type, event_time);
 
-- 3. 搬运用的物化视图
CREATE MATERIALIZED VIEW demo.events_consumer TO demo.events AS
SELECT
    event_time,
    user_id,
    event_type,
    page,
    upper(country) AS country,      -- 可以顺手做清洗
    device,
    duration_ms
FROM demo.events_queue;

4.2 关键设置解释

设置作用
kafka_num_consumers并发消费线程数,设为 Kafka 分区数(不要超过)
kafka_max_block_size攒够这么多行才触发一次写入,直接决定 part 大小
kafka_flush_interval_ms攒不够行数时的超时刷新
kafka_skip_broken_messages解析失败时跳过的消息数,防止一条脏数据卡死消费
kafka_handle_error_mode设为 stream 时坏消息可通过虚拟列捕获
⚠️Kafka 引擎表不能直接 SELECT
SELECT * FROM events_queue LIMIT 10;   -- 会真的消费掉这些消息!

Kafka 表是一次性读取的流,SELECT 会推进 offset,这些消息就不会再进目标表了。调试时绝对不要直接查它。

4.3 虚拟列与错误处理

CREATE MATERIALIZED VIEW demo.events_consumer TO demo.events AS
SELECT
    event_time, user_id, event_type, page, country, device, duration_ms,
    _topic, _partition, _offset, _timestamp     -- 虚拟列,可用于排查
FROM demo.events_queue;
 
-- 捕获坏消息到单独的表
CREATE TABLE demo.events_errors
(
    topic String, partition UInt64, offset UInt64,
    raw String, error String
) ENGINE = MergeTree ORDER BY (topic, partition, offset);
 
CREATE MATERIALIZED VIEW demo.events_error_mv TO demo.events_errors AS
SELECT _topic AS topic, _partition AS partition, _offset AS offset,
       _raw_message AS raw, _error AS error
FROM demo.events_queue
WHERE length(_error) > 0;

4.4 监控消费

-- 消费状态
SELECT * FROM system.kafka_consumers FORMAT Vertical;
Row 1:
──────
database:            demo
table:               events_queue
consumer_id:         ClickHouse-node1-demo-events_queue-0
assignments.topic:   ['web_events']
assignments.partition_id: [0,1,2,3]
assignments.current_offset: [1284012, 1283998, 1284105, 1283877]
num_messages_read:   5136992
last_poll_time:      2024-07-30 12:03:41
num_rebalance_revocations: 0
exceptions.text:     []
-- 暂停/恢复消费(运维必备)
DETACH TABLE demo.events_consumer;    -- 停止搬运
ATTACH TABLE demo.events_consumer;    -- 恢复

5. 导出

# CSV
clickhouse-client -q "SELECT * FROM demo.events WHERE country='CN'" \
  --format CSVWithNames > cn_events.csv
 
# JSON
clickhouse-client -q "SELECT * FROM demo.events LIMIT 100" \
  --format JSONEachRow > sample.jsonl
 
# Parquet(推荐给数据科学同事)
clickhouse-client -q "SELECT * FROM demo.events WHERE event_time >= '2024-07-01'" \
  --format Parquet > july.parquet
 
# 压缩导出
clickhouse-client -q "SELECT * FROM demo.events" --format Native \
  | zstd -3 > events.native.zst

也可以用 INTO OUTFILE(客户端侧写文件):

SELECT * FROM demo.events WHERE country = 'CN'
INTO OUTFILE 'cn_events.parquet'
FORMAT Parquet;
 
-- 自动按扩展名推断格式和压缩
SELECT * FROM demo.events INTO OUTFILE 'events.csv.gz';

6. 大规模导入的实战建议

-- 导入时临时调大的设置
INSERT INTO demo.events
SELECT * FROM s3('...*.parquet', 'Parquet')
SETTINGS
    max_insert_threads = 8,                   -- 并行写入线程
    max_insert_block_size = 1048576,          -- 每个 block 的行数
    min_insert_block_size_rows = 1048576,     -- 攒够 100 万行才形成 part
    min_insert_block_size_bytes = 268435456,  -- 或攒够 256 MB
    max_memory_usage = 20000000000;

大批量导入的完整流程:

1. 建表时先不加 TTL、不建物化视图(减少写入开销)
2. 按分区分批导入,每批一个分区
   INSERT ... WHERE toYYYYMM(t) = 202406
3. 每批导完检查 part 数量,必要时 OPTIMIZE
4. 全部导完后再建 MV 并回填
5. 最后加上 TTL
💡导入前先小批量验证
-- 先只导 1000 行看看类型对不对
INSERT INTO demo.events
SELECT * FROM s3('...', 'Parquet') LIMIT 1000;
 
-- 检查关键列有没有变成默认值
SELECT
    count(),
    countIf(country = '') AS empty_country,
    countIf(duration_ms = 0) AS zero_duration,
    min(event_time), max(event_time)
FROM demo.events;

发现问题只需 TRUNCATE TABLE 重来,比导了 10 亿行才发现时间字段全是 1970 好得多。

🎯练习
  1. 用 generateRandom 生成 10 万行数据导出成 CSV,再用 file() 表函数读回来,验证往返一致。
  2. 写一个循环脚本,每次 INSERT 100 行、连续 500 次,观察 system.parts 的数量增长。然后加上 async_insert = 1 重跑,对比最终 part 数。
  3. 如果你手边有 MySQL,用 mysql() 表函数读一张表并导入 ClickHouse,用 EXPLAIN 或 MySQL 的 general_log 观察 WHERE 是否被下推。
  4. 用 Docker 起一个单节点 Kafka,按第 4 节配置 Kafka 引擎表 + MV,向 topic 灌 JSON 消息,验证数据自动进入 events 表。然后 DETACH MV,继续灌消息,再 ATTACH,验证 offset 没丢。

小结

  • 写入原则:每批 1 万100 万行、每秒 12 次,做不到就开 async_insert
  • Too many parts 的根因永远是小批量高频写入
  • 格式选择:文本交换用 CSV/JSONEachRow,数据湖用 Parquet,集群间搬运用 Native
  • 表函数 file / url / s3 / mysql 让你不建表就能读外部数据,配合 INSERT ... SELECT 就是完整 ETL
  • mysql() 表函数只下推 WHERE,不下推聚合,大表必须分片拉取
  • Kafka 接入是三件套:Kafka 引擎表 + MergeTree 存储表 + 物化视图搬运
  • Kafka 引擎表绝不能直接 SELECT,会消费掉消息
  • 大规模导入先小批量验证 schema,导完再建 MV 和 TTL
  • 下一章开始进入查询侧,先讲 SELECT 语法与常用函数 →