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_filter | Bloom 过滤器判断存在性 | 等值查询(字符串/UUID) |
| tokenbf_v1 | 支持分词的 Bloom Filter | 全文检索 |
| ngrambf_v1 | N-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 + 1Mutation:可选的 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 View | Projection |
|---|---|---|
| 存储位置 | 独立目标表 | 表内嵌(共享存储) |
| 自动路由 | 否(需显式查询目标表) | 是(优化器自动选择) |
| 灵活性 | 高(可多表 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 |
| 内存溢出 | 大聚合 / 大 JOIN | system.query_log.memory_usage | 加 GROUP 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 语义。
UPDATE和DELETE通过异步的 Mutation 实现:
ALTER TABLE ... DELETE WHERE ...生成 mutation 任务,后台重写 affected parts- 执行期间,旧数据仍可读;mutation 完成后旧 part 被替换
- 大规模 mutation 开销大,推荐用
ReplacingMergeTree/CollapsingMergeTree替代频繁更新
Q2: 为什么 COUNT(DISTINCT) 比 uniqExact 慢,两者有何区别?
COUNT(DISTINCT)是标准 SQL,ClickHouse 内部映射为uniqExact。但直接写uniqExact(col)可配合-State在AggregatingMergeTree中预计算。
uniq(col):近似去重(HyperLogLog++),误差 < 1%,内存占用小uniqExact(col):精确去重,内存消耗大,大数据量建议用物化视图预聚合
Q3: 分布式表查询为什么比本地表慢很多?
常见原因:
- 网络开销:协调节点需向所有 shard 广播查询、接收结果、二次聚合
- 单点瓶颈:
rand()分片导致同一维度数据分散,GROUP BY 无法本地完成- 副本不均衡:某 shard 数据量/查询负载远高于其他
- IN/JOIN 下推问题:复杂 IN 子句未拆分,导致全量数据传输
优化方案:改用业务键哈希分片、使用
GLOBAL IN、添加prefer_localhost_replica=1。
Q4: 如何选择 index_granularity?
默认 8192 行适合大部分场景。若查询模式为高 Cardinality 点查(如 UUID 精确匹配),调至 512
1024 可减少无效扫描 90%;若以时序范围扫描为主,40968192 更平衡(太细会增加索引内存)。百亿级大表可尝试 16384 降低索引开销。
Q5: ClickHouse 能否完全替代 Elasticsearch 做日志分析?
视场景而定:
- 结构化日志 + 聚合分析:ClickHouse 优势明显(存储成本低 5
10 倍、聚合快 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 极致性能的关键。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。