Skip to content

第 7 章 写入与读取最佳实践

学习目标:彻底搞懂「为什么 ClickHouse 怕高频小批量写」「如何写得又快又稳」「读取并发与内存怎么控」。学完这一章,你能说出 async_insertBuffer 引擎、INSERT ... SELECTFORMAT 全家桶分别在什么时候用,并能用 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_insert0总开关,置 1 启用
async_insert_max_data_size10 MBbuffer 大小阈值,达到就 flush
async_insert_busy_timeout_ms200 msbuffer 时间阈值,超过就 flush
async_insert_max_query_number450buffer 中允许的最大 INSERT 语句数
wait_for_async_insert1客户端是否等 flush 完成才返回
wait_for_async_insert_timeout120 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 主战场」语法。要点:

  1. 服务端流式:数据不会先回客户端再写回去,而是 server 内部 pipeline,吞吐极高。

  2. 支持 SETTINGS:可以单独控制并发、内存。

    sql
    INSERT INTO target SELECT ... FROM source
    SETTINGS max_insert_threads = 8,
             max_insert_block_size = 1048576;
  3. 跨集群:搭配 remote() / cluster() 表函数可以在两个 ClickHouse 之间搬。

    sql
    INSERT INTO local_events
    SELECT * FROM remote('10.0.0.1:9000', learn_ck, events, 'default', '');
  4. 跨库类型:搭配 mysql() / postgresql() / s3() / url() 表函数可以从外部数据源直接灌进来。

    sql
    INSERT 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写入耗时写入吞吐文件大小适用场景
Values38 s~26 万行/s小数据 / 调试
JSONEachRow22 s~45 万行/s跨语言、易读
CSVWithNames11 s~91 万行/s数据交换
TabSeparated9 s~111 万行/s简单 ETL
RowBinary4 s~250 万行/s自有 pipeline
Native2.6 s~385 万行/s最小ClickHouse 间互导
Parquet3.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.parquet

HTTP 端:

bash
curl 'http://127.0.0.1:8123/?query=INSERT%20INTO%20events%20FORMAT%20JSONEachRow' \
     --data-binary @data.ndjson

📌 一条经验:服务端之间互导 → Native;落数据湖 → Parquet;和外部脚本交换 → CSVWithNamesJSONEachRow别用 Values,它是给人看不是给机器跑的。


7.7 读取并发与资源控制

写入控好了,读取这边也别拉胯。ClickHouse 的并发模型是单查询多线程,所以同样要控阈值。

7.7.1 关键参数

参数默认作用
max_threads物理核数单条查询使用的最大线程数
max_block_size65 505单次 pipeline 处理的最大行数
max_insert_block_size1 048 576INSERT 切块大小
max_memory_usage10 GB单条查询最大内存
max_memory_usage_for_user0(无限)单个用户所有查询合计内存
max_bytes_before_external_group_by0GROUP BY 内存阈值,超过就 spill 到磁盘
max_bytes_before_external_sort0ORDER BY 同上
max_execution_time0查询超时秒数
priority0查询优先级(数字越小越高)

7.7.2 怎么调

单查询临时调

sql
SELECT count() FROM events
SETTINGS max_threads = 16, max_memory_usage = 20000000000;

