凌晨一点半,手机告警把我从床上拽起来:集群里某个BE的数据目录使用率到了80%,另外两台才不到20%。打开监控一看,磁盘占用曲线从三个月前就开始分叉,晚高峰查询越来越慢,最夸张的一条SQL跑了四十多秒还没出结果。这不是网络问题,也不是机器问题,是典型的数据倾斜——而且根子,早在建表那一刻就埋下了。
DorIS集群里的数据倾斜,十有八九是分区分桶设计没想清楚导致的。分区管数据的"收纳范围",分桶管数据的"落点分布",这两件事经常被放在一起说,但职责完全不同。本文就围绕Doris数据倾斜处理与查询优化,从分区分桶的原理讲起,到怎么排查倾斜、怎么选分桶键、怎么改查询,最后完整复盘一次线上故障的处理过程。适合正在维护Doris集群、经常被慢查询和磁盘不均衡折磨的兄弟们,也适合准备建表但还没想好分桶策略的新手。
1. 数据倾斜的本质:先知道数据在Doris里是怎么落的
很多人在建表时直接把PARTITION BY和DISTRIBUTED BY当两个必填项抄过去,完全没想过它们各自干了什么。等线上出了倾斜再回头补课,成本就高了。
1.1 分区管"好删快查",分桶管"均匀散开"
打个比方。分区就像档案室按年份分柜子:2023年的材料放一个柜子,2024年的放另一个。到了年底要销毁过期档案,直接把整个柜子拖走就行,不需要翻每一页;查询时也只需要打开对应年份的柜子,不用把整个档案室翻一遍。这就是分区带来的生命周期管理和分区裁剪能力。
分桶则是柜子内部的隔层。假设2024年这个柜子里有几百万份材料,如果不做分隔,查询时只能整柜翻找。现在按"姓氏首字母哈希"分成26个格子,找人时直接去对应格子,速度就上去了。分桶决定了数据在物理上的分散程度,也决定了查询和聚合能开多大并行度。所以:
- 分区解决的是"数据管理边界"问题,按时间切分最实用。
- 分桶解决的是"数据分布形态"问题,决定数据落在哪个tablet上。
一个分区内可以有很多个分桶,每个分桶对应一个tablet,tablet的副本分布在不同的BE节点上。Doris的查询是以tablet为最小扫描单位的,分桶数越多,扫描并行度越高,但tablet总量也会膨胀。
1.2 数据倾斜是建表时期埋下的雷
分桶的原理是哈希取模:hash(分桶键) % 分桶数得到桶号。哈希算法本身是均匀的,但均匀的前提是分桶键的值足够多样。
如果分桶键只有少数几个值,情况就麻烦了。比如订单表按order_status分桶,状态只有"待支付、已支付、已取消"三个值,取3个桶。哈希之后,大量数据会集中落在某一个桶里,另外两个桶几乎空闲。表现在集群层面,就是某些BE磁盘疯涨,某些BE整天摸鱼;表现在查询层面,就是某个tablet扫描耗时特别长,整个SQL被这一个桶拖死。
这就是为什么Doris社区一直强调分桶键要选高基数字段。高基数意味着每个值对应的数据量相对少,哈希后各个桶的数据量趋于均衡。低基数键做分桶,等于把数据倾斜这件事提前钦定了。
1.3 分桶数不是越大越好:tablet数量是个约束
有人觉得分桶数越大并行度越高,于是上来就搞128个桶512个桶。但分桶数一旦和分区数、副本数乘起来,tablet总量会非常恐怖。
tablet总数 = 分区数 × 分桶数 × 副本数。
简单算一笔账:按天建分区,一年365个分区,32个分桶,3副本,总量是35040个tablet。如果分桶数提到128,就变成14万个tablet。每个tablet在BE上都要有对应的目录、元数据、compaction任务,tablet太多会让BE启动变慢、元数据加载变慢、调度开销变大,反而拖垮集群。
经验上,单个tablet的数据量控制在1GB到10GB之间比较合理。假设单日数据量约1亿行,表宽度中等大约20GB,那一天的分区建议分20到40个桶。分桶数取2的幂次,32基本够用。同时还要参考BE节点数:分桶数至少要比BE数大,不然会出现某些BE完全分不到tablet副本的情况。更稳的做法是分桶数取BE数的整数倍,比如3台BE就取12、24、48这样的值,让每个节点的副本尽量均衡。
2. 现场排查:三步定位数据歪到哪里去了
数据倾斜不是靠感觉判断的,要用命令和指标把"歪"的程度量化出来。我排查倾斜的一般路径是:先看节点,再看表,最后看查询。
2.1 先看节点:BE磁盘差异是最大的警报
倾斜最直观的表现就是BE之间磁盘占用差距大。直接在MySQL协议连接FE后执行:
SHOW BACKENDS;输出里有DataUsedCapacity字段,能直接看到每个BE的数据占用。如果3台BE分别是1.2T、400G、350G,那就不用犹豫了,必然存在数据分布不均。
这里要注意,BE磁盘不均也可能来自副本修复、迁移任务没跑完等临时状态。所以看磁盘之前,先看一眼SHOW PROC '/cluster_balance',确认没有正在进行的副本均衡任务。如果集群本身在忙,等它跑完再看,不然容易被表象误导。
2.2 再看表:用 SHOW TABLET 找出异常桶
确认节点级别有倾斜后,下一步定位到具体表。Doris提供了直接查看tablet分布的命令:
SHOW TABLET FROM db_name.table_name;输出会包含TabletId、ReplicaId、BackendId、DataSize、RowCount、Version等字段。用它来比较某个分区下不同tablet的行数差异。
我通常的做法是查出来之后按RowCount排序,看前几名和后几名的差距。如果某个tablet有几千万行,另一个tablet只有几万行,差距两个数量级以上,那这张表的分桶策略基本可以判定有问题。这里要补一句:SHOW TABLET默认输出全部分区的tablet,数据量大时可以先根据表名和分区,用类似WHERE的条件过滤,或者直接用程序拉取后分析,避免在客户端一次性刷几千行。
2.3 最后看查询:Profile 里每个 instance 的扫描行数藏着真相
存储层面的倾斜已经确认,但为了在汇报时拿出硬核证据,最好再把查询Profile甩出来。Doris中开启查询Profile的方式是:
SET enable_profile = true;然后重新执行那条慢SQL,结束后去FE的Web页面(通常是8030端口)找到对应的Query Profile,展开ScanNode部分,看每个instance的rows和execution time。
如果发现某个instance扫描的行数是其他instance的5倍、10倍,扫描耗时也成倍放大,这就把倾斜的影响从"磁盘不均"延伸到了"查询被拖慢"的直接证据链。Profile里的信息很多,第一次用容易看花眼,先盯 ScanNode 的 rows 和 exec time 就行,其他性能问题后面再细分析。
3. 分桶键选型的实战:三原则和一个翻车案例
排查倾斜只是事后擦屁股,真正值钱的是建表阶段想清楚。分桶键的选择没有银弹,但有明确的判断标准。
3.1 选键三原则:高基数、分布均匀、贴合查询
第一原则,高基数。分桶键的取值数量要远大于分桶数,至少是分桶数的几十倍到几百倍。比如32个桶,分桶键至少有上千个不重复值,哈希之后才会散得开。用户ID、订单ID、设备ID这类字段天然具备高基数,是分桶键的首选。
第二原则,分布均匀。有些字段基数很高,但分布极其不均匀。比如"店铺ID",头部大卖家的订单量可能占全平台的30%,按它分桶同样会倾斜。选键前一定要跑一遍分布探查SQL:
SELECT shop_id, COUNT(*) AS cnt FROM order_table GROUP BY shop_id ORDER BY cnt DESC LIMIT 10;如果TOP1的占比超过5%到10%,这个键就要谨慎了。注意,分桶是哈希,不是取模范围,单纯看TOP值不够,还得看整体分布的偏态,但TOP占比是最快的初筛手段。
第三原则,贴合查询。分桶键最好能覆盖高频查询中的等值条件。比如订单表最常见的是按user_id查订单,那分桶键选user_id,查询时就能通过哈希直接定位到少数几个tablet,减少扫描量。
3.2 翻车案例:按 order_status 建表后的三个月之痛
前面提到的那个故障表,建表语句简化后长这样:
CREATE TABLE order_table ( order_id BIGINT, user_id BIGINT, order_status VARCHAR(20), amount DECIMALV3(12, 2), dt DATE ) DUPLICATE KEY(order_id) PARTITION BY RANGE(dt)() DISTRIBUTED BY HASH(order_status) BUCKETS 16 PROPERTIES ( "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "DAY", "dynamic_partition.end" = "3", "dynamic_partition.start" = "-30" );当时选order_status当分桶键的理由是"查询经常按状态过滤",这个逻辑乍一听没毛病,但完全忽略了基数问题。状态就3个值,哈希后数据必然聚集在部分桶里。上线头一个月数据量小,问题不明显,等累计到一定规模,倾斜开始放大:某些tablet的RowCount是其他tablet的几百倍,BE磁盘差距越拉越大,按状态维度的聚合查询慢到无法接受。
这个案例给我最大的教训是:分桶键和过滤条件不是一回事。过滤条件只需要建合适的索引或者靠分区裁剪,分桶键是用来决定数据物理分布的,优先保证均衡,其次才考虑裁剪。
3.3 没有好键怎么办:随机分桶和动态分区兜底
有一种情况是:表里确实找不到一个同时满足高基数、均匀、贴合查询的字段。比如一张日志大宽表,查询条件五花八门,没有哪个字段能扛起"高频等值过滤"的担子。这时候可以考虑随机分桶:
DISTRIBUTED BY RANDOM BUCKETS 32随机分桶的做法是把数据按导入顺序随机打散到各个桶,数据分布稳定性不错,但失去了按分桶键裁剪的能力,适合以全表扫描为主的场景。如果你的查询都是范围扫描加复杂过滤,随机分桶是可以接受的兜底方案。
另外,分区键的选择也要配套。数据仓库类表按时间分区没问题,但如果你的表是维表、配置表、小码表,数据量不大还按天分区,反而制造大量空分区和空tablet,白白增加元数据负担。维表这类小表不需要动态分区,直接单分区配合合适的分桶数就行。
3.4 分桶键不能改,分桶数可以调
建表后如果发现分桶数不够,可以调整:
ALTER TABLE order_table DISTRIBUTED BY HASH(user_id) BUCKETS 64;但有个魔鬼细节:分桶键本身是不允许修改的。上文这个SQL把分桶数从16改到64可以,但想把分桶键从order_status换成user_id,直接ALTER是行不通的,只能通过新建表+数据迁移来重建。所以建表前想清楚分桶键,比上线后做任何优化都省钱。
调整分桶数也要挑低峰期做,因为这会让所有分区的tablet全部重新分布,涉及大量副本复制和迁移,对IO和网络的消耗不小。改完之后盯着BE的tablet数量变化,等到SHOW PROC '/cluster_balance'里不再有迁移任务,才算真正完事。
4. 查询侧的倾斜与慢查询优化:group by、join、count distinct 逐个拆
建表再合理,查询写得不讲究,照样会把性能做成灾难。这一节说几个我在Doris慢查询优化里最常碰到的场景。
4.1 group by 倾斜:二次聚合把热点摊开
现象很典型:一条SQL按"城市"分组统计订单量,某一线城市的数据量占全表40%,负责这个城市的聚合节点忙到飞起,其他节点早干完活等它。整体耗时被这个热点键拉满。
Doris默认有两阶段聚合机制,会自动做本地聚合再全局聚合,但如果热点key数据量实在大,自动两阶段也扛不住。可以手动把聚合拆细:
SELECT city, SUM(cnt) FROM ( SELECT city, user_id, COUNT(*) AS cnt FROM order_table GROUP BY city, user_id ) t GROUP BY city;思路是先按city + user_id做一层细粒度聚合,把同一个用户的多条记录压成一条,数据量缩小一个量级;然后再按city聚合。这个改写对"用户行为日志"这类同用户多次记录的场景特别管用。
如果SQL能容忍近似值,也可以用Doris的approx函数(如APPROX_COUNT_DISTINCT)替代精确count distinct,性能能快一个档次,这点在流量分析类报表里经常用到。
4.2 join 倾斜:Runtime Filter、Colocate Join 和拆key三选一
两张表join时,如果大表的join键分布不均,同样会出现某个节点处理的数据量远超其他节点的情况。排查方式还是看Profile里ScanNode之后的ExchangeNode:哪个instance收到的行数特别多,哪个就有热点。
常规优化三板斧:
- Runtime Filter。Doris默认开启,但如果关过或调过参数,可以检查
runtime_filter_mode。它会用左表(驱动表)的结果动态过滤右表的扫描,让右表少扫很多无关数据。适用场景是小表join大表,效果最明显。 - Colocate Join。两张表如果分桶键、分桶数、副本数都一致,并把它们放到同一个colocate group里,join时数据可以在本地完成,完全不需要shuffle。这是Doris处理多张大表频繁join的最优解,前提是两张表的DDL从一开始就对齐设计。
- 拆key。倾斜key就一两个,但它们占了海量数据。可以先把倾斜key的行拆出来单独join,再和正常部分union合并。这个方案实施成本高,一般作为最后手段。我遇到过一次,用SQL Case When配合两个子查询实现的,调了一个下午,效果明显但代码维护确实麻烦,非必要不上。
4.3 count distinct 别硬扛:bitmap 与 HLL
Doris里最容易被写坏的SQL就是COUNT(DISTINCT user_id)。这种写法会用精确去重,内存和网络消耗极大,数据量上来之后慢得让人怀疑人生。
精确去重场景,如果用户ID可以编码成整型,改用bitmap方式:
SELECT bitmap_count(bitmap_union(to_bitmap(user_id))) FROM order_table WHERE dt = '2025-01-01';bitmap_union是聚合算子,to_bitmap把user_id转成bitmap,多个分桶的bitmap做按位或再计数,计算量小非常多。如果建表时就把user_id定义为BITMAP类型,查询时直接用BITMAP_UNION就行,性能进一步上升。
如果指标对精度要求不高,比如PV/UV类报表能接受千分之几的误差,直接上HLL:
SELECT hll_union_agg(hll_hash(user_id)) FROM order_table WHERE dt = '2025-01-01';HLL的内存占用比bitmap还小,但存在一定误差率,适合大促实时大屏这类场景。一句话总结:能近似就别精确,能bitmap就别count distinct。
4.4 版本堆积导致的隐性慢查询:compaction 能救
还有一种查询变慢,跟倾斜没关系,但同样让人抓狂——表没有任何结构问题,数据分布也均衡,但就是越来越慢。这种情况十有八九是版本堆积。
Doris的存储引擎是类似LSM的追加写模式,每次导入都会生成一个新的版本。查询时要读取一个Tablet的快照,如果版本数太多,读路径会合并多个版本的数据,开销自然变大。小批量高频写入最容易攒版本,特别是几百KB一次、一天写几千次的那种。
排查方式还是用SHOW TABLET FROM table_name,看Version字段的数值。正常情况一个tablet的版本数应该在几十以内,如果看到几百上千,就得处理了。
Doris会自动做compaction,但积压严重时手动触发一次更直接:
ALTER TABLE order_table COMPACT;不同版本的Doris对COMPACT语法支持略有差异,执行前确认下当前版本的支持情况。触发后可以用SHOW PROC '/compactions'观察合并进度。更长效的办法是控制导入节奏,比如把高频小文件合并成大文件再导入,或者调大cumulative_compaction相关参数,让自动压缩跑得更快。这个属于运维长期优化项了。
5. 一次实时告警背后的分区分桶改造复盘
把前面所有知识点串起来,复盘一次完整的线上故障处理过程。这个案例就是我开头说的那个凌晨告警,整个过程从发现问题到修复完成花了大概两天。
5.1 告警到根因:磁盘不均衡的完整定位
告警内容是某个BE磁盘使用率超过80%。我没有直接去查大表,而是按先节点、再表、再查询的顺序走。
第一步,SHOW BACKENDS,确认三台BE的DataUsedCapacity分别是1.1T、360G、380G。差距超过3倍。
第二步,SHOW PROC '/cluster_balance',确认没有正在执行的副本均衡任务,排除临时迁移因素。
第三步,用脚本拉取全部门店表的SHOW TABLET信息,按DataSize排序,发现RowCount大的tablet集中在特定几个桶。
第四步,翻了建表记录,分桶键就是order_status。到这里根因已经很清晰:低基数分桶键导致哈希分布严重不均,热点数据全压在某几个tablet上。
5.2 改造方案对比:重建表回灌 vs 单分区替换
分桶键不能直接改,只能重建表。方案有两个:
- 重建整表:新建表时把分桶键改成
user_id,然后把老表数据全量INSERT INTO ... SELECT到新表。优点是结构彻底干净;缺点是表太大时回灌时间长,业务需要停写或双写,影响大。 - 只重导受影响分区:保留近期数据在老表,只对倾斜严重的历史分区做迁移。Doris分区可以单独替换,
ALTER TABLE DROP PARTITION加ALTER TABLE ADD PARTITION,再把对应分区的数据重新导入。优点是窗口短,缺点是无法根治未来分区的问题。
考虑到业务不能长时间停写,我们采用折中方案:新表直接上线,把当天数据切到新表写入;历史数据通过后台任务按天INSERT INTO ... SELECT回灌,同时用CCR Syncer把数据变更同步过去。整个过程业务影响控制在分钟级。
5.3 最终建表语句与优化后效果
重建后的核心表结构:
CREATE TABLE order_table_new ( order_id BIGINT, user_id BIGINT, order_status VARCHAR(20), amount DECIMALV3(12, 2), dt DATE ) DUPLICATE KEY(order_id) PARTITION BY RANGE(dt)() DISTRIBUTED BY HASH(user_id) BUCKETS 48 PROPERTIES ( "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "DAY", "dynamic_partition.end" = "3", "dynamic_partition.start" = "-60", "dynamic_partition.buckets" = "48", "replication_num" = "3" );关键改动就一处:分桶键从order_status换成user_id,分桶数从16提到48。user_id基数高、分布均衡,而且订单查询几乎都带用户维度,一箭三雕。
改造完成后的效果对比:
| 指标 | 改造前 | 改造后 |
|---|---|---|
| BE磁盘最大差距 | 1.1T vs 360G | 480G vs 420G |
| 某高频订单查询耗时 | 40秒以上 | 2秒以内 |
| 按状态聚合查询耗时 | 15秒左右 | 1.5秒 |
| 晚高峰集群CPU | 常驻85%以上 | 峰值65% |
数据回灌完成后,磁盘差距收敛到15%以内,查询全部回到秒级。
6. 日常运维中值得长期坚持的几个习惯
Doris跑得稳不稳,一半靠架构设计,一半靠日常巡检。分享几个我一直坚持的习惯。
6.1 每张新表上线前做一次数据分布探查
建表前花五分钟跑一个分布探查SQL,能避免80%的倾斜事故。重点看两个指标:分桶键的基数、TOP值的占比。如果TOP1占比超过总行数的5%,我会换一个候选键再测;如果实在找不到合适的,就用随机分桶。这个习惯成本极低,收效极高。我现在审批任何业务的建表DDL,都会先要这个探查结果。
6.2 慢查询巡检:轻量级方案
开启FE的慢查询审计日志,定期grep出超过阈值(比如2秒)的SQL,按出现频率排序。对TOP N的慢查询逐个看执行计划,能用分区裁剪解决的加分区条件,能用索引的加索引,能用物化视图的建物化视图。每周花一小时处理一轮,比月底一次性清理要轻松得多。查询超时时间也建议统一设置,比如在JDBC连接串里配置queryTimeout=30,配合Doris的query_timeout会话变量,防止个别慢SQL把连接池占死。
6.3 定期看版本数和Compaction状态
磁盘使用率只反映容量,不反映健康度。我习惯每周随机抽查几张活跃表的SHOW TABLET信息,看Version字段是否异常增长。如果发现某张表版本数持续偏高,就去查它的导入频率:是不是经常小批量写入?是不是有业务在循环insert单条数据?从源头调整写入节奏,比天天手动触发COMPACT省心得多。
我在实际运维Doris的过程中,最大的体会就是:很多看起来玄乎的线上问题,追到根上都是建表时的几个小决定没做对。分桶键选得好,后面所有查询优化都事半功倍;分桶键选得糙,后续花多少精力擦屁股都找补不回来。所以如果你只记住一句话,那就是:分区分桶不是建表模板里的填空项,而是Doris性能的命门。