ClickHouse 数据接入与 ETL

将海量数据高效写入 ClickHouse 是一门技术。本文详解 INSERT 优化、批量写入策略、Kafka 引擎集成、MaterializedMySQL、数据格式支持和分区设计。

1. 数据写入的核心原则

ClickHouse 的写入优化遵循一条核心原则:批量写入优于单条写入。

原因:

  • ClickHouse 的 part 是按批组织的,每一批写入产生一个新 part
  • 大量小 part 会导致后台 merge 压力剧增,影响查询性能
  • 批量写入能更好利用压缩和向量化处理

推荐批量大小:10,000 到 100,000 行/批次

2. 高效 INSERT 策略

2.1 基本批量插入

-- ❌ 单条插入(极慢)
INSERT INTO events VALUES (now(), 1, 'click', 10.0);
INSERT INTO events VALUES (now(), 2, 'view', 5.0);

-- ✅ 批量插入(推荐)
INSERT INTO events (event_time, user_id, event_type, value) VALUES
    (now(), 1, 'click', 10.0),
    (now(), 2, 'view', 5.0),
    (now(), 3, 'click', 15.0),
    ... -- 一次插入 1-10 万条

2.2 文件导入

# CSV 文件导入
cat data.csv | clickhouse-client --query="INSERT INTO events FORMAT CSV"

# JSON 导入
cat data.json | clickhouse-client --query="INSERT INTO events FORMAT JSONEachRow"

# 并行导入(加速)
cat large.csv | clickhouse-client \
    --query="INSERT INTO events FORMAT CSV" \
    --max_insert_block_size=100000

# 本地文件直接查询后插入
clickhouse-client --query="
    INSERT INTO events
    SELECT * FROM file('data.csv', CSV, 'event_time DateTime, user_id UInt64')
"

2.3 HTTP 接口写入

# 使用 HTTP 批量写入
curl -X POST 'http://localhost:8123/?query=INSERT%20INTO%20events%20FORMAT%20JSONEachRow' \
    -d '{"event_time":"2024-01-01 00:00:00","user_id":1,"event_type":"click"}'

Python 批量写入示例:

import clickhouse_driver
import json

client = clickhouse_driver.Client('localhost')

# 缓冲写入器
batch = []
BATCH_SIZE = 50000

with open('events.jsonl', 'r') as f:
    for line in f:
        event = json.loads(line)
        batch.append(event)
        
        if len(batch) >= BATCH_SIZE:
            client.execute(
                'INSERT INTO events (event_time, user_id, event_type, value) VALUES',
                batch
            )
            batch = []
            print(f"Inserted {BATCH_SIZE} rows")
    
    # 写入剩余数据
    if batch:
        client.execute(
            'INSERT INTO events (event_time, user_id, event_type, value) VALUES',
            batch
        )

2.4 异步插入

对于写吞吐量极高的场景,可以使用异步插入:

SET async_insert = 1;
SET async_insert_max_data_size = 1000000;  -- 缓冲 1MB 后批量写入

INSERT INTO events VALUES (now(), 1, 'click', 10.0);
INSERT INTO events VALUES (now(), 2, 'view', 5.0);
-- 数据先进入缓冲区,满足条件后自动批量写入

3. Kafka 引擎集成

Kafka 引擎是 ClickHouse 最常用的实时数据接入方式。

3.1 创建 Kafka 引擎表

-- 1. Kafka 消费表(接收消息)
CREATE TABLE events_kafka (
    user_id UInt64,
    event_type String,
    event_time DateTime,
    value Float64
) ENGINE = Kafka()
SETTINGS
    kafka_broker_list = 'kafka1:9092,kafka2:9092',
    kafka_topic_list = 'user-events',
    kafka_group_name = 'clickhouse-consumer',
    kafka_format = 'JSONEachRow',
    kafka_row_delimiter = '\n',
    kafka_num_consumers = 4,         -- 并行消费者数
    kafka_max_block_size = 65536;    -- 每批大小

