将数据加载到 ClickHouse:如何告别逐行插入,将数据摄入速度提升 500 倍
我曾参与一个实时博彩分析项目。每天涌入 5000 万条事件。我天真地以为 INSERT INTO bets VALUES (...) 逐条插入没问题。ClickHouse 能处理每秒 500 次插入,但我们需要 2000 次。服务器不堪重负:磁盘队列暴涨,后台合并产生了成千上万个小 parts。生产环境直接宕机。
事实证明,ClickHouse 不是 MySQL。它专为批量插入而设计,而非单点 INSERT。将加载程序重写为每次批量插入 10 万条记录后,我实现了每秒 1 万次插入且性能无损。
下面介绍所有加载方法,从小型 CSV 到生产级批处理程序,以及那些让我彻夜难眠的陷阱。
1. 为什么单行插入是恶魔(即使你想这么做)
ClickHouse 每次 INSERT 都会在磁盘上创建一个新的 part。Part 包含 8192 行的颗粒,但即使只插入 1 行,也会物理创建包含微小 .bin 文件的新 part。
后果:
- 1 万次
INSERT,每次 1 行 = 1 万个 parts SELECT查询必须打开 1 万个文件- 后台合并尝试合并它们,导致 CPU 紧张
- 磁盘空间浪费在 part 元数据上
我的经验法则: 每批至少 1000 行或 1 MB 压缩数据。
2. INSERT ... VALUES —— 仅用于测试,不用于生产
-- 仅用于本地实验
INSERT INTO betting.bets (user_id, created_at, amount, odds, sport, outcome)
VALUES
(1001, now(), 50.00, 2.1, 'football', 'win'),
(1002, now(), 100.00, 1.8, 'basketball', 'loss'),
(1003, now(), 75.00, 3.0, 'tennis', 'win');
在生产环境中,永远不要这样做。即使你在 VALUES 中批量插入 1000 行,SQL 解析也会成为瓶颈。
3. INSERT ... FORMAT CSV / JSONEachRow / TSV —— 文件的救星
ClickHouse 支持多种格式的即时解析。
CSV 文件示例 bets.csv:
1001,2024-03-15 10:00:00,50.00,2.1,football,win
1002,2024-03-15 10:01:00,100.00,1.8,basketball,loss
1003,2024-03-15 10:02:00,75.00,3.0,tennis,win
加载:
INSERT INTO betting.bets FORMAT CSV
但你需要将数据通过标准输入传入。因此在脚本中,使用:
clickhouse-client --query="INSERT INTO betting.bets FORMAT CSV" < bets.csv
JSONEachRow —— 日志的理想选择:
{"user_id":1001,"created_at":"2024-03-15 10:00:00","amount":50.00,"odds":2.1,"sport":"football","outcome":"win"}
{"user_id":1002,"created_at":"2024-03-15 10:01:00","amount":100.00,"odds":1.8,"sport":"basketball","outcome":"loss"}
加载:
clickhouse-client --query="INSERT INTO betting.bets FORMAT JSONEachRow" < bets.json
TabSeparated —— 最快(解析最少):
clickhouse-client --query="INSERT INTO betting.bets FORMAT TSV" < bets.tsv
4. 通过 clickhouse-client 从文件加载 —— 我最喜欢的备份方式
# 小文件
clickhouse-client --query="INSERT INTO betting.bets FORMAT CSV" < /tmp/daily_bets.csv
# 大文件并显示进度
clickhouse-client \
--query="INSERT INTO betting.bets FORMAT CSV" \
--progress \
--max_insert_block_size=100000 \
< /data/historical_bets_2023.csv
重要标志: --max_insert_block_size=100000 —— 将文件分割成 10 万行的块。这样可以在不内存溢出的情况下插入 TB 级文件。
我学到的: FORMAT 定义输入数据的结构,而非输出。注意查询本身不包含数据,数据通过标准输入传入。
5. 使用 POST 的 HTTP API —— 当没有命令行访问权限时
# 通过 curl 使用 CSV
curl -X POST "http://localhost:8123/?query=INSERT+INTO+betting.bets+FORMAT+CSV" \
--data-binary @bets.csv
# JSONEachRow
curl -X POST "http://localhost:8123/?query=INSERT+INTO+betting.bets+FORMAT+JSONEachRow" \
--data-binary @bets.json \
-u analyst:password
# 即时 GZIP(节省带宽)
gzip -c bets.csv | curl -X POST "http://localhost:8123/?query=INSERT+INTO+betting.bets+FORMAT+CSV" \
--data-binary @- \
--header "Content-Encoding: gzip"
为什么用 POST 而非 GET: GET 请求会被完整记录,包括数据。如果你通过 ?query=INSERT... 插入,密码和数据可能会出现在代理日志中。
6. INSERT SELECT —— 在 ClickHouse 内部移动数据
最快的大规模重载方式不是导出到外部。
-- 将一月份数据复制到归档表
INSERT INTO betting.bets_archive
SELECT * FROM betting.bets
WHERE created_at >= '2024-01-01' AND created_at < '2024-02-01';
-- 即时聚合
INSERT INTO betting.daily_stats (date, total_bets, total_amount)
SELECT
toDate(created_at) AS date,
count() AS total_bets,
sum(amount) AS total_amount
FROM betting.bets
WHERE created_at >= today() - 7
GROUP BY date;
速度: INSERT SELECT 避免了客户端-服务器交换,是本地传输。数十亿行可以在几分钟内移动。
7. 异步插入:当每秒来自玩家的数千个事件时
传统的 INSERT 是同步的:客户端等待数据写入磁盘。在每秒 1 万次插入时,这很痛苦。
async_insert 在内存中缓冲数据并批量插入。
在 config.xml 或通过 SETTINGS 启用:
-- 会话级别
SET async_insert = 1;
SET wait_for_async_insert = 0; -- 不等待完成
SET async_insert_max_data_size = 10000000; -- 10 MB 缓冲区
SET async_insert_busy_timeout_ms = 200; -- 每 200 毫秒刷新
INSERT INTO betting.bets FORMAT JSONEachRow
{"user_id":1001,"amount":50,"odds":2.1}
{"user_id":1002,"amount":100,"odds":1.8}
工作原理:
- 客户端发送小批量插入
- ClickHouse 将它们累积在内存缓冲区中
- 当缓冲区满(10 MB)或超时(200 毫秒)时,执行一次批量插入
让我吃亏的地方: wait_for_async_insert=0 意味着客户端不知道插入是否失败。网络问题可能导致数据丢失。如果可靠性重要,设置 wait_for_async_insert=1。性能会下降,但不会太严重。
8. 设置:max_insert_block_size 和 max_batch_size
<!-- config.xml -->
<max_insert_block_size>1048576</max_insert_block_size> <!-- 1M 行 -->
<max_block_size>65536</max_block_size>
max_insert_block_size —— ClickHouse 一次处理的最大块大小。如果你的插入更大,它会被分割成块。
生产规则: 将 max_insert_block_size 设置为预期批次的 2-4 倍小。这样,即使意外的大 INSERT 也不会耗尽内存。
# 为单次插入强制设置
clickhouse-client --query="INSERT INTO bets FORMAT CSV" \
--max_insert_block_size=500000 < huge_file.csv
9. 通过 system.query_log 监控插入
永远不要猜测插入了多少行。检查 system.query_log。
-- 最近的 INSERT 查询
SELECT
event_time,
query_duration_ms,
read_rows,
written_rows,
result_rows,
memory_usage,
query
FROM system.query_log
WHERE type = 'QueryFinish'
AND query LIKE '%INSERT INTO betting.bets%'
AND event_time >= now() - INTERVAL 1 HOUR
ORDER BY event_time DESC
LIMIT 20;
我的监控仪表板: 显示过去一小时的平均批次大小、失败插入次数、行/秒速度等指标。
SELECT
toStartOfMinute(event_time) AS minute,
count() AS inserts,
avg(written_rows) AS avg_batch_size,
sum(written_rows) AS total_rows,
avg(query_duration_ms) AS avg_latency_ms
FROM system.query_log
WHERE type = 'QueryFinish'
AND query LIKE '%INSERT INTO betting.bets%'
AND event_time >= now() - INTERVAL 1 HOUR
GROUP BY minute
ORDER BY minute DESC;
10. 用于批量加载历史数据的 Python 脚本
在一个真实项目中,我们从 S3 的 parquet 文件加载了 3 年的博彩历史数据。以下脚本处理了 20 亿行:
#!/usr/bin/env python3
from clickhouse_driver import Client
import pandas as pd
import glob
from tqdm import tqdm
# 连接到 ClickHouse
client = Client(
host='clickhouse.prod.internal',
port=9000,
user='loader',
password='strong_password',
database='betting'
)
# 针对大插入的优化
client.execute("SET max_insert_block_size = 1000000")
client.execute("SET min_insert_block_size_rows = 500000")
client.execute("SET min_insert_block_size_bytes = 100000000") # 100 MB
def insert_batch(df, table='bets'):
"""通过 Native 协议批量插入 DataFrame"""
# clickhouse-driver 自动转换类型
client.execute(f"INSERT INTO {table} FORMAT TabSeparated",
df.to_csv(sep='\t', header=False, index=False),
types_check=True)
def stream_parquet_files(pattern='/data/bets_*.parquet'):
files = glob.glob(pattern)
batch_size = 500000
buffer = []
for file in tqdm(files, desc="处理文件"):
# 分块读取 parquet 以避免内存占用过多
for chunk in pd.read_parquet(file, chunksize=batch_size):
# 为 ClickHouse 转换类型
chunk['created_at'] = pd.to_datetime(chunk['created_at'])
chunk['amount'] = chunk['amount'].astype('float64')
chunk['odds'] = chunk['odds'].astype('float64')
buffer.append(chunk)
if len(buffer) >= 5: # 5 个 50 万行的块 = 250 万行
df_batch = pd.concat(buffer, ignore_index=True)
insert_batch(df_batch)
buffer = []
# 剩余数据
if buffer:
df_batch = pd.concat(buffer, ignore_index=True)
insert_batch(df_batch)
if __name__ == '__main__':
stream_parquet_files()
通过 requests 的替代方案(无需 pandas):
import requests
import gzip
import io
def insert_via_http_csv_gz(filepath):
with open(filepath, 'rb') as f:
compressed = gzip.compress(f.read())
response = requests.post(
'http://localhost:8123/?query=INSERT+INTO+betting.bets+FORMAT+CSV',
data=compressed,
headers={'Content-Encoding': 'gzip'},
auth=('loader', 'password')
)
if response.status_code != 200:
raise Exception(f"插入失败: {response.text}")
print(f"已插入 {filepath}")
# 批量发送多个文件的示例
for file in glob.glob('/data/bets_*.csv.gz'):
insert_via_http_csv_gz(file)
常见加载错误(以及我如何修复它们)
错误:Memory limit (total) exceeded
原因: 插入块过大。
解决方案: 减小 max_insert_block_size 或将文件分割成小块。
错误:Too many parts
原因: 过多的小 INSERT 或高吞吐量下 async_insert 被禁用。
解决方案: 启用 async_insert,增加 min_rows_for_wide_part。
错误:Code: 252. DB::Exception: Too many simultaneous inserts
原因: 并行插入同一分区。
解决方案: 为该表设置 max_insert_threads = 1。
错误:分区不合并
原因: 插入的数据不符合 ORDER BY 顺序。
解决方案: 确保数据按表的 ORDER BY 排序。
下一步
现在你知道了如何通过任何方法加载数据——从小型 CSV 到生产流。下一篇文章将介绍查询优化和高级聚合。
临别建议:如果你在生产环境中工作,任何 INSERT 都不应缺少 max_insert_block_size 和 async_insert 这对组合。我在这上面栽了三次跟头——够了。
← 上一篇: ClickHouse中的MergeTree:引擎如何将分析数据切分为颗粒并合并分区
→ 下一篇: ClickHouse中的SELECT查询:在PostgreSQL十年后如何重塑思维
— Editorial Team
暂无评论。