07. ClickHouse 高性能分析引擎

深入 ClickHouse 列式存储引擎:MergeTree 引擎家族、 sparse primary index、向量化执行、物化视图、集群分片与查询优化实战。

ClickHouse 是由 Yandex 开源的列式 OLAP 数据库,凭借其极致的单机查询性能和高压缩比,已成为实时分析场景的事实标准。本文系统讲解其存储引擎、索引机制、集群架构与生产优化。

1. 列式存储与向量化执行

1.1 行存 vs 列存对比

特性行式存储 (MySQL)列式存储 (ClickHouse)
存储方式按行存储按列存储
读取列需读取整行只读目标列
压缩率低(数据类型混杂)高(同类型连续)
OLTP 查询优秀较差
OLAP 聚合较差极优秀
单行写入
批量写入一般极快

1.2 数据物理存储

表目录结构:
/var/lib/clickhouse/data/db/table/
├── detached/                  ← 已分离的分区
├── format_version.txt
└── 202401_1_3_1/             ← 数据分区目录 (Partition + MinBlock + MaxBlock + Level)
    ├── checksums.txt          ← 文件校验和
    ├── columns.txt            ← 列定义
    ├── count.txt              ← 行数
    ├── primary.idx            ← 主键稀疏索引
    ├── skp_idx_*.idx/mrk      ← 跳数索引
    ├── minmax_timestamp.idx   ← 分区键索引
    ├── data.bin               ← 列数据文件 (LZ4/ZSTD 压缩)
    ├── data.mrk2              ← 数据标记文件 (offset)
    └── ... 其他列文件

1.3 向量化执行引擎

ClickHouse 采用 SIMD(单指令多数据) + 列式批量处理 实现极致性能:

传统火山模型(逐行处理):
for row in table.rows:
    a = col1[row]
    b = col2[row]
    if a > 100:
        sum += b

ClickHouse 向量化执行(批量处理):
chunk = [10000 行]
mask = col1[chunk] > 100        ← SIMD 并行比较
filtered = compress(col2[chunk], mask)  ← 只保留有效数据
sum = reduce(filtered)          ← 快速聚合
-- 查看是否使用向量化执行
EXPLAIN PIPELINE SELECT sum(amount) FROM orders WHERE amount > 100;
-- 输出中的 ExpressionTransform 即向量化算子

2. MergeTree 引擎家族

2.1 引擎对比

引擎特点适用场景写入性能查询性能
MergeTree基础引擎,按主键排序合并通用场景
ReplacingMergeTree自动去重(按主键去重)幂等写入、CDC
SummingMergeTree自动聚合数值列预聚合指标极快
AggregatingMergeTree自动聚合聚合函数状态实时聚合极快
CollapsingMergeTree行级更新(Sign 标记)状态更新
VersionedCollapsingMergeTree版本化折叠带版本的状态
GraphiteMergeTree专为 Graphite 优化监控时序

2.2 MergeTree 核心原理