-- 2. 目标表(存储数据)
CREATE TABLE events (
    user_id UInt64,
    event_type String,
    event_time DateTime,
    value Float64
) ENGINE = MergeTree()
ORDER BY (event_time, user_id);

-- 3. 物化视图(自动搬运数据)
CREATE MATERIALIZED VIEW events_mv TO events AS
SELECT * FROM events_kafka;

-- 现在 Kafka 消息会自动消费并写入 events 表

3.2 Kafka 配置优化

-- 提高消费吞吐量
ALTER TABLE events_kafka MODIFY SETTING
    kafka_max_block_size = 1048576,     -- 1MB 批次
    kafka_num_consumers = 8,            -- 8 个消费者
    kafka_poll_timeout_ms = 100;        -- 100ms 轮询超时

-- 错误处理
-- kafka_skip_broken_messages = 1000  -- 跳过格式错误的 JSON

3.3 监控 Kafka 消费

-- 查看消费延迟
SELECT
    database,
    table,
    num_messages,
    last_used,
    last_error
FROM system.kafka_consumers;

4. 数据格式支持

ClickHouse 支持丰富的数据导入格式:

-- CSV
INSERT INTO events FORMAT CSV
2024-01-01 00:00:00,1,click,10.0
2024-01-01 00:00:01,2,view,5.0

-- TSV(Tab 分隔)
INSERT INTO events FORMAT TabSeparated

-- JSONEachRow(每行一个 JSON 对象)
INSERT INTO events FORMAT JSONEachRow
{"event_time":"2024-01-01 00:00:00","user_id":1}

-- Parquet(列式存储格式)
INSERT INTO events SELECT * FROM file('data.parquet', Parquet)

-- ORC
INSERT INTO events SELECT * FROM file('data.orc', ORC)

-- 自定义分隔符
INSERT INTO events FORMAT CustomSeparated
SETTINGS format_custom_row_before_delimiter='', format_custom_field_delimiter='|'

5. 分区设计策略

5.1 分区键选择

-- 按天分区(推荐,最常用)
PARTITION BY toYYYYMMDD(event_time)

-- 按月分区(数据量极大,分区数不宜过多)
PARTITION BY toYYYYMM(event_time)

-- 按事件类型分区(类型有限且查询常按类型过滤)
PARTITION BY event_type

-- 组合分区
PARTITION BY (toYYYYMM(event_time), event_type)

-- 使用城市/地区(地理分布式查询)
PARTITION BY city

5.2 分区数量控制

-- 查看分区数量和大小
SELECT
    partition,
    count() as parts,
    sum(bytes) as total_bytes,
    sum(rows) as total_rows
FROM system.parts
WHERE table = 'events' AND active
GROUP BY partition
ORDER BY partition;

-- 表分区数不宜超过几千个(每个 part 都有文件句柄开销)

6. MaterializedMySQL

ClickHouse 可以直接同步 MySQL 数据:

-- 创建 MaterializedMySQL 数据库
CREATE DATABASE mysql_replica
ENGINE = MaterializedMySQL('mysql-host:3306', 'mydb', 'user', 'password')
SETTINGS
    allows_query_when_mysql_lost = true,
    max_wait_time_when_mysql_unavailable = 10000;

-- 自动同步 MySQL 的所有表
-- MySQL binlog → ClickHouse 实时同步

-- 查询同步过来的数据
SELECT * FROM mysql_replica.users WHERE id = 12345;

7. 数据类型与 Schema 设计

7.1 优化数据类型

