返回首页

将数据加载到ClickHouse:批量、async_insert、监控

将数据加载到ClickHouse的实用指南,采用正确的批处理策略。解释为什么小INSERT会创建数千个parts并扼杀性能。展示所有方法:INSERT VALUES(仅用于测试),通过clickhouse-client和HTTP API(包括gzip)加载CSV/JSONEachRow/TSV,INSERT SELECT用于数据库内复制,async_insert用于高频流(来自玩家的每秒数千事件)及缓冲区设置,参数max_insert_block_size和max_batch_size。提供用于通过system.query_log监控的SQL查询,以及一个现成的Python脚本,用于使用clickhouse-driver从parquet批量加载。所有示例基于betting.bets表。

ClickHouse:如何批量加载数据且不搞崩生产环境
Advertisement 728x90

将数据加载到 ClickHouse:如何告别逐行插入,将数据摄入速度提升 500 倍

我曾参与一个实时博彩分析项目。每天涌入 5000 万条事件。我天真地以为 INSERT INTO bets VALUES (...) 逐条插入没问题。ClickHouse 能处理每秒 500 次插入,但我们需要 2000 次。服务器不堪重负:磁盘队列暴涨,后台合并产生了成千上万个小 parts。生产环境直接宕机。

事实证明,ClickHouse 不是 MySQL。它专为批量插入而设计,而非单点 INSERT。将加载程序重写为每次批量插入 10 万条记录后,我实现了每秒 1 万次插入且性能无损。

下面介绍所有加载方法,从小型 CSV 到生产级批处理程序,以及那些让我彻夜难眠的陷阱。

Google AdInline article slot

1. 为什么单行插入是恶魔(即使你想这么做)

ClickHouse 每次 INSERT 都会在磁盘上创建一个新的 part。Part 包含 8192 行的颗粒,但即使只插入 1 行,也会物理创建包含微小 .bin 文件的新 part。

后果:

  • 1 万次 INSERT,每次 1 行 = 1 万个 parts
  • SELECT 查询必须打开 1 万个文件
  • 后台合并尝试合并它们,导致 CPU 紧张
  • 磁盘空间浪费在 part 元数据上

我的经验法则: 每批至少 1000 行或 1 MB 压缩数据。

Google AdInline article slot

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

Google AdInline article slot
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}

工作原理:

  1. 客户端发送小批量插入
  2. ClickHouse 将它们累积在内存缓冲区中
  3. 当缓冲区满(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_sizeasync_insert 这对组合。我在这上面栽了三次跟头——够了。


上一篇:
下一篇: ClickHouse中的SELECT查询:在PostgreSQL十年后如何重塑思维

— Editorial Team

Advertisement 728x90

继续阅读