
做ClickHouse运维和开发这些年最常被问的问题不是“怎么安装”而是“为什么我上了集群查询反而变慢了”。明明节点一堆数据也放进去了分布式查询一跑就是几秒甚至几十秒排查到最后十有八九问题出在数据分片策略上。ClickHouse的分片设计直接影响分布式查询能不能命中更少的节点、能不能减少网络传输而这恰恰是很多人刚接触集群时最容易忽略的部分。这篇文章就围绕数据分片怎么定、分布式表怎么查、集群怎么配、坑怎么踩把ClickHouse分布式查询性能这摊事讲透。内容适合正在规划ClickHouse集群的架构师、被线上慢查询折磨的DBA也适合想系统理解分布式表原理的开发者。全文基于我实际运维和压测过的生产环境经验给出的配置和排查思路都可以直接参考。1. 分片策略设计性能在写入那一刻就决定了1.1 分片键选择是分布式查询性能的分水岭很多人以为ClickHouse集群就是把数据均匀撒到每台机器上查询时全部节点一起算肯定比单机快。这个想法对了一半分布式查询确实会并行下发到所有分片但“一起算”不代表“快”。查询能不能裁剪到更少的分片取决于分片键和查询条件的匹配程度。ClickHouse分片的原理很简单写入时按分片键的哈希值路由到某个分片查询时分布式表会把SQL发到每个分片各分片算完一部分结果再汇总起来。如果分片键选得好查询条件里带了分片键那么查询就可以只路由到少数几个分片其他节点完全不参与速度自然快。如果分片键跟查询条件对不上那就只能全量广播到所有分片每个节点都扫一遍自己的数据中间结果还要通过网络汇总性能必然暴跌。我见过一个典型案例一张用户行为表按device_id做了哈希分片线上查询却几乎都按app_version过滤结果每次查询都广播到全集群20个分片单次查询扫描的数据量是实际需要的十几倍。后来重建表的时候把分片键换成了app_version相关的字段虽然单个分片内部的数据量不均衡了但查询裁剪效果极其明显P99延迟从3秒降到200毫秒。所以选择分片键的第一原则尽量选择高频查询条件里的等值过滤字段让查询能够命中尽量少的分片。常见的可选方案有三种分片键方式数据分布查询裁剪能力适用场景rand()最均匀冷热均衡几乎无法裁剪日志、事件流水等全量分析场景业务键哈希均匀可控按该键过滤时能裁剪用户ID、订单ID等强业务标识低基数字段直接分片极不均匀按该字段过滤时命中精准租户、地域等隔离查询场景有一类典型陷阱是用低基数字段直接做分片键比如按省份、按机房分片。表面上看每个分片对应一个省份查询单个省份时很爽但只要有一个省份的数据量是别的省份的十倍这个节点就会成为热节点其他节点闲着它自己被拖垮。低基数字段更适合做WHERE里的过滤条件而不是分片键。如果业务查询确实没有固定维度比如纯日志分析那用rand()也没问题至少数据分布均匀。但这时候就要接受所有查询都是全分片扫描集群的优势主要体现在CPU和磁盘并行上而不是智能裁剪上。1.2 分片数量与副本数量怎么定分片数量决定集群的扩展粒度和查询并行度。在ClickHouse里一个分片可以理解为一份数据的一个物理副本集合分片之间数据不重叠。分片数量通常由数据总量和单机处理能力决定。经验值上单个副本的服务器建议配置不低于16核CPU、64GB内存数据量在单机1TB到5TB之间比较舒服。如果总数据量是20TB那8个分片左右是底线加上磁盘冗余和查询并发10到12个分片更稳妥。分片太少会导致单节点压力大分片太多则协调节点合并结果的开销变大小查询也可能被放大成大规模并行任务。副本数量的选择也很关键。ClickHouse的副本默认依赖ZooKeeper或ClickHouse Keeper进行元数据同步副本不是越多越好。生产环境我一般建议一个分片至少2个副本保证节点宕机时查询不受影响。副本数量超过3个后收益会明显递减反而增加了ZooKeeper的通信压力。如果一个分片的副本间数据完全一样写入时需要注意internal_replication这个配置它设置为true时分布式表只向该分片的一个副本写入由集群内部同步到其他副本设置为false时分布式表会向该分片的所有副本各自写入适合底层存储不支持自动复制的场景。从性能角度说查询时分布式表会优先选择一个本地副本prefer_localhost_replica默认开启这样当前节点如果正好持有数据就不用走网络。对多副本集群这个默认行为能省不少内网带宽。1.3 分区和分片不是一回事别再搞混这是新手最容易混淆的两个概念。分片解决的是“数据存在哪台机器上”的问题分区解决的是“一块数据在单台机器内怎么组织”的问题。分片是集群维度分区是表内维度。比如说一张订单表按order_id哈希分片到3台机器每台机器上的本地表又按toYYYYMM(create_time)做了月度分区。查询的时候分布式表先通过分片键裁剪到命中节点节点内部再通过分区裁剪跳过无关月份的数据。两套机制是叠加生效的缺一不可。很多慢查询优化的第一刀不是加索引而是核对分片和分区是否都做到了裁剪。著名的经验法则是如果查询扫描的分区数量从12个降到1个往往比加索引更有效。分区字段一般选择时间天、月、年配合TTL还能自动清理过期数据这属于另一块大话题但和分片策略是配合使用的。2. 分布式表查询机制与性能要点2.1 Distributed表引擎的写入与查询路径分布式表引擎为Distributed本身不存储任何数据它只是一个路由层。写入时分布式表根据分片键计算出目标分片把数据异步转发给对应节点上的本地表。查询时分布式表把SQL改写后下发到所有需要参与的分片各分片的本地表执行完局部查询再把结果返回给发起节点由发起节点做最终聚合、排序、LIMIT等操作。这套机制决定了分布式查询的性能瓶颈通常不在计算而在两个地方一是网络传输量二是发起节点的合并开销。比如一个大查询每个分片返回几百万行中间结果汇聚节点要先把这些结果全部收进来再排序内存压力极大。理解这个路径之后很多优化手段就顺理成章了。首先尽量让WHERE条件能够下推到分片内部执行不要等所有数据汇到协调节点再做过滤其次LIMIT条件也会下推分片内先各自截断可以大幅减少传输量再次ORDER BY同样会下推各分片局部排序后再合并协调节点的压力会小很多。还有一点容易被忽略分布式表查询每个分片使用的连接数默认是max_threads如果查询并发高每个节点打开的连接数会非常多。适当控制查询并发和分片数比盲目扩容分片更实际。2.2 GLOBAL JOIN为什么能救命又能害人分布式表上做JOIN是ClickHouse的经典痛点。普通JOIN在分布式表上的行为是每个分片执行查询时发现自己需要关联的数据在别的分片就会去远端拉取结果产生近似于NMFN倍网络放大的效果。分片越多放大越严重很多分布式慢查询就是这么来的。GLOBAL JOIN的原理是先把右表的数据全部拉到协调节点做成一个全局哈希表再把这个哈希表分发到每个参与查询的分片。这样每个分片只需要一次本地哈希查找网络传输次数被限制住了。GLOBAL JOIN不是银弹。它要求右表数据量足够小能放进内存哈希表。如果右表是几亿行的分布式大表做GLOBAL JOIN会把协调节点直接打爆。更合理的做法是先把右表在本地节点物化成一张聚合后的紧凑表再用于关联。或者使用GLOBAL IN加子查询配合子查询内的聚合减少数据传输量。实际经验在千万级用户维表关联场景下GLOBAL JOIN配合右表按分片键预先分布能获得接近单机的性能但如果右表数据量过亿就必须考虑改用字典表或者预聚合纯SQL层面的JOIN优化空间已经不大了。2.3 分布式查询的几个关键参数参数作用我的建议prefer_localhost_replica查询优先命中本机副本默认开启不要关optimize_skip_unused_shards对分片键做等值过滤时跳过无关分片建议开启前提是分片键与查询条件匹配load_balancing多副本时选择查询副本的策略random或nearest_hostname均可max_threads单查询并发线程数根据CPU核数调整不要盲目拉高max_memory_usage单查询最大内存分布式聚合场景务必设置合理上限optimize_skip_unused_shards这个参数很多人不知道。它开启后如果查询的WHERE条件里包含了分片键的等值表达式协调节点可以在下发查询前就判断出只需要访问哪几个分片其他分片完全不唤醒。对分片数量多的集群这个优化效果非常明显。还有个小技巧如果集群里所有节点查询都会命中本地副本可以开启allow_experimental_parallel_reading_from_replicas新版中已经改名或成为默认让一个查询同时并发读取同一个分片的多个副本每个副本只读一部分数据这个特性在单分片性能吃紧时很管用。3. 从零搭一个三节点ClickHouse集群并验证分片效果3.1 配置文件与网络规划纸上谈兵说再多不如跑一个最小集群。我用三个节点ch01、ch02、ch03搭建一个3分片集群每节点只部署一个分片的一个副本总分片3总副本3。这套方案足够验证分片策略对查询性能的影响也方便后续扩展。ClickHouse的集群定义在config.xml的remote_servers段生产环境一般单独抽成metrika.xml引入。关键配置如下clickhouse remote_servers clickhouse_cluster shard replica hostch01/host port9000/port /replica /shard shard replica hostch02/host port9000/port /replica /shard shard replica hostch03/host port9000/port /replica /shard /clickhouse_cluster /remote_servers macros shard01/shard replicach01/replica /macros /clickhousemacros段非常重要它让分布式DDL可以在所有节点上自动替换成对应的本地宏变量建表、删表就不需要逐节点手工执行了。三个节点的macros分别写成01/ch01、02/ch02、03/ch03。集群建好后用一条SQL确认连通性SELECT hostName() AS node, count() FROM cluster(clickhouse_cluster, system.one) GROUP BY node;能返回三条记录说明集群网络层已经打通。3.2 建表与分片效果验证建一张分布式表和对应的本地表。本地表负责实际存储分布式表用来路由。两条DDL都通过ON CLUSTER下发保证集群内所有节点执行一致。-- 本地表 CREATE TABLE user_events_local ON CLUSTER clickhouse_cluster ( user_id UInt64, event_time DateTime, event_type String, amount Decimal(18,2) ) ENGINE MergeTree() PARTITION BY toYYYYMM(event_time) ORDER BY (user_id, event_time); -- 分布式表 CREATE TABLE user_events_dist ON CLUSTER clickhouse_cluster AS user_events_local ENGINE Distributed( clickhouse_cluster, default, user_events_local, cityHash64(user_id) );写入测试数据时分布式表会自动根据cityHash64(user_id)分流到对应分片。如果想验证数据分布是否均匀可以用下面的SQL查看每个节点各自处理的行数SELECT hostName() AS node, count() AS rows FROM user_events_dist GROUP BY node ORDER BY node;由于user_id本身是一个分布均匀的字段三个节点的行数应该非常接近。再看另一种情况如果把分片键换成event_type这种基数很低的字段三者行数就可能出现严重倾斜。实际测试时我曾经往一张按月份字符串分片的表里灌数据结果70%的数据全落在同一个节点上分布式集群活生生用成了单机。建表之后建议顺手验证一下optimize_skip_unused_shards的效果把三层配置文件里optimize_skip_unused_shards1/optimize_skip_unused_shards打开然后分别执行不带分片键过滤和带分片键过滤的同一条聚合查询观察system.query_log里的read_rows和read_bytes。带分片键过滤的查询扫描行数应该会明显下降这验证了分片裁剪在实际查询中确实生效。3.3 性能验证思路单节点 vs 分布式我在验证分片效果时通常会在同一份数据上做三组对比本地表查询、分布式表全扫描、分布式表按分片键裁剪。对比指标只看两个扫描行数read_rows和总耗时。一组典型的对比结果是这样的10亿行数据3节点集群查询方式扫描行数耗时单节点本地表全表聚合10亿8.2s分布式表全分片聚合10亿2.9s分布式表按分片键过滤聚合约333万0.4s分布式表全分片聚合只比单节点快约三倍这符合三节点并行的预期。按分片键过滤后由于裁剪掉两个分片扫描量下降两个数量级耗时自然大幅缩短。这说明分片键和查询条件的匹配程度对性能影响远比“集群多不多”更大。如果手里没有现成的ClickHouse环境可以用以下方式模拟在本地写一个简单的Python脚本插入100万行user_id均匀的测试数据然后用三种不同的分片键分别建表对比。实践下来这个过程中的体验比看任何文档都深刻。4. 典型故障排查与部署经验4.1 重启报错“failed to flush system log already exists”怎么办这个话题在社区里搜的人非常多几乎每个经历过ClickHouse异常重启的人都见过类似日志。报错原因是ClickHouse重启时需要把系统日志表system.query_log、system.query_thread_log等flush到磁盘但元数据发现某些日志表已经存在发生冲突导致启动流程中断。常见诱因有三个异常kill导致日志表数据文件不完整、多副本节点同时操作了系统表、或者手动改过system库的表结构。处理思路分为两步。第一步尝试通过SQL清理日志表前提是实例能起来SYSTEM FLUSH LOGS; TRUNCATE TABLE system.query_log; TRUNCATE TABLE system.query_thread_log; TRUNCATE TABLE system.trace_log;如果实例已经起不来可以先用临时配置禁用系统日志表把query_log、trace_log、query_thread_log等配置项临时置为removetrue/remove或用空配置覆盖让实例跳过系统表的初始化直接启动然后清理旧日志表再恢复原配置重启。更稳妥的做法是在config.xml里对系统日志表加上engineEngine MergeTree/engine或使用Buffer时注意刷盘策略减少异常重启时日志表文件损坏的概率。关键教训是不要随意DROPsystem库下的表这些表和系统表引擎强关联处理不当会导致更多元数据不一致问题。还有一个容易被忽略的点如果集群里多副本共用同一个system.path也就是数据目录被多个实例共享也会触发类似报错。确认每个节点数据目录独立是避免这类问题的基础前提。4.2 Windows、银河麒麟等环境下的部署体验说到部署很多人会在Windows上想当然地跑ClickHouse服务端。官方长期并不提供Windows原生服务端支持Windows下常见的做法是WSL2或Docker。WSL2跑ClickHouse做学习验证是没问题的但性能损失明显而且Windows文件系统和Linux文件系统之间的IO语义差异可能导致意外损坏生产环境千万不要这么干。银河麒麟这类国产系统上部署ClickHouse思路和CentOS/RHEL类似核心是离线安装包。ClickHouse官方提供tgz离线包解压后配置好环境变量和systemd服务即可。有几个细节需要注意下载包时选择与操作系统架构匹配的版本x86和ARM的二进制不能混用离线安装时提前准备好libicu等动态库依赖另外单机二进制方式启动时clickhouse-server会自动探测CPU指令集个别老CPU上启动会报非法指令解决办法是换用兼容版本或在编译参数上调整。银河麒麟上部署我还踩过一个包管理器源冲突的坑卸载旧版本时残留了/etc/clickhouse-server下的配置目录导致新版本起不来。卸载后务必确认配置目录和/var/lib/clickhouse数据目录完全清理干净再重装。4.3 数据倾斜的发现与缓解分片数据倾斜是分布式集群的头号性能杀手。倾斜的发现办法很简单用分布式表按节点分组统计行数、字节数、分区数对比差异。也可以在system.parts表上按hostName()和partition_id聚合看得更细。SELECT hostName() AS node, sum(bytes_on_disk) AS total_bytes FROM cluster(clickhouse_cluster, system.parts) WHERE active 1 GROUP BY node ORDER BY total_bytes DESC;如果某个节点的磁盘用量是其他节点的两倍以上基本可以断定分片键选择存在倾斜问题。缓解倾斜需要区分是“键本身倾斜”还是“数据随时间自然倾斜”。键本身倾斜比如按租户ID分片但租户体量差异巨大最有效的手段是引入复合分片键。生产上我常用的一种方案是cityHash64(concat(tenant_id, _, cityHash64(user_id))) % shards_num这样既保留租户维度的部分亲和性又避免单个租户压垮单个节点。如果业务上做不到可以在应用层按租户大小打散后再写入。数据随时间自然倾斜例如历史数据集中在一段时间段跨节点后某些分片存储被塞满这种情况更适合用ALTER TABLE ... MOVE PARTITION把部分分区从热节点迁移到冷节点。ClickHouse支持跨分片移动分区到任意节点的本地表实现手工再平衡。不过要提醒一句倾斜一旦写入就很难自动恢复重建表是成本最低的纠正手段。在数据量可控的时候重新建表、分片键优化、重新灌数往往比重做搬迁逻辑更省心。最后分享两个实操体会第一个体会是分片键和查询条件的匹配比集群规模重要得多。一个分片键设计合理的3节点集群性能往往好过一个分片键糟糕的10节点集群。规划ClickHouse集群时先想清楚线上查询会按什么字段过滤再决定怎么分片顺序不要反。第二个体会是分布式查询性能瓶颈多数在网络传输和汇总开销而不是单节点计算能力。所以做性能测试时不要只盯单条SQL的执行计划还要观察system.query_log里的read_bytes、memory_usage、result_rows这三个字段它们能直观反映查询是否走了最优路径。ClickHouse这种工具配置项再多也不可怕关键是找到真正影响性命的那几个开关然后持续用数据说话。