-- 使用最合适的类型
CREATE TABLE events_optimized (
    event_time DateTime CODEC(Delta, LZ4),
    -- DateTime 比 String 存储时间更高效
    
    user_id UInt32,
    -- 如果用户数 < 40 亿,UInt32 比 UInt64 省一半空间
    
    event_type LowCardinality(String),
    -- LowCardinality: 内部字典编码,重复值多时省 90% 空间
    
    is_premium UInt8,
    -- UInt8 代替 Boolean/String
    
    revenue Decimal(10, 2),
    -- Decimal 代替 Float 存金额(精确计算)
    
    tags Array(LowCardinality(String)),
    -- 数组类型存储标签列表
    
    metadata Map(String, String),
    -- Map 类型存储键值对
) ENGINE = MergeTree()
ORDER BY (event_time, user_id);

7.1a Nested 嵌套数据类型

对于具有层级关系的数据,ClickHouse 的 Nested 类型是理想选择:

-- 存储订单和订单项
CREATE TABLE orders (
    order_id UInt64,
    order_time DateTime,
    customer_id UInt64,
    -- Nested: 每个订单包含多个商品项
    items.id Array(UInt64),
    items.name Array(String),
    items.price Array(Decimal(10, 2)),
    items.quantity Array(UInt32)
) ENGINE = MergeTree()
ORDER BY (order_time, order_id);

-- 查询嵌套数据
SELECT
    order_id,
    arraySum((price, qty) -> price * qty, items.price, items.quantity) AS total
FROM orders
WHERE order_time > '2024-01-01';

嵌套类型在内部实际上是多个同名的 Array 列,ClickHouse 会自动维护它们的行对齐关系。

7.2 预设默认值

CREATE TABLE events (
    event_time DateTime DEFAULT now(),
    user_id UInt64,
    event_type String DEFAULT 'unknown',
    value Float64 DEFAULT 0.0
) ENGINE = MergeTree()
ORDER BY (event_time, user_id);

-- 插入时省略带默认值的列
INSERT INTO events (user_id, event_type) VALUES (1, 'click');
-- event_time 自动填充为 now(),value 为 0.0

7.3 写入性能调优

-- 调整批量写入缓冲区大小
SET min_insert_block_size_rows = 1048576;      -- 100万行
SET min_insert_block_size_bytes = 268435456;   -- 256MB

-- 异步插入模式(前端低延迟)
SET async_insert = 1;
SET async_insert_max_data_size = 10000000;     -- 10MB 缓冲
SET async_insert_busy_timeout_ms = 1000;       -- 最多 1 秒缓冲

-- 关闭 fsync 提高写入速度(有数据丢失风险)
SET insert_deduplicate = 0;  -- 如果不需要去重

在实际生产环境中,写入性能调优应该根据硬件配置进行。SSD 上的 ClickHouse 写入吞吐量通常可达 100MB/s 到 500MB/s,具体取决于数据压缩率、CPU 速度和并发写入数。

8. 数据质量检查

-- 写入后验证数据量
SELECT count() FROM events WHERE event_time >= today();

-- 检查重复数据
SELECT user_id, event_time, count() AS cnt
FROM events
GROUP BY user_id, event_time
HAVING cnt > 1
LIMIT 10;

-- 检查 Null 值(ClickHouse 默认 Nullable 列)
SELECT count() FROM events WHERE user_id IS NULL;

9. 总结

高效数据接入的关键要点:

策略方法效果
批量写入1-10 万行/批次10x+ 吞吐量提升
Kafka 集成Kafka Engine + MV实时流式接入
分区设计按天/月分区快速删除 + 查询剪枝
异步写入async_insert = 1前端低延迟
数据类型优化UInt32、LowCardinality50%+ 空间节省
MySQL 同步MaterializedMySQL零代码迁移

数据接入是 ClickHouse 应用的第一步,也是最容易出现性能瓶颈的环节。遵循"批量优先、异步缓冲、格式对齐"的原则,可以确保后续分析查询的高效运行。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「数据库」更多文章

  1. ClickHouse 表引擎详解
  2. ClickHouse 监控与运维
  3. ClickHouse 生产案例与最佳实践