给某用户固定调(在 users.xmlCREATE 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=1async_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 INSERTPostgreSQL COPYClickHouse INSERT
推荐批次几千 ~ 几万几万 ~ 几十万≥ 10 万
协议层优化多 VALUES + LOAD DATACOPY 二进制流Native 列式协议
批小的代价binlog 体积膨胀WAL 翻倍Part 暴涨 → 拒绝写入
服务端攒批无原生支持无原生支持async_insert / Buffer
列存?
跨实例搬mysqldump / pt-archiverpg_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 单条」的成本结构。

标准答案

  1. ClickHouse 的写入单位是 Part,每次 INSERT(不管 1 行还是 100 万行)都会走「内存排序 → 列拆分 → 压缩 → 写 primary.idx → 写 .mrk → 落盘 → 等待 Merge」一整套流水线,成本接近常数。
  2. 单条插入会让磁盘上瞬间出现海量小 Part,后台 Merge 跟不上,触发 parts_to_delay_insert(默认 150)和 parts_to_throw_insert(默认 300),轻则限流重则拒写。
  3. 列存压缩需要「一段连续的同类值」才能发挥效率,单条插入意味着每列只有 1 个值,压缩比退化。
  4. 正确姿势是客户端攒批(≥10 万行),或开启 async_insert 让服务端帮你攒批,或用 Buffer 引擎做内存缓冲层。

加分项:能补一句「Part 数对查询也有副作用」—— 查询时 ClickHouse 要并行扫描所有满足分区裁剪条件的 active Part,Part 越多元数据越重,schedule 开销越大。

易错点:把 parts_to_throw_insert 直接调高来压制错误。这是治标不治本,根本问题是写入姿势错了。


Q2:async_insert 的工作原理与适用场景?

考察点:是否真用过 async_insert,能不能讲清 buffer / flush 模型。

标准答案

  1. async_insert = 1 时,服务端为每对(query SQL + 用户 + settings)维护一个内存 query buffer
  2. 客户端的 INSERT 进入 buffer 后立即返回 ack(取决于 wait_for_async_insert),不等真实落盘。
  3. 当 buffer 满足任一条件时整批 flush 成一个 Part:
    • 字节量 ≥ async_insert_max_data_size(默认 10 MB)
    • 时间 ≥ async_insert_busy_timeout_ms(默认 200 ms)
    • INSERT 语句数 ≥ async_insert_max_query_number(默认 450)
  4. 适用场景:客户端无法攒批(埋点 SDK / IoT 设备 / 网关日志)、需要尽快 ack。
  5. 不适用场景:需要「写入立即 SELECT 看到」、不能容忍服务端 crash 时的微量数据丢失(除非 wait_for_async_insert=1)。

加分项

  • 提到 system.asynchronous_insertssystem.asynchronous_insert_log 用于排错。
  • 23.x 版本起支持 deduplicate_blocks_in_dependent_materialized_views 与 async_insert 配合。

易错点:把 async_insert 当成「批量插入的替代品」。它是「把攒批从客户端搬到服务端」,没有改变批的本质


Q3:Buffer 引擎和 async_insert 的区别?

考察点:两个相似工具的取舍。

标准答案

维度async_insertBuffer 引擎
形态一个 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-clientclickhouse-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 后引入了字典编码进一步压缩。

易错点:把 NativeRowBinary 混淆,前者列式后者行式。


Q5:max_threads 一定是越大越好吗?

考察点:对 ClickHouse 单查询并发模型的理解。

标准答案

不是。max_threads 控制的是单条查询能并行的 pipeline 线程数,过大会带来:

  1. 上下文切换 overhead:当机器只有 8 核而 max_threads=64 时,OS 调度成本反超并行收益。
  2. 内存放大:每个线程都要维护自己的 hash table / sort 状态,内存占用 ≈ N×。
  3. 小查询反而变慢:scan 量小时,调度多线程的固定成本就吃掉了红利。
  4. 多用户场景互相挤占:单查询吃满 32 核,其他用户的查询排队。

正确做法

  • 默认 = 物理核数;
  • 大查询(亿级 scan)保持默认或略调高;
  • 小查询(百万级以内)反而适当调小(如 4),降低调度成本;
  • 在多租户环境中,给不同 profile 设不同的 max_threads 上限,配合 max_concurrent_queries_for_user

加分项:能提到 max_block_sizemax_threads 的相互作用 —— block 太小会放大调度成本,block 太大会放大内存峰值。

易错点:以为 max_threads 是「全实例并发」,其实是「单查询并发」;总并发要看 max_concurrent_queries


📌 下一章预告:第 8 章我们进入 ClickHouse 真正的看家本领 —— 聚合、窗口与数组uniquniqExact 的取舍、-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()

insert_modes.py ↗