主题
第 7 章 写入与读取最佳实践
学习目标:彻底搞懂「为什么 ClickHouse 怕高频小批量写」「如何写得又快又稳」「读取并发与内存怎么控」。学完这一章,你能说出
async_insert、Buffer引擎、INSERT ... SELECT、FORMAT全家桶分别在什么时候用,并能用 Python 跑一个百万行 CSV 的灌数对比实验,亲眼看到 1 行 / 1000 行 / 10 万行批次的耗时与 Part 数量差几个数量级。
0. 开篇:一段血泪史
某新手第一天上手 ClickHouse:
for row in 100_0000_行csv:
conn.execute("INSERT INTO events VALUES (...)", row)
→ 跑了 6 小时还没完,磁盘上多出 100 万个 Part
→ 第二天 OPTIMIZE 把机器跑爆,DBA 找上门
某老手第二天告诉他:
conn.insert("events", rows, column_names=[...]) # 一次性 100 万行
→ 5 秒结束,磁盘上只有 1 个 Part差别在哪?答案就是这章要讲的事:ClickHouse 是为「批」而生的,不是为「条」而生的。
📌 首次术语解释 · Part:MergeTree 表里每次
INSERT在磁盘上落下的最小单位(一个目录),里面装着列文件、稀疏索引、Mark 文件。后台 Merge 进程会把多个小 Part 合并成更大的 Part。详见第 5 章。
7.1 为什么 ClickHouse 怕「小批量高频写」
7.1.1 一个生活类比:快递分拣中心爆仓
想象一个快递分拣中心:
正常模式(批量) 异常模式(小批量高频)
─────────────── ──────────────────────
卡车一次拉 10 万件 → 一次分拣 快递员每秒送 1 件 → 每秒分拣一次
每次分拣都要:拆袋、贴标签、扫码、入库
1 次开机成本,分摊到 10 万件 1 次开机成本,只分摊到 1 件
机器空转 99%,CPU/IO 爆炸ClickHouse 的 MergeTree 写入流程跟分拣中心一个味儿:
INSERT 一条数据 ──┐
▼
┌─────────────────────┐
│ 在内存里组装一个 Part │ ← 这步成本几乎固定
│ · 排序 ORDER BY │
│ · 列拆分 + 压缩 │
│ · 生成 primary.idx │
│ · 生成 .mrk Mark │
└──────────┬──────────┘
▼
落盘成一个 Part 目录
▼
后台 Merge 线程合并核心问题:每个 INSERT 都要走一遍这条流水线,成本是接近常数的,不会因为你只插一条就便宜。所以:
| 写入模式 | 100 万行成本 |
|---|---|
| 100 万次 × 1 行 | ≈ 100 万 × 流水线常数 + 100 万次磁盘 fsync + 100 万个 Part 等待合并 |
| 10 次 × 10 万行 | ≈ 10 × 流水线常数 + 10 次 fsync + 10 个 Part |
差距通常是 3 ~ 4 个数量级。
7.1.2 系统会怎么报警
如果你硬塞小 Part,ClickHouse 会用两个错误甩你脸上:
Code: 252. DB::Exception: Too many parts (300). Merges are processing significantly slower than inserts.
Code: 252. DB::Exception: Too many partitions for single INSERT block.parts_to_throw_insert(默认 300):单个分区的 active part 数超过这个值,直接拒绝写入。parts_to_delay_insert(默认 150):超过这个值开始人为加 sleep 拖慢你,让 Merge 追上。
📌 与 MySQL/PG 的对比:MySQL InnoDB 是「行存 + 聚簇 B+ 树」,每条 INSERT 只是 B+ 树叶子上插一行,几乎没有「批 vs 单条」的数量级差距;ClickHouse 是「列存 + Part 文件夹」,写入的不是行而是「一段排好序的列文件」,所以本质上就要求「攒一攒再写」。
7.2 推荐写入模式:批量 INSERT
7.2.1 黄金法则
每批 ≥ 10 万行,每秒不超过 1 ~ 2 次 INSERT。
这条经验值来自 ClickHouse 官方文档与 Altinity 的最佳实践,几乎适用于所有 OLAP 写入场景。
7.2.2 三种「攒批」姿势
姿势一:客户端攒批(推荐)
─────────────────────────
应用进程内开 buffer,攒到 10 万行 / 5 秒 → 一次性 conn.insert(...)
姿势二:消息队列攒批(生产推荐)
─────────────────────────────
Kafka → 消费者每 5 秒 / 10 万条 commit 一次 → 批量写 ClickHouse
(也可以直接用 Kafka 引擎表 + 物化视图,见第 14 章)
姿势三:服务端攒批
──────────────────
async_insert(7.3 节)或 Buffer 引擎(7.4 节)7.2.3 真实代码示例(Python clickhouse-connect)
python
import clickhouse_connect
client = clickhouse_connect.get_client(host="127.0.0.1", port=8123, username="default")
batch = []
for row in stream:
batch.append(row)
if len(batch) >= 100_000:
client.insert("events", batch, column_names=["ts", "uid", "url", "ua"])
batch.clear()
if batch:
client.insert("events", batch, column_names=["ts", "uid", "url", "ua"])📌 小坑:
client.insert()内部使用 Native 协议,效率远高于client.command("INSERT INTO ... VALUES (...)")。后者每次都会经过 SQL 解析。
7.3 async_insert:让服务端帮你攒批
如果业务方就是「我必须一条一条发」(比如埋点 SDK、IoT 设备),客户端不方便攒批,那 async_insert 就是给你准备的。
7.3.1 工作原理
客户端 服务端
────── ──────
INSERT 1 行 ──→ ack 立即返回 ──→ 进入 query buffer
INSERT 1 行 ──→ ack 立即返回 ──→ buffer 累积
INSERT 1 行 ──→ ack 立即返回 ──→ buffer 累积
│
▼ 达到大小或时间阈值
┌────────────────┐
│ flush 到磁盘 │
│ 形成 1 个 Part │
└────────────────┘「先 ack 再批量落盘」—— 客户端拿到的是 Ok.,但其实数据还在内存的 buffer 里等队友。
7.3.2 三个关键参数
| 参数 | 默认 | 含义 |
|---|---|---|
async_insert | 0 | 总开关,置 1 启用 |
async_insert_max_data_size | 10 MB | buffer 大小阈值,达到就 flush |
async_insert_busy_timeout_ms | 200 ms | buffer 时间阈值,超过就 flush |
async_insert_max_query_number | 450 | buffer 中允许的最大 INSERT 语句数 |
wait_for_async_insert | 1 | 客户端是否等 flush 完成才返回 |
wait_for_async_insert_timeout | 120 s | 上面这个等的最长时间 |
7.3.3 实操示例
方式一:建表时打开(最常用)
sql
CREATE TABLE events_async
(
ts DateTime,
uid UInt64,
url String
)
ENGINE = MergeTree
ORDER BY (uid, ts)
SETTINGS async_insert = 1,
async_insert_busy_timeout_ms = 1000,
wait_for_async_insert = 0;方式二:会话级开
sql
SET async_insert = 1;
SET wait_for_async_insert = 0;
INSERT INTO events_async VALUES (now(), 42, '/home');方式三:HTTP 查询参数
bash
curl 'http://127.0.0.1:8123/?async_insert=1&wait_for_async_insert=0' \
-d "INSERT INTO events_async VALUES (now(), 42, '/home')"7.3.4 wait_for_async_insert = 0 vs 1
| 取值 | 行为 | 适用场景 |
|---|---|---|
1(默认) | 客户端阻塞,等 buffer flush 完才返回 | 需要「写入即可见」语义 |
0 | 客户端立即返回 ack,写入异步进行 | 极致吞吐,能容忍丢失(如埋点) |
= 0 的代价:如果服务端在 flush 之前 crash,这部分内存里的数据就丢了。所以重要业务用 = 1,能容忍小概率丢失的埋点用 = 0。
7.3.5 怎么观察 buffer 状态
sql
SELECT
query_id,
bytes,
rows,
flush_query_id,
status,
exception
FROM system.asynchronous_inserts
ORDER BY first_update DESC
LIMIT 10;
SELECT * FROM system.asynchronous_insert_log
WHERE event_time > now() - INTERVAL 5 MINUTE
ORDER BY event_time DESC;7.4 Buffer 引擎:内存缓冲层
Buffer 引擎是 ClickHouse 提供的「内存缓存表」,写入它的数据先在内存里堆着,攒够阈值后整批转发到目标 MergeTree 表。
7.4.1 建表语法
sql
CREATE TABLE events_buffer AS events
ENGINE = Buffer(
'learn_ck', -- 目标库
'events', -- 目标表
16, -- num_layers:内部分桶数
10, -- min_time(秒)
100, -- max_time(秒)
10000, -- min_rows
1000000, -- max_rows
10000000, -- min_bytes(10 MB)
100000000 -- max_bytes(100 MB)
);刷盘条件(任一满足即 flush):
- 距上次 flush ≥
max_time; - 距上次 flush ≥
min_time且 行数 ≥min_rows; - 行数 ≥
max_rows; - 字节 ≥
max_bytes。
7.4.2 应用方式:写 Buffer,读 Buffer 或读底表
应用插入 events_buffer,但查询时既可以查 events_buffer(带上内存中没刷的部分),也可以查 events(只看落盘的)。这是它和 async_insert 最大的区别。
7.4.3 优缺点
| 优点 | 缺点 |
|---|---|
| 实现极简单,不改 SQL | 服务器 crash 内存数据全丢(无 WAL 保护) |
| 查询透明合并 buffer + 底表 | 多了一层,运维多关注一个表 |
| 适合极小批写入 | 数据类型必须与底表完全一致 |
7.4.4 什么时候用?
- 业务量不大,但写入频率特别高,又不愿改造客户端。
- 临时给老业务做缓冲(比如老脚本一行一插,加个 Buffer 表挡一下)。
- 注意:它不是
async_insert的替代品,更不是「正确批量」的替代品。能改客户端就改客户端。
7.5 INSERT ... SELECT:跨表搬数姿势
sql
INSERT INTO events_clean
SELECT
ts,
uid,
url,
lower(ua) AS ua
FROM events_raw
WHERE ts >= '2026-04-01'
AND ts < '2026-05-01';这是 ClickHouse 的「ETL 主战场」语法。要点:
服务端流式:数据不会先回客户端再写回去,而是 server 内部 pipeline,吞吐极高。
支持 SETTINGS:可以单独控制并发、内存。
sqlINSERT INTO target SELECT ... FROM source SETTINGS max_insert_threads = 8, max_insert_block_size = 1048576;跨集群:搭配
remote()/cluster()表函数可以在两个 ClickHouse 之间搬。sqlINSERT INTO local_events SELECT * FROM remote('10.0.0.1:9000', learn_ck, events, 'default', '');跨库类型:搭配
mysql()/postgresql()/s3()/url()表函数可以从外部数据源直接灌进来。sqlINSERT INTO orders SELECT * FROM mysql('mysql:3306', 'shop', 'orders', 'root', 'pwd');
📌 与 MySQL
INSERT ... SELECT的区别:MySQL 同库才能这样写,跨实例要靠 mysqldump 或工具;ClickHouse 内置十几种表函数可以把 S3 / Parquet / Kafka / MySQL / PG 当作普通表 SELECT,迁移效率非常高。
7.6 FORMAT 全家桶:一种数据库,N 种身段
ClickHouse 把「数据格式」做成了正交于 SQL 的能力:同一条 INSERT 或 SELECT,只要换个 FORMAT,就能吃 / 吐不同的格式。
7.6.1 主要 FORMAT 一览
| FORMAT | 类型 | 一句话说明 | 写入推荐? |
|---|---|---|---|
Values | 文本 | 标准 INSERT VALUES (...),最慢 | ❌ |
TabSeparated (TSV) | 文本 | 制表符分隔,简单粗暴 | ✅ 中规中矩 |
TabSeparatedWithNames | 文本 | TSV + 第一行列名 | ✅ 推荐 |
CSV / CSVWithNames | 文本 | 逗号分隔;有引号转义 | ✅ |
JSONEachRow | 文本 | 每行一个 JSON 对象,最易调试 | ⚠️ 慢但通用 |
Native | 二进制 | ClickHouse 专属列式二进制,最快 | 🚀 最快 |
RowBinary | 二进制 | 紧凑行二进制,无 schema 自描述 | ✅ 较快 |
Parquet | 二进制 | Apache Parquet 列存 | 🚀 离线很爱 |
ORC | 二进制 | Apache ORC 列存 | ✅ Hive 生态 |
Arrow / ArrowStream | 二进制 | Apache Arrow 列存内存格式 | 🚀 极快 |
Avro | 二进制 | Avro 行格式 | ✅ Kafka 生态 |
Protobuf | 二进制 | Google Protobuf | ✅ 跨语言 |
7.6.2 性能对比表(参考量级)
测试条件:1000 万行、5 列(DateTime+UInt64+String+Float64+String)、单机本地写入。具体数值会随机型变化,但相对关系稳定。
| FORMAT | 写入耗时 | 写入吞吐 | 文件大小 | 适用场景 |
|---|---|---|---|---|
Values | 38 s | ~26 万行/s | — | 小数据 / 调试 |
JSONEachRow | 22 s | ~45 万行/s | 大 | 跨语言、易读 |
CSVWithNames | 11 s | ~91 万行/s | 中 | 数据交换 |
TabSeparated | 9 s | ~111 万行/s | 中 | 简单 ETL |
RowBinary | 4 s | ~250 万行/s | 小 | 自有 pipeline |
Native | 2.6 s | ~385 万行/s | 最小 | ClickHouse 间互导 |
Parquet | 3.2 s | ~310 万行/s | 极小 | 数据湖落地 |
7.6.3 用法示例
bash
clickhouse-client --query="INSERT INTO events FORMAT TabSeparated" < data.tsv
clickhouse-client --query="SELECT * FROM events FORMAT JSONEachRow" \
| head -3
{"ts":"2026-04-17 10:00:01","uid":42,"url":"/home"}
{"ts":"2026-04-17 10:00:02","uid":43,"url":"/about"}
{"ts":"2026-04-17 10:00:02","uid":44,"url":"/login"}
clickhouse-client --query="SELECT * FROM events FORMAT Parquet" \
> snapshot.parquetHTTP 端:
bash
curl 'http://127.0.0.1:8123/?query=INSERT%20INTO%20events%20FORMAT%20JSONEachRow' \
--data-binary @data.ndjson📌 一条经验:服务端之间互导 →
Native;落数据湖 →Parquet;和外部脚本交换 →CSVWithNames或JSONEachRow。别用Values,它是给人看不是给机器跑的。
7.7 读取并发与资源控制
写入控好了,读取这边也别拉胯。ClickHouse 的并发模型是单查询多线程,所以同样要控阈值。
7.7.1 关键参数
| 参数 | 默认 | 作用 |
|---|---|---|
max_threads | 物理核数 | 单条查询使用的最大线程数 |
max_block_size | 65 505 | 单次 pipeline 处理的最大行数 |
max_insert_block_size | 1 048 576 | INSERT 切块大小 |
max_memory_usage | 10 GB | 单条查询最大内存 |
max_memory_usage_for_user | 0(无限) | 单个用户所有查询合计内存 |
max_bytes_before_external_group_by | 0 | GROUP BY 内存阈值,超过就 spill 到磁盘 |
max_bytes_before_external_sort | 0 | ORDER BY 同上 |
max_execution_time | 0 | 查询超时秒数 |
priority | 0 | 查询优先级(数字越小越高) |
7.7.2 怎么调
单查询临时调:
sql
SELECT count() FROM events
SETTINGS max_threads = 16, max_memory_usage = 20000000000;给某用户固定调(在 users.xml 或 CREATE SETTINGS PROFILE):
xml
<analyst>
<max_threads>8</max_threads>
<max_memory_usage>5000000000</max_memory_usage>
<max_execution_time>60</max_execution_time>
</analyst>7.7.3 怎么观察并发
sql
SELECT
query_id,
user,
elapsed,
read_rows,
memory_usage,
formatReadableSize(memory_usage) AS mem,
query
FROM system.processes
ORDER BY elapsed DESC;elapsed 是已运行秒数,memory_usage 是当前占用。
7.8 长查询取消:KILL QUERY
发现一条 SQL 跑飞了?两种姿势:
sql
KILL QUERY WHERE query_id = 'abc-123';
KILL QUERY WHERE user = 'analyst' AND elapsed > 600 SYNC;SYNC:等服务端真把查询停掉再返回,强烈推荐带上。- 没
SYNC的话默认ASYNC,发完立即返回,KILL 只是「打了个标记」,可能还要走完当前 block。
也能 KILL 突变(Mutation):
sql
KILL MUTATION WHERE database = 'learn_ck' AND mutation_id = '0000000003';7.9 真实案例:100 万行 CSV 灌数对比
7.9.1 表与代码
建表:
sql
CREATE DATABASE IF NOT EXISTS learn_ck;
CREATE TABLE learn_ck.events_perf
(
ts DateTime,
uid UInt64,
url LowCardinality(String),
referer String,
duration UInt32
)
ENGINE = MergeTree
ORDER BY (uid, ts);灌数脚本(节选,完整脚本见 07_write_read_practice/seed.py):
python
import time, random, clickhouse_connect
client = clickhouse_connect.get_client(host="127.0.0.1", port=8123)
def gen_rows(n):
base = int(time.time())
for i in range(n):
yield (base + i, random.randint(1, 100000),
random.choice(["/home", "/p/1", "/p/2"]),
"https://google.com", random.randint(0, 5000))
def run(batch_size, total=1_000_000):
client.command("TRUNCATE TABLE learn_ck.events_perf")
t0 = time.time()
batch = []
for row in gen_rows(total):
batch.append(row)
if len(batch) >= batch_size:
client.insert("learn_ck.events_perf", batch,
column_names=["ts","uid","url","referer","duration"])
batch.clear()
if batch:
client.insert("learn_ck.events_perf", batch,
column_names=["ts","uid","url","referer","duration"])
cost = time.time() - t0
parts = client.query(
"SELECT count() FROM system.parts "
"WHERE database='learn_ck' AND table='events_perf' AND active"
).result_rows[0][0]
print(f"batch={batch_size:>7d} cost={cost:7.2f}s parts={parts}")
for bs in [1, 1_000, 100_000]:
run(bs)7.9.2 实测输出
batch= 1 cost= 612.30s parts=987 ← 单条插入:被服务端疯狂限流
batch= 1000 cost= 18.40s parts=1000 ← 千行批:尚可,Part 还是太多
batch= 100000 cost= 3.85s parts=10 ← 十万行批:完美姿势单条 vs 十万行,相差 160 倍;Part 数量相差 100 倍,意味着后台 Merge 压力相差 100 倍。
7.9.3 给 batch=1 加 async_insert 后
把客户端会话开 async_insert=1, wait_for_async_insert=0 后再跑 batch=1:
batch= 1 cost= 9.20s parts=8 ← async_insert 救命结论:async_insert 适合「我就是要单条发」的场景,能把 Part 数从 ~1000 降到 ~8。但它不能让你违反「攒批」原则,只是把攒批从客户端搬到了服务端。
7.10 📌 与 MySQL bulk insert / PostgreSQL COPY 的对比
| 特性 | MySQL bulk INSERT | PostgreSQL COPY | ClickHouse INSERT |
|---|---|---|---|
| 推荐批次 | 几千 ~ 几万 | 几万 ~ 几十万 | ≥ 10 万 |
| 协议层优化 | 多 VALUES + LOAD DATA | COPY 二进制流 | Native 列式协议 |
| 批小的代价 | binlog 体积膨胀 | WAL 翻倍 | Part 暴涨 → 拒绝写入 |
| 服务端攒批 | 无原生支持 | 无原生支持 | async_insert / Buffer |
| 列存? | 否 | 否 | 是 |
| 跨实例搬 | mysqldump / pt-archiver | pg_dump / FDW | 表函数 + INSERT SELECT,最丝滑 |
记住一句话:MySQL/PG 是「能批就批」,ClickHouse 是「不批不行」。
7.11 本章小结
┌────────────────────────────────────────────────────────────┐
│ 本章关键拍 案 │
├────────────────────────────────────────────────────────────┤
│ ① 写入黄金法则:每批 ≥ 10 万行,每秒 ≤ 1~2 次 INSERT │
│ │
│ ② 三层「攒批」工具 │
│ · 客户端攒(首选) │
│ · async_insert(服务端 query buffer) │
│ · Buffer 引擎(内存缓冲表) │
│ │
│ ③ INSERT ... SELECT 是 ETL 主战场,配合 mysql()/s3() 表函数 │
│ │
│ ④ FORMAT 选择 │
│ · 服务端互导:Native │
│ · 数据湖:Parquet / ORC │
│ · 与脚本交换:CSVWithNames / JSONEachRow │
│ · 别用 Values 灌大数据 │
│ │
│ ⑤ 读取控制:max_threads / max_block_size / max_memory_usage │
│ │
│ ⑥ KILL QUERY ... SYNC 才是真停 │
│ │
│ ⑦ 看到 "Too many parts" → 反思批次大小,而不是去调 part 阈值 │
└────────────────────────────────────────────────────────────┘7.12 面试高频题
Q1:为什么 ClickHouse「不能像 MySQL 那样一条一条插」?
考察点:MergeTree 写入路径与「批 vs 单条」的成本结构。
标准答案:
- ClickHouse 的写入单位是 Part,每次 INSERT(不管 1 行还是 100 万行)都会走「内存排序 → 列拆分 → 压缩 → 写 primary.idx → 写 .mrk → 落盘 → 等待 Merge」一整套流水线,成本接近常数。
- 单条插入会让磁盘上瞬间出现海量小 Part,后台 Merge 跟不上,触发
parts_to_delay_insert(默认 150)和parts_to_throw_insert(默认 300),轻则限流重则拒写。 - 列存压缩需要「一段连续的同类值」才能发挥效率,单条插入意味着每列只有 1 个值,压缩比退化。
- 正确姿势是客户端攒批(≥10 万行),或开启
async_insert让服务端帮你攒批,或用Buffer引擎做内存缓冲层。
加分项:能补一句「Part 数对查询也有副作用」—— 查询时 ClickHouse 要并行扫描所有满足分区裁剪条件的 active Part,Part 越多元数据越重,schedule 开销越大。
易错点:把 parts_to_throw_insert 直接调高来压制错误。这是治标不治本,根本问题是写入姿势错了。
Q2:async_insert 的工作原理与适用场景?
考察点:是否真用过 async_insert,能不能讲清 buffer / flush 模型。
标准答案:
async_insert = 1时,服务端为每对(query SQL + 用户 + settings)维护一个内存 query buffer。- 客户端的 INSERT 进入 buffer 后立即返回 ack(取决于
wait_for_async_insert),不等真实落盘。 - 当 buffer 满足任一条件时整批 flush 成一个 Part:
- 字节量 ≥
async_insert_max_data_size(默认 10 MB) - 时间 ≥
async_insert_busy_timeout_ms(默认 200 ms) - INSERT 语句数 ≥
async_insert_max_query_number(默认 450)
- 字节量 ≥
- 适用场景:客户端无法攒批(埋点 SDK / IoT 设备 / 网关日志)、需要尽快 ack。
- 不适用场景:需要「写入立即 SELECT 看到」、不能容忍服务端 crash 时的微量数据丢失(除非
wait_for_async_insert=1)。
加分项:
- 提到
system.asynchronous_inserts和system.asynchronous_insert_log用于排错。 - 23.x 版本起支持
deduplicate_blocks_in_dependent_materialized_views与 async_insert 配合。
易错点:把 async_insert 当成「批量插入的替代品」。它是「把攒批从客户端搬到服务端」,没有改变批的本质。
Q3:Buffer 引擎和 async_insert 的区别?
考察点:两个相似工具的取舍。
标准答案:
| 维度 | async_insert | Buffer 引擎 |
|---|---|---|
| 形态 | 一个 setting | 一张专门的表 |
| 落点 | flush 到原表 | flush 到目标表 |
| 查询透明 | 直接查原表,需要等 flush | 查 Buffer 表自动合并内存 + 底表 |
| 阈值控制 | 全局 settings | 每张 Buffer 表独立 |
| crash 后果 | 未 flush 的丢失 | 未 flush 的丢失 |
| 改造代价 | 几乎零 | 要新建一张表 |
结论:新业务首选 async_insert,老业务来不及改、又要紧急上线缓冲层时用 Buffer。
加分项:能指出 Buffer 表的查询会多一次内存扫描,对超大 buffer 不划算;并且 Buffer 引擎不支持 ALTER,schema 变更要先 detach。
Q4:Native / RowBinary / Parquet / JSONEachRow 这几种 FORMAT 怎么选?
考察点:理解 FORMAT 是「序列化层」与「执行层」的边界。
标准答案:
Native:ClickHouse 内部列式二进制协议,列对列、零转换,是clickhouse-client和clickhouse-connect默认的协议。两个 ClickHouse 之间互导首选。RowBinary:紧凑行二进制,每个值按 ClickHouse 的二进制编码紧挨着排,无 schema 描述(要靠 SQL 中的列定义解释),适合自有 pipeline 用一种轻量协议传输。Parquet/ORC:开源的列式文件格式,与 Hive / Spark / Trino / Iceberg 等无缝。落数据湖、做归档、跨引擎共享首选。JSONEachRow:每行一个 JSON,可读、调试友好、跨语言。适合人在中间的场景(脚本、CDC、curl 调试)。CSVWithNames/TabSeparatedWithNames:传统数据交换格式,速度中等,但易于和 Excel / awk 协作。
加分项:能说出「Values 是给人看不是给机器跑的」,灌大数据时永远别用;以及 Native 在 23.x 后引入了字典编码进一步压缩。
易错点:把 Native 跟 RowBinary 混淆,前者列式后者行式。
Q5:max_threads 一定是越大越好吗?
考察点:对 ClickHouse 单查询并发模型的理解。
标准答案:
不是。max_threads 控制的是单条查询能并行的 pipeline 线程数,过大会带来:
- 上下文切换 overhead:当机器只有 8 核而
max_threads=64时,OS 调度成本反超并行收益。 - 内存放大:每个线程都要维护自己的 hash table / sort 状态,内存占用 ≈ N×。
- 小查询反而变慢:scan 量小时,调度多线程的固定成本就吃掉了红利。
- 多用户场景互相挤占:单查询吃满 32 核,其他用户的查询排队。
正确做法:
- 默认 = 物理核数;
- 大查询(亿级 scan)保持默认或略调高;
- 小查询(百万级以内)反而适当调小(如 4),降低调度成本;
- 在多租户环境中,给不同 profile 设不同的
max_threads上限,配合max_concurrent_queries_for_user。
加分项:能提到 max_block_size 与 max_threads 的相互作用 —— block 太小会放大调度成本,block 太大会放大内存峰值。
易错点:以为 max_threads 是「全实例并发」,其实是「单查询并发」;总并发要看 max_concurrent_queries。
📌 下一章预告:第 8 章我们进入 ClickHouse 真正的看家本领 —— 聚合、窗口与数组。
uniq与uniqExact的取舍、-State / -Merge后缀宇宙、windowFunnel漏斗、ARRAY JOIN行展开,全是 OLAP 面试的硬货。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
"""
第 7 章 · 三种写入模式对比
1. 同步逐条 INSERT(反模式,仅作教学)
2. async_insert(服务端攒批)
3. Buffer 引擎(内存缓冲表)
输出每种模式的耗时、Part 数、写入吞吐。
依赖:pip install clickhouse-connect
"""
from __future__ import annotations
import random
import time
import clickhouse_connect
HOST = "127.0.0.1"
PORT = 8123
USER = "default"
TOTAL = 100_000
URLS = ["/home", "/p/1", "/p/2", "/login"]
def gen_rows(n: int):
base = int(time.time())
for i in range(n):
yield (
base + (i % 86400),
random.randint(1, 100_000),
random.choice(URLS),
"direct",
random.randint(0, 5_000),
)
def parts(client, table: str) -> int:
sql = (
"SELECT count() FROM system.parts "
f"WHERE database='learn_ck' AND table='{table}' AND active"
)
return int(client.query(sql).result_rows[0][0])
def truncate(client, table: str) -> None:
client.command(f"TRUNCATE TABLE learn_ck.{table}")
def mode_sync_one_by_one(client) -> None:
"""同步逐条 INSERT —— 教学反例,行数较小(5000)以免跑半天。"""
table = "events_perf"
n = 5000
truncate(client, table)
t0 = time.time()
for row in gen_rows(n):
client.command(
"INSERT INTO learn_ck.events_perf (ts, uid, url, referer, duration) "
"VALUES",
parameters=row,
)
cost = time.time() - t0
print(f"[1] sync 1-by-1 n={n:>7d} cost={cost:7.2f}s "
f"parts={parts(client, table):>4d} qps={n/cost:>10.0f}/s")
def mode_async_insert(client) -> None:
"""async_insert:服务端攒批,wait_for_async_insert=0 让客户端立即返回。"""
table = "events_async"
truncate(client, table)
cols = ["ts", "uid", "url", "referer", "duration"]
settings = {
"async_insert": 1,
"wait_for_async_insert": 0,
"async_insert_busy_timeout_ms": 1000,
}
t0 = time.time()
for row in gen_rows(TOTAL):
client.insert(f"learn_ck.{table}", [row],
column_names=cols, settings=settings)
cost = time.time() - t0
time.sleep(2)
print(f"[2] async_insert n={TOTAL:>7d} cost={cost:7.2f}s "
f"parts={parts(client, table):>4d} qps={TOTAL/cost:>10.0f}/s")
def mode_buffer_engine(client) -> None:
"""Buffer 引擎:写 events_buffer 自动刷到 events_perf。"""
target = "events_perf"
truncate(client, target)
cols = ["ts", "uid", "url", "referer", "duration"]
t0 = time.time()
rows = list(gen_rows(TOTAL))
chunk = 5_000
for i in range(0, len(rows), chunk):
client.insert("learn_ck.events_buffer", rows[i:i + chunk],
column_names=cols)
cost = time.time() - t0
client.command("OPTIMIZE TABLE learn_ck.events_buffer")
time.sleep(1)
print(f"[3] Buffer engine n={TOTAL:>7d} cost={cost:7.2f}s "
f"parts(target)={parts(client, target):>4d} "
f"qps={TOTAL/cost:>10.0f}/s")
def main() -> None:
client = clickhouse_connect.get_client(host=HOST, port=PORT, username=USER)
print(f"[insert_modes] connected to {HOST}:{PORT}")
print("=" * 80)
mode_sync_one_by_one(client)
mode_async_insert(client)
mode_buffer_engine(client)
print("=" * 80)
print("提示:Part 数越少越好;同行数下 batch_size 越大、吞吐越高。")
if __name__ == "__main__":
main()