数据导入与导出
数据进不来,前面九章都是空谈。这一章覆盖从 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──▶ 服务端缓冲区
◀─立即返回─┘
延迟低,但服务端崩溃时缓冲区数据丢失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,5600clickhouse-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.csvdate_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.parquetParquet 是数据湖交换的事实标准,和 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';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; -- 每次查询都实时打到 MySQL3.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 时坏消息可通过虚拟列捕获 |
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 好得多。
- 用
generateRandom生成 10 万行数据导出成 CSV,再用file()表函数读回来,验证往返一致。 - 写一个循环脚本,每次 INSERT 100 行、连续 500 次,观察
system.parts的数量增长。然后加上async_insert = 1重跑,对比最终 part 数。 - 如果你手边有 MySQL,用
mysql()表函数读一张表并导入 ClickHouse,用EXPLAIN或 MySQL 的 general_log 观察 WHERE 是否被下推。 - 用 Docker 起一个单节点 Kafka,按第 4 节配置 Kafka 引擎表 + MV,向 topic 灌 JSON 消息,验证数据自动进入
events表。然后DETACHMV,继续灌消息,再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 语法与常用函数 →