CREATE TABLE orders (
    order_id UInt64,
    user_id UInt32,
    amount Decimal(18,2),
    status UInt8,
    create_time DateTime,
    city String
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(create_time)        -- 按月分区
ORDER BY (city, create_time, order_id)    -- 主键顺序(稀疏索引)
PRIMARY KEY (city, create_time)           -- 主键(用于索引)
SETTINGS index_granularity = 8192;         -- 索引粒度

数据合并机制

写入 → 生成新 Part (202401_10_10_0)
         ↓
后台 Merge → 合并相邻 Part (202401_1_10_1)
                  ↓
         同一分区内的 Part 会自动合并
         合并过程: 排序 → 去重(Replacing) → 聚合(Summing)

2.3 ReplacingMergeTree(去重)

-- 按 order_id 去重,保留最新 version
CREATE TABLE user_events (
    user_id UInt64,
    event_type String,
    event_time DateTime,
    properties String,
    version UInt32           -- 版本号,用于冲突解决
) ENGINE = ReplacingMergeTree(version)
PARTITION BY toYYYYMMDD(event_time)
ORDER BY (user_id, event_type, event_time);

-- 查询时通常需要 FINAL 强制合并
SELECT * FROM user_events FINAL WHERE user_id = 123;
-- 或依赖后台 merge(非实时去重)

2.4 SummingMergeTree(预聚合)

-- 自动按主键汇总数值列
CREATE TABLE orders_daily (
    dt Date,
    city String,
    category String,
    order_count UInt64,       -- 会被聚合
    total_amount Decimal(18,2), -- 会被聚合
    -- 非数值列保留首次插入值
    first_order_id UInt64
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(dt)
ORDER BY (dt, city, category);

-- 批量写入明细(小文件暂存)
INSERT INTO orders_daily VALUES ('2024-01-01', 'Beijing', 'Electronics', 1, 999.00, 1001);
INSERT INTO orders_daily VALUES ('2024-01-01', 'Beijing', 'Electronics', 1, 599.00, 1002);

-- 后台合并后自动聚合为:
-- ('2024-01-01', 'Beijing', 'Electronics', 2, 1598.00, 1001)

-- 查询时务必加 FINAL
SELECT * FROM orders_daily FINAL;

2.5 AggregatingMergeTree(聚合状态)

-- 存储聚合函数中间状态,实时聚合查询
CREATE TABLE user_metrics (
    user_id UInt64,
    dt Date,
    -- 使用 AggregateFunction 类型存储中间状态
    total_amount AggregateFunction(sum, Decimal(18,2)),
    order_count AggregateFunction(count, UInt64),
    unique_cities AggregateFunction(uniqExact, String),
    max_amount AggregateFunction(max, Decimal(18,2))
) ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(dt)
ORDER BY (user_id, dt);

-- 写入需使用 -State 后缀
INSERT INTO user_metrics
SELECT
    user_id,
    today() as dt,
    sumState(amount) as total_amount,
    countState() as order_count,
    uniqExactState(city) as unique_cities,
    maxState(amount) as max_amount
FROM orders
GROUP BY user_id;

-- 查询需使用 -Merge 后缀合并状态
SELECT
    user_id,
    sumMerge(total_amount) as total,
    countMerge(order_count) as cnt,
    uniqExactMerge(unique_cities) as cities,
    maxMerge(max_amount) as max_amt
FROM user_metrics
GROUP BY user_id;

3. 索引与查询优化

3.1 主键稀疏索引

ClickHouse 使用 稀疏索引(非 B+ 树),每 index_granularity(默认 8192)行存储一个索引标记。

数据文件 (按主键排序):
[Beijing, 2024-01-01 00:00:00] row 1
[Beijing, 2024-01-01 00:01:00] row 2
...
[Beijing, 2024-01-01 01:00:00] row 8192  ← 索引标记点
[Beijing, 2024-01-01 02:00:00] row 16384 ← 索引标记点
[Shanghai, 2024-01-01 00:00:00] row ...

主键索引 (primary.idx):
[Beijing, 2024-01-01 01:00:00] → data offset 8192
[Beijing, 2024-01-01 02:00:00] → data offset 16384

查询 WHERE city = 'Shanghai' → 二分查找索引 → 直接定位数据范围

3.2 跳数索引(Skip Index)

CREATE TABLE events (
    user_id UInt64,
    event_time DateTime,
    url String,
    response_time UInt32,
    
    -- 跳数索引:加速特定过滤条件
    INDEX idx_url url TYPE bloom_filter GRANULARITY 4,
    INDEX idx_resp response_time TYPE minmax GRANULARITY 4,
    INDEX idx_set url TYPE set(100) GRANULARITY 4
) ENGINE = MergeTree()
ORDER BY (user_id, event_time);
索引类型说明适用场景
minmax存储每 granule 的 min/max数值/时间范围查询
set(n)存储每 granule 最多 n 个唯一值枚举值过滤
bloom_filterBloom 过滤器判断存在性等值查询(字符串/UUID)
tokenbf_v1支持分词的 Bloom Filter全文检索
ngrambf_v1N-gram Bloom Filter模糊匹配

3.3 查询优化实践

-- 1. 利用主键过滤(最左前缀)
SELECT * FROM orders WHERE city = 'Beijing' AND create_time > '2024-01-01';
-- 有效 ✓

SELECT * FROM orders WHERE order_id = 12345;
-- 低效 ✗(order_id 不在主键最左列)

-- 2. 分区裁剪
SELECT * FROM orders WHERE create_time >= '2024-01-01' AND create_time < '2024-02-01';
-- 只扫描 202401 分区

-- 3. 预过滤(Projection)
CREATE TABLE orders (
    ...
) ENGINE = MergeTree()
ORDER BY (city, create_time)
PROJECTION projection_by_user
(SELECT user_id, sum(amount) GROUP BY user_id);

-- 4. LIMIT 优化
SELECT * FROM orders ORDER BY create_time DESC LIMIT 10;
-- 若主键包含 create_time,无需全表排序

-- 5. 避免 SELECT *
SELECT user_id, amount FROM orders;  -- 只读取两列
SELECT * FROM orders;                 -- 读取所有列

4. 物化视图

4.1 同步物化视图(Projection)

-- 创建带物化视图的表
CREATE TABLE orders (
    order_id UInt64,
    user_id UInt64,
    amount Decimal(18,2),
    create_time DateTime,
    city String,
    
    PROJECTION projection_user_daily
    (
        SELECT 
            toDate(create_time) as dt,
            user_id,
            sum(amount) as total_amount,
            count() as order_count
        GROUP BY dt, user_id
    ),
    
    PROJECTION projection_city_monthly
    (
        SELECT 
            toYYYYMM(create_time) as month,
            city,
            sum(amount) as total_amount
        GROUP BY month, city
    )
) ENGINE = MergeTree()
ORDER BY (city, create_time, order_id);

4.2 物化视图(Materialized View)

-- 目标表
CREATE TABLE orders_summary (
    dt Date,
    city String,
    total_amount AggregateFunction(sum, Decimal(18,2)),
    order_count AggregateFunction(count, UInt64)
) ENGINE = AggregatingMergeTree()
ORDER BY (dt, city);

-- 物化视图(写入 orders 时自动触发)
CREATE MATERIALIZED VIEW orders_summary_mv
TO orders_summary
AS SELECT
    toDate(create_time) as dt,
    city,
    sumState(amount) as total_amount,
    countState() as order_count
FROM orders
GROUP BY dt, city;

-- 查询物化视图
SELECT 
    dt, city,
    sumMerge(total_amount) as gmv,
    countMerge(order_count) as cnt
FROM orders_summary
GROUP BY dt, city;

5. 集群部署与高可用

5.1 分片 + 副本架构

┌────────────────────────────────────────────────────┐
│                  Distributed Table                  │
│           CREATE TABLE dist ON CLUSTER '{cluster}' │
└──────────────┬──────────────────┬──────────────────┘
               │                  │
      ┌────────▼────────┐  ┌─────▼──────┐
      │   Shard 1       │  │  Shard 2   │
      │ ┌─────┐ ┌─────┐ │  │┌─────┐┌────┐│
      │ │Rep 1│ │Rep 2│ │  ││Rep 1││Rep2││
      │ │ :9000│ │ :9001│ │  ││:9000││:9001││
      │ └─────┘ └─────┘ │  │└─────┘└────┘│
      │  1/2 数据        │  │  1/2 数据   │
      └─────────────────┘  └─────────────┘
      
配置: <shard>01</shard><replica>01</replica>
     <shard>01</shard><replica>02</replica>
     <shard>02</shard><replica>01</replica>
     <shard>02</shard><replica>02</replica>

5.2 集群表配置

-- 1. 本地表(每个节点存储实际数据)
CREATE TABLE orders_local ON CLUSTER '{cluster}' (
    order_id UInt64,
    user_id UInt64,
    amount Decimal(18,2),
    create_time DateTime
) ENGINE = ReplicatedMergeTree('/clickhouse/{cluster}/tables/{shard}/orders', '{replica}')
PARTITION BY toYYYYMM(create_time)
ORDER BY (user_id, create_time);

-- 2. 分布式表(路由层)
CREATE TABLE orders_dist ON CLUSTER '{cluster}' AS orders_local
ENGINE = Distributed('{cluster}', 'default', 'orders_local', rand());

-- 3. 写入分布式表
INSERT INTO orders_dist VALUES (...);
-- 自动按 sharding_key (rand()) 分发到各 shard

-- 4. 查询分布式表
SELECT user_id, sum(amount) FROM orders_dist GROUP BY user_id;
-- 自动下发到各 shard → 本地聚合 → 汇总结果

5.3 集群写入优化

-- 方式一:直接写本地表(需客户端自己分片)
-- 适合离线批量导入
INSERT INTO orders_local VALUES (...);

-- 方式二:写分布式表 + 本地表双写
-- BI 查分布式表,实时应用写本地表

-- 方式三:异步 distributed_directory_monitor
-- clickhouse 自动将分布式目录数据分发到 shard

6. 生产运维

6.1 监控指标

指标说明告警阈值
ClickHouseProfileEvents_Query查询执行次数关注增长率
ClickHouseMetrics_Query当前执行查询数> 80% max_concurrent_queries
ClickHouseAsyncMetrics_DiskUsage磁盘使用率> 80%
ReplicatedMaxAbsoluteDelay副本延迟> 300s
Merge正在合并的 Part 数> 20
-- 查询当前正在执行的查询
SELECT * FROM system.processes WHERE is_cancelled = 0;

-- 查询慢查询日志
SELECT * FROM system.query_log WHERE event_time > now() - 3600 ORDER BY query_duration_ms DESC LIMIT 20;

-- 查看 Part 合并情况
SELECT database, table, partition, name, bytes_on_disk, modification_time 
FROM system.parts WHERE active = 1;

6.2 关键配置

<!-- config.xml 关键配置 -->
<max_concurrent_queries>100</max_concurrent_queries>
<max_memory_usage>50000000000</max_memory_usage>  <!-- 50GB -->
<max_execution_time>300</max_execution_time>       <!-- 5分钟超时 -->
<max_partitions_per_insert_block>100</max_partitions_per_insert_block>

<!-- 压缩配置 -->
<merge_tree>
    <min_compress_block_size>65536</min_compress_block_size>
    <max_compress_block_size>1048576</max_compress_block_size>
</merge_tree>

总结

决策场景推荐方案
通用分析表MergeTree
幂等写入/去重ReplacingMergeTree + FINAL
维度聚合分析SummingMergeTree / AggregatingMergeTree
实时聚合查询AggregatingMergeTree + 物化视图
高频更新场景CollapsingMergeTree
集群部署ReplicatedMergeTree + Distributed
索引优化主键有序 + 分区键 + 跳数索引
写入性能批量写入、适当调大 index_granularity

7. MergeTree 底层原理

7.1 分区目录结构

ClickHouse 将数据按 PARTITION BY 划分存储。每个分区由多个 Part 组成,Part 是数据不可变的最小单元。

/var/lib/clickhouse/data/default/orders/
├── detached/                          ← 已分离分区(ALTER DETACH 后存放)
├── format_version.txt
└── 202401_1_3_1_42/                  ← 分区目录格式解析:
    │                                   │   {Partition}_{MinBlock}_{MaxBlock}_{Level}_{Mutation}
    ├── checksums.txt                   ← 各文件校验和(防损坏)
    ├── columns.txt                     ← 列定义元数据
    ├── count.txt                       ← 当前 Part 行数
    ├── primary.idx                     ← 稀疏主键索引(每 N 行一个标记)
    ├── minmax_create_time.idx          ← 分区键 min/max 索引(快速裁剪)
    ├── data.bin                        ← 列数据(按列存储 + LZ4/ZSTD 压缩)
    ├── data.mrk2                       ← 数据标记(granule 到文件偏移映射)
    ├── skp_idx_url.idx / .mrk          ← 跳数索引文件
    └── default_compression_codec.txt   ← 压缩算法声明

目录命名规则Partition_MinBlock_MaxBlock_Level

  • Partition:分区值(如 202401
  • MinBlock / MaxBlock:写入 block 的编号范围
  • Level:合并次数,每次 merge 后 Level + 1
  • Mutation:可选的 mutation 版本号

7.2 Part 合并策略

写入过程:
INSERT batch → 生成 Part (202401_5_5_0)
                ↓
多个 Part 积累 → 后台 Merge 线程触发
                ↓
合并策略(按尺寸):
    Level 0 (0-1MB)    → 合并阈值为 4 个 Part
    Level 1 (1-10MB)   → 合并阈值为 4 个 Part
    Level 2 (10-100MB) → 合并阈值为 8 个 Part
    Level 3 (>100MB)   → 合并阈值为 8 个 Part
                ↓
合并后生成新 Part (202401_1_5_1)
                ↓
旧 Part 标记为 inactive → 后续清理
-- 查看 Part 合并状态
SELECT 
    database, table, partition, name,
    level, bytes_on_disk, rows,
    modification_time,
    active                                  -- 1=活跃, 0=待清理
FROM system.parts 
WHERE table = 'orders' AND active = 1
ORDER BY partition, name;

-- 手动触发合并(慎用,生产环境避免高峰期执行)
OPTIMIZE TABLE orders PARTITION '202401' FINAL;

-- 查看合并任务队列
SELECT * FROM system.merges;

调优参数users.xml profile):

<merge_tree>
    <!-- 触发合并的最小 Part 数 -->
    <parts_to_delay_insert>300</parts_to_delay_insert>
    <parts_to_throw_insert>600</parts_to_throw_insert>
    <!-- 最大 Part 尺寸 -->
    <max_bytes_to_merge_at_max_space_in_pool>107374182400</max_bytes_to_merge_at_max_space_in_pool>
    <!-- 旧数据合并阈值(降低冷数据合并频率) -->
    <min_age_to_force_merge_seconds>86400</min_age_to_force_merge_seconds>
    <min_age_to_force_merge_on_partition_only>false</min_age_to_force_merge_on_partition_only>
</merge_tree>

7.3 TTL 自动过期

-- 按时间自动删除旧数据
CREATE TABLE logs (
    event_time DateTime,
    message String,
    level String
) ENGINE = MergeTree()
ORDER BY event_time
TTL event_time + INTERVAL 90 DAY;           -- 90 天后自动删除

-- 按时间自动转移到冷存储(S3 / 另一磁盘卷)
CREATE TABLE logs_tiered (
    event_time DateTime,
    message String
) ENGINE = MergeTree()
ORDER BY event_time
TTL event_time + INTERVAL 7 DAY TO VOLUME 's3_cold',
    event_time + INTERVAL 30 DAY DELETE;    -- 7 天后转 S3,30 天后删除

-- 查看 TTL 任务状态
SELECT 
    table,
    name as partition,
    delete_ttl_info_min,
    delete_ttl_info_max,
    move_ttl_info.expression
FROM system.parts 
WHERE table = 'logs' AND active = 1;

7.4 索引粒度调优

index_granularity 控制稀疏索引的采样间隔,默认 8192 行。

场景推荐粒度理由
大宽表、低 Cardinality 过滤8192(默认)减少索引体积,提升扫描效率
高 Cardinality 点查(如 user_id 精确匹配)512 / 1024更精准定位,减少无效数据扫描
时序数据、范围扫描为主4096 / 8192范围查询以顺序读取为主
超大数据量(百亿级)16384降低索引内存占用
-- 建表时指定索引粒度
CREATE TABLE high_cardinality_events (
    event_id UUID,
    user_id UInt64,
    event_time DateTime
) ENGINE = MergeTree()
ORDER BY event_id
SETTINGS index_granularity = 512;           -- 更细粒度索引

-- 运行时查看表设置
SELECT * FROM system.tables WHERE name = 'orders' \G

8. 物化视图与 Projection

8.1 物化视图(Materialized View)异步刷新机制

物化视图是触发器式的。当源表 INSERT 时,数据自动转投到 MV 的目标表中,不占用实时查询时间。

-- 1. 创建目标表(存储聚合状态)
CREATE TABLE events_agg (
    dt Date,
    domain String,
    pv AggregateFunction(count, UInt64),
    uv AggregateFunction(uniqExact, UInt64),
    total_latency AggregateFunction(sum, UInt64)
) ENGINE = AggregatingMergeTree()
ORDER BY (dt, domain);

-- 2. 创建物化视图(自动触发)
CREATE MATERIALIZED VIEW events_agg_mv
TO events_agg
AS SELECT
    toDate(event_time) as dt,
    domain,
    countState() as pv,
    uniqExactState(user_id) as uv,
    sumState(latency_ms) as total_latency
FROM events
GROUP BY dt, domain;

-- 3. 查询物化视图(极速)
SELECT 
    dt, domain,
    countMerge(pv) as pv,
    uniqExactMerge(uv) as uv,
    sumMerge(total_latency) / countMerge(pv) as avg_latency
FROM events_agg
WHERE dt = today()
GROUP BY dt, domain;

注意事项

  • MV 只处理 INSERT,不处理 UPDATE/DELETE/ALTER
  • 源表历史数据不会自动回填 MV,需手动写入目标表
  • 多 MV 写入同一目标表时需避免数据冲突

8.2 Projection 查询自动路由

Projection 是表内建的预聚合数据结构,ClickHouse 查询优化器会自动选择最优 Projection 执行查询。

CREATE TABLE events_with_proj (
    event_time DateTime,
    user_id UInt64,
    domain String,
    page String,
    latency_ms UInt32,
    
    -- Projection 1: 按 domain + date 聚合 PV/UV
    PROJECTION proj_domain_daily
    (
        SELECT 
            toDate(event_time) as dt,
            domain,
            count() as pv,
            uniqExact(user_id) as uv,
            sum(latency_ms) as total_latency
        GROUP BY dt, domain
    ),
    
    -- Projection 2: 按 page 聚合访问情况
    PROJECTION proj_page_stats
    (
        SELECT 
            domain,
            page,
            count() as pv,
            avg(latency_ms) as avg_latency
        GROUP BY domain, page
    )
    
) ENGINE = MergeTree()
ORDER BY (domain, event_time);

-- 自动路由:以下查询会自动使用 proj_domain_daily
SELECT 
    toDate(event_time) as dt,
    domain,
    count() as pv,
    uniqExact(user_id) as uv
FROM events_with_proj
WHERE dt = today()
GROUP BY dt, domain;

-- 验证是否命中 Projection(查看执行计划)
EXPLAIN ACTIONS 
SELECT domain, count() FROM events_with_proj GROUP BY domain;
-- 输出中含 "ReadFromProjection" 即命中预聚合数据

MV vs Projection 对比

特性Materialized ViewProjection
存储位置独立目标表表内嵌(共享存储)
自动路由否(需显式查询目标表)是(优化器自动选择)
灵活性高(可多表 Join 后聚合)低(只能单表列)
维护成本中(需管理目标表结构)低(随表 DDL 自动管理)
适用场景复杂多源聚合、跨表计算单表多维预聚合、查询自动加速

9. 分布式表与副本

9.1 Distributed 引擎分片规则

-- 分布式表定义:基于 sharding_key 分发
CREATE TABLE orders_dist ON CLUSTER '{cluster}' AS orders_local
ENGINE = Distributed('{cluster}', 'default', 'orders_local', 
    cityHash64(user_id)        -- 分片键:保证同一 user_id 发到同一分片
);

-- 常见分片策略对比
-- 1. rand()           → 均匀分布,但不利于按维度聚合
-- 2. cityHash64(id)   → 按业务键哈希,利于本地 JOIN/聚合
-- 3. toYYYYMM(dt)     → 按时间分片,适合时序场景
-- 4. jumpConsistentHash(user_id, 8) → 一致性哈希,扩容友好

分片查询执行流程

客户端 → SELECT ... FROM orders_dist
            ↓
    ┌─────────────────┬─────────────────┐
    │  协调节点        │                 │
    │ 发送查询到各shard│                 │
    └────────┬────────┘                 │
             │                          │
    ┌────────▼────────┐        ┌───────▼────────┐
    │  Shard 1 (local)│        │  Shard 2 (remote)│
    │  本地聚合        │        │  本地聚合        │
    │  sum(amount)    │        │  sum(amount)    │
    └────────┬────────┘        └────────┬───────┘
             │                          │
             └──────────┬───────────────┘
                        ↓
                 协调节点合并结果
                        ↓
                 返回客户端
-- 查看分布式查询执行状态
SELECT 
    initial_query_id,
    host_name,
    type,
    event_time,
    query_duration_ms,
    read_rows,
    read_bytes
FROM clusterAllReplicas('{cluster}', system.query_log)
WHERE initial_query_id = 'xxx'
ORDER BY event_time;

9.2 ReplicatedMergeTree 副本同步

-- 本地副本表(每个节点独立运行)
CREATE TABLE orders_local ON CLUSTER '{cluster}' (
    order_id UInt64,
    user_id UInt64,
    amount Decimal(18,2),
    create_time DateTime
) ENGINE = ReplicatedMergeTree(
    '/clickhouse/{cluster}/tables/{shard}/orders',   -- ZooKeeper 路径
    '{replica}'                                        -- 副本标识
)
PARTITION BY toYYYYMM(create_time)
ORDER BY (user_id, create_time);

ZooKeeper / ClickHouse Keeper 协调机制

写入流程:
INSERT INTO orders_local (Shard 1, Replica A)
        ↓
 Replica A 写入本地 Part
        ↓
 通知 ZooKeeper(/clickhouse/.../orders/log)
        ↓
 Replica B 监听 log 变更 → 拉取 Part → 验证 checksum → 激活 Part
        ↓
 副本 A/B 数据最终一致

Keeper 路径结构

/clickhouse/{cluster}/tables/{shard}/orders/
├── replicas/
│   ├── replica_01/           ← 各副本注册自身信息
│   │   ├── is_active
│   │   ├── host
│   │   ├── log_pointer       ← 当前同步到的 log 位置
│   │   └── parts/
│   └── replica_02/
├── queue/                     ← 待执行的复制任务队列
├── log/                       ← 全局操作日志(INSERT/ALTER/MERGE)
├── leader_election/           ← 主副本选举
└── columns                    ← 表结构元数据
-- 查看副本同步延迟
SELECT 
    database, table,
    is_leader,
    can_become_leader,
    is_readonly,
    future_parts,
    parts_to_check,
    zookeeper_path,
    replica_name,
    queue_size,                         -- 待处理队列大小
    inserts_in_queue,                   -- 待插入 Part 数
    merges_in_queue,                    -- 待合并任务数
    absolute_delay,                     -- 绝对延迟(秒)
    total_replicas,
    active_replicas
FROM system.replicas 
WHERE table = 'orders_local';

-- 查看 ZooKeeper 操作统计
SELECT * FROM system.zookeeper WHERE path = '/clickhouse';

10. 查询优化实战

10.1 PREWHERE 优化

ClickHouse 默认启用 optimize_move_to_prewhere,将 WHERE 条件中能高效过滤的列推到 PREWHERE 阶段,先过滤再读取其他列。

-- 原始查询
SELECT user_id, amount, city 
FROM orders 
WHERE city = 'Beijing' AND amount > 1000;

-- 优化器自动重写为:
SELECT user_id, amount, city
FROM orders
PREWHERE city = 'Beijing'       -- 先只用 city 列过滤,减少后续数据量
WHERE amount > 1000;

-- 手动控制 PREWHERE(当优化器判断不准时)
SELECT user_id, amount
FROM orders
PREWHERE create_time > '2024-01-01'
WHERE status = 1;

10.2 向量化执行与 SIMD 加速

-- 确认查询是否使用向量化执行(查看 PIPELINE)
EXPLAIN PIPELINE 
SELECT sum(amount), avg(latency_ms) 
FROM events 
WHERE domain = 'api.example.com';

-- 期望输出包含以下算子(表示向量化路径):
-- ExpressionTransform
-- FilterTransform
-- AggregatingTransform

SIMD 加速条件

  • 数据类型为定宽类型(UInt8/16/32/64, Float32/64)
  • 过滤条件为简单比较(=, <, >, BETWEEN)
  • 聚合函数为 sum/count/avg/min/max 等
-- 使用 Int32 而非 String 编码状态,触发 SIMD
SELECT count() FROM events WHERE status_code > 400;
-- 若 status_code 为 UInt16,ClickHouse 会用 SSE/AVX2 批量比较

10.3 查询 Profile 分析与慢查询排查

-- 开启 Profile 日志(users.xml)
<trace_log>
    <database>system</database>
    <table>trace_log</table>
</trace_log>

-- 查看最近慢查询 Top 20
SELECT 
    event_time,
    query_id,
    query,
    query_duration_ms,
    read_rows,
    read_bytes,
    result_rows,
    memory_usage,
    Settings['max_threads'] as threads
FROM system.query_log
WHERE event_time > now() - INTERVAL 1 HOUR
  AND type = 'QueryFinish'
ORDER BY query_duration_ms DESC
LIMIT 20;

-- 分析单次查询详细 Stage
SELECT 
    event_name,
    duration_ms,
    read_rows,
    read_bytes
FROM system.query_execution_log        -- ClickHouse 24.3+
WHERE query_id = 'xxx'
ORDER BY event_time;

慢查询常见原因与排查

症状根因排查命令优化方案
扫描行数远大于结果行数未命中主键/分区裁剪EXPLAIN indexes调整 ORDER BY / PARTITION BY
内存溢出大聚合 / 大 JOINsystem.query_log.memory_usageGROUP BY 分桶、限制 max_memory_usage
高 CPU 低 IO复杂正则 / 函数计算Trace 日志预计算、物化视图
高 IO 低 CPU读取列过多 / 压缩率差system.parts.bytes_on_disk减少 SELECT *、更换压缩算法
副本查询延迟分布式查询等待慢副本system.replicas.absolute_delay读写分离、prefer_localhost_replica
-- 使用 EXPLAIN 分析索引命中情况
EXPLAIN indexes = 1 
SELECT * FROM orders WHERE order_id = 12345;

-- 若 order_id 不在 ORDER BY 最左前缀,输出将显示
-- "Condition(order_id = 12345) is not analyzed" 或无索引信息

11. 外部集成

11.1 Kafka Engine 实时摄入

-- 创建 Kafka 消费表(只做消费端接入)
CREATE TABLE events_kafka (
    event_time DateTime,
    user_id UInt64,
    domain String,
    page String,
    latency_ms UInt32
) ENGINE = Kafka()
SETTINGS 
    kafka_broker_list = 'kafka-1:9092,kafka-2:9092',
    kafka_topic_list = 'clickhouse-events',
    kafka_group_name = 'ch-events-consumer',
    kafka_format = 'JSONEachRow',
    kafka_num_consumers = 4,            -- 并发消费者数
    kafka_max_block_size = 1048576,     -- 每批次最大行数
    kafka_skip_broken_messages = 10;    -- 跳过错误消息阈值

-- 创建目标 MergeTree 表
CREATE TABLE events (
    event_time DateTime,
    user_id UInt64,
    domain String,
    page String,
    latency_ms UInt32
) ENGINE = MergeTree()
ORDER BY (domain, event_time);

-- 创建物化视图桥接 Kafka 表 → 目标表
CREATE MATERIALIZED VIEW events_consumer TO events
AS SELECT * FROM events_kafka;

Kafka Engine 消费架构

Kafka Topic (clickhouse-events)
        ↓
    ┌───┴───┬───┬───┐
    │ C0    │C1 │C2 │C3    ← kafka_num_consumers 并行消费
    └───┬───┴───┴───┘
        ↓
  ClickHouse Buffer / Direct Insert
        ↓
  MergeTree 表(异步 merge 优化存储)

11.2 MySQL / PostgreSQL 数据库引擎联邦查询

-- MySQL 引擎:直接查询远端 MySQL 表
CREATE TABLE mysql_users (
    id UInt64,
    name String,
    email String
) ENGINE = MySQL('mysql-host:3306', 'mydb', 'users', 'reader', 'password');

-- 直接查询(数据不落地 ClickHouse)
SELECT * FROM mysql_users WHERE id > 1000;

-- PostgreSQL 引擎
CREATE TABLE pg_orders (
    order_id UInt64,
    amount Decimal(18,2)
) ENGINE = PostgreSQL('pg-host:5432', 'shop', 'orders', 'reader', 'password');

-- 联邦 JOIN:ClickHouse 拉取远端小表做本地 JOIN
SELECT 
    o.order_id,
    u.name,
    o.amount
FROM pg_orders o
JOIN mysql_users u ON o.user_id = u.id
WHERE o.amount > 100;

11.3 S3 表函数与冷存

-- 直接查询 S3 上的 Parquet 文件(无需预先加载)
SELECT 
    user_id,
    count() as cnt,
    sum(amount) as total
FROM s3(
    'https://bucket.s3.amazonaws.com/data/orders/*.parquet',
    'AKIA...',                          -- Access Key
    'secret...',                        -- Secret Key
    'Parquet'                           -- 文件格式
)
WHERE create_time > '2024-01-01'
GROUP BY user_id;

-- 插入数据到 S3(导出冷存)
INSERT INTO FUNCTION s3(
    'https://bucket.s3.amazonaws.com/archive/orders_202401.parquet',
    'AKIA...', 
    'secret...',
    'Parquet'
)
SELECT * FROM orders WHERE create_time < '2024-02-01';

-- S3 作为外部卷挂载到表 TTL
CREATE TABLE logs (
    event_time DateTime,
    message String
) ENGINE = MergeTree()
ORDER BY event_time
TTL event_time + INTERVAL 30 DAY TO DISK 's3_disk';

12. 集群部署架构

12.1 分片副本矩阵

生产环境推荐最少 2 Shard × 2 Replica(共 4 节点),兼顾性能与高可用。

                    ┌─────────────────────────────────────────┐
                    │           ClickHouse Cluster              │
                    │              "prod_cluster"               │
                    └─────────────────────────────────────────┘
                                      │
              ┌───────────────────────┼───────────────────────┐
              │ Shard 1 (1/2 数据)    │      │ Shard 2 (1/2 数据)
        ┌─────┴──────┐            ┌───┴────┐            ┌────┴──────┐
        │  Replica   │            │ Replica│            │  Replica  │
        │   01       │◄──────────►│   02   │            │   01      │
        │ (Primary)  │   互为副本  │(Backup)│            │ (Primary) │
        │ 192.168.1.1│            │192.168.│            │ 192.168.  │
        │  :9000     │            │1.2:9000│            │ 1.3:9000  │
        └─────┬──────┘            └────────┘            └────┬──────┘
              │                                              │
              └──────────────────┬───────────────────────────┘
                                 │
                        ┌────────▼─────────┐
                        │   Distributed    │
                        │    分布式表       │
                        │  路由 + 结果合并   │
                        └──────────────────┘
<!-- remote_servers.xml 配置示例 -->
<remote_servers>
    <prod_cluster>
        <shard>
            <internal_replication>true</internal_replication>
            <replica>
                <host>ch-shard1-replica1</host>
                <port>9000</port>
            </replica>
            <replica>
                <host>ch-shard1-replica2</host>
                <port>9000</port>
            </replica>
        </shard>
        <shard>
            <internal_replication>true</internal_replication>
            <replica>
                <host>ch-shard2-replica1</host>
                <port>9000</port>
            </replica>
            <replica>
                <host>ch-shard2-replica2</host>
                <port>9000</port>
            </replica>
        </shard>
    </prod_cluster>
</remote_servers>

12.2 读写分离架构

写入层 ──► 直接写 Local ReplicatedMergeTree(各节点)
             ↑               ↑               ↑
        Application      Flink/Spark     Kafka Engine
             
查询层 ──► ch-proxy / clickhouse-operator 负载均衡
             │
        ┌────┴────┬────────┬────────┐
        │ Replica 1 │ Replica 2 │ Replica 3 │  ← 只读查询分发
        └─────────┴────────┴────────┘

ch-proxy 配置示例

# ch-proxy.yml
server:
  http:
    listen_addr: ":9090"
  
clusters:
  - name: "prod_cluster"
    scheme: "http"
    nodes:
      - "192.168.1.1:8123"
      - "192.168.1.2:8123"
      - "192.168.1.3:8123"
      - "192.168.1.4:8123"
    users:
      - name: "reader"
        password: "xxx"
        max_concurrent_queries: 50
        max_execution_time: 120s
    # 负载均衡策略
    kill_query_user:
      name: "admin"
      password: "xxx"
    heartbeat:
      interval: 5s
      timeout: 3s
      request: "/ping"

clickhouse-operator(Kubernetes)部署

apiVersion: clickhouse.altinity.com/v1
kind: ClickHouseInstallation
metadata:
  name: prod-cluster
spec:
  configuration:
    clusters:
      - name: prod
        layout:
          shardsCount: 2
          replicasCount: 2
    zookeeper:
      nodes:
        - host: zookeeper-0.zk
          port: 2181
  templates:
    podTemplates:
      - name: clickhouse-pod
        spec:
          containers:
            - name: clickhouse
              image: clickhouse/clickhouse-server:24.3
              resources:
                requests:
                  memory: "8Gi"
                  cpu: "4"
                limits:
                  memory: "32Gi"
                  cpu: "16"
    volumeClaimTemplates:
      - name: data
        spec:
          accessModes:
            - ReadWriteOnce
          resources:
            requests:
              storage: 500Gi

13. 运维实战

13.1 磁盘监控

-- 查看各数据库/表磁盘占用
SELECT 
    database,
    table,
    formatReadableSize(sum(bytes_on_disk)) as disk_size,
    sum(rows) as total_rows,
    count() as parts_count,
    max(modification_time) as last_modified
FROM system.parts
WHERE active = 1
GROUP BY database, table
ORDER BY sum(bytes_on_disk) DESC
LIMIT 20;

-- 查看各磁盘卷使用情况
SELECT 
    name,
    path,
    formatReadableSize(free_space) as free,
    formatReadableSize(total_space) as total,
    round((1 - free_space / total_space) * 100, 2) as usage_pct
FROM system.disks;

-- 监控 Merge 导致的临时磁盘膨胀
SELECT 
    database, table,
    round(100 * sum(bytes_on_disk * (1 - active)) / sum(bytes_on_disk), 2) as inactive_pct
FROM system.parts
GROUP BY database, table
HAVING inactive_pct > 20;               -- 活跃数据占比过低,说明旧 Part 堆积

13.2 BACKUP / RESTORE

-- ClickHouse 24.3+ 原生备份(推荐)
BACKUP TABLE orders TO File('/backups/orders_20240101.zip');

-- 备份整个数据库
BACKUP DATABASE default TO S3('https://bucket.s3.amazonaws.com/ch-backup/default_20240101', 'AK', 'SK');

-- 恢复表
RESTORE TABLE orders FROM File('/backups/orders_20240101.zip');

-- 使用 clickhouse-backup 工具(生产推荐)
# 1. 创建备份
clickhouse-backup create orders_backup_20240101

# 2. 上传到远程存储
clickhouse-backup upload orders_backup_20240101 --storage=s3

# 3. 恢复(先停写入)
clickhouse-backup restore orders_backup_20240101 --table=orders

13.3 升级策略

# 滚动升级流程(副本保证可用性)
# 1. 标记副本为只读(避免写入到待升级节点)
echo "SYSTEM STOP REPLICATED SENDS" | clickhouse-client

# 2. 检查副本同步延迟(确保 <= 10s)
SELECT absolute_delay FROM system.replicas WHERE replica_name = 'replica_01';

# 3. 停止 ClickHouse 服务
systemctl stop clickhouse-server

# 4. 替换二进制文件(保留 config)
apt install clickhouse-server=24.8.1 clickhouse-client=24.8.1

# 5. 启动并验证
systemctl start clickhouse-server
clickhouse-client --query "SELECT version()"

# 6. 恢复副本同步
echo "SYSTEM START REPLICATED SENDS" | clickhouse-client

# 7. 逐节点重复(确保每轮至少一个副本可用)

13.4 常见故障排查

-- 故障 1:副本同步中断(queue_size 持续增长)
-- 解决:检查 ZooKeeper 连接、磁盘空间、手动恢复
SYSTEM RESTART REPLICA orders_local;

-- 故障 2:Too many parts(写入过于频繁,小文件堆积)
-- 解决:增大 batch_size、启用 Buffer 表、调大 parts_to_delay_insert
SELECT 
    table,
    count() as part_count
FROM system.parts
WHERE active = 1
GROUP BY table
HAVING part_count > 300;

-- 故障 3:查询卡死(max_execution_time 不生效)
-- 根因:某些算子(如复杂 JOIN)不响应取消信号
-- 解决:Kill 查询、限制 Join 表大小、改用 GLOBAL JOIN
KILL QUERY WHERE query_id = 'xxx';

-- 故障 4:内存溢出(OOM Killer)
-- 解决:调低 max_memory_usage、开启内存溢出转磁盘(allow_experimental_memory_bound_aggregator)
SET max_memory_usage = 10_000_000_000;      -- 10GB 上限
SET max_bytes_before_external_group_by = 5_000_000_000;
SET max_bytes_before_external_sort = 5_000_000_000;

14. FAQ

Q1: ClickHouse 是否支持 UPDATE 和 DELETE?

支持,但非传统 OLTP 语义。UPDATEDELETE 通过异步的 Mutation 实现:

  • ALTER TABLE ... DELETE WHERE ... 生成 mutation 任务,后台重写 affected parts
  • 执行期间,旧数据仍可读;mutation 完成后旧 part 被替换
  • 大规模 mutation 开销大,推荐用 ReplacingMergeTree / CollapsingMergeTree 替代频繁更新

Q2: 为什么 COUNT(DISTINCT) 比 uniqExact 慢,两者有何区别?

COUNT(DISTINCT) 是标准 SQL,ClickHouse 内部映射为 uniqExact。但直接写 uniqExact(col) 可配合 -StateAggregatingMergeTree 中预计算。

  • uniq(col):近似去重(HyperLogLog++),误差 < 1%,内存占用小
  • uniqExact(col):精确去重,内存消耗大,大数据量建议用物化视图预聚合

Q3: 分布式表查询为什么比本地表慢很多?

常见原因:

  1. 网络开销:协调节点需向所有 shard 广播查询、接收结果、二次聚合
  2. 单点瓶颈rand() 分片导致同一维度数据分散,GROUP BY 无法本地完成
  3. 副本不均衡:某 shard 数据量/查询负载远高于其他
  4. IN/JOIN 下推问题:复杂 IN 子句未拆分,导致全量数据传输

优化方案:改用业务键哈希分片、使用 GLOBAL IN、添加 prefer_localhost_replica=1

Q4: 如何选择 index_granularity?

默认 8192 行适合大部分场景。若查询模式为高 Cardinality 点查(如 UUID 精确匹配),调至 5121024 可减少无效扫描 90%;若以时序范围扫描为主,40968192 更平衡(太细会增加索引内存)。百亿级大表可尝试 16384 降低索引开销。

Q5: ClickHouse 能否完全替代 Elasticsearch 做日志分析?

视场景而定:

  • 结构化日志 + 聚合分析:ClickHouse 优势明显(存储成本低 510 倍、聚合快 10100 倍)
  • 全文检索 + relevance scoring:ES 更优,ClickHouse 仅支持 LIKE / hasToken / multiSearchAny 等基础文本匹配
  • 混合方案:日志写入 Kafka → 结构化字段入 ClickHouse、原始文本入 ES,由业务需求决定查询路由

15. 总结

ClickHouse 凭借列式存储、稀疏索引、向量化执行和丰富的 MergeTree 引擎家族,在 OLAP 场景建立了显著的性价比优势。本文从引擎选型到底层存储、从单机优化到集群部署进行了系统梳理,核心要点如下:

决策维度推荐方案
通用分析表(追加型日志)MergeTree + 主键有序 + 按月分区
幂等写入 / 去重查询ReplacingMergeTree + FINAL / 或物化视图
数值指标预聚合SummingMergeTree(简单求和)或 AggregatingMergeTree(复杂聚合)
状态变更 / 需要更新语义CollapsingMergeTree(Sign 折叠)
查询加速自动路由Projection(表内建)优先,复杂场景用 Materialized View
集群高可用ReplicatedMergeTree 副本 + Distributed 分布式表
实时数据摄入Kafka Engine + MV 转存 MergeTree
联邦查询MySQL / PostgreSQL Engine 或 S3 表函数
查询性能瓶颈排查EXPLAIN + system.query_log + Profile 分析
生产运维保障ch-proxy 负载均衡、磁盘监控、滚动升级、BACKUP/RESTORE

在生产环境中,ClickHouse 的强项是大批量写入 + 聚合分析,弱项是高频单条更新 + 复杂事务。合理选择引擎、设计分区与主键、利用物化视图预计算,是发挥 ClickHouse 极致性能的关键。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「data-engineering」更多文章

  1. 数据工程深度指南:Modern Data Stack 全栈实践
  2. 数据平台工程:Data Mesh、FinOps 与 DataOps 生产实践
  3. Kafka Connect CDC 实战:Debezium 数据同步与变更捕获