十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Doris数据倾斜治理:分区分桶设计与查询优化实战

Doris数据倾斜治理:分区分桶设计与查询优化实战 凌晨一点半手机告警把我从床上拽起来集群里某个BE的数据目录使用率到了80%另外两台才不到20%。打开监控一看磁盘占用曲线从三个月前就开始分叉晚高峰查询越来越慢最夸张的一条SQL跑了四十多秒还没出结果。这不是网络问题也不是机器问题是典型的数据倾斜——而且根子早在建表那一刻就埋下了。DorIS集群里的数据倾斜十有八九是分区分桶设计没想清楚导致的。分区管数据的收纳范围分桶管数据的落点分布这两件事经常被放在一起说但职责完全不同。本文就围绕Doris数据倾斜处理与查询优化从分区分桶的原理讲起到怎么排查倾斜、怎么选分桶键、怎么改查询最后完整复盘一次线上故障的处理过程。适合正在维护Doris集群、经常被慢查询和磁盘不均衡折磨的兄弟们也适合准备建表但还没想好分桶策略的新手。1. 数据倾斜的本质先知道数据在Doris里是怎么落的很多人在建表时直接把PARTITION BY和DISTRIBUTED BY当两个必填项抄过去完全没想过它们各自干了什么。等线上出了倾斜再回头补课成本就高了。1.1 分区管好删快查分桶管均匀散开打个比方。分区就像档案室按年份分柜子2023年的材料放一个柜子2024年的放另一个。到了年底要销毁过期档案直接把整个柜子拖走就行不需要翻每一页查询时也只需要打开对应年份的柜子不用把整个档案室翻一遍。这就是分区带来的生命周期管理和分区裁剪能力。分桶则是柜子内部的隔层。假设2024年这个柜子里有几百万份材料如果不做分隔查询时只能整柜翻找。现在按姓氏首字母哈希分成26个格子找人时直接去对应格子速度就上去了。分桶决定了数据在物理上的分散程度也决定了查询和聚合能开多大并行度。所以分区解决的是数据管理边界问题按时间切分最实用。分桶解决的是数据分布形态问题决定数据落在哪个tablet上。一个分区内可以有很多个分桶每个分桶对应一个tablettablet的副本分布在不同的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%按它分桶同样会倾斜。选键前一定要跑一遍分布探查SQLSELECT 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 与 HLLDoris里最容易被写坏的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类报表能接受千分之几的误差直接上HLLSELECT 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 360G480G 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连接串里配置queryTimeout30配合Doris的query_timeout会话变量防止个别慢SQL把连接池占死。6.3 定期看版本数和Compaction状态磁盘使用率只反映容量不反映健康度。我习惯每周随机抽查几张活跃表的SHOW TABLET信息看Version字段是否异常增长。如果发现某张表版本数持续偏高就去查它的导入频率是不是经常小批量写入是不是有业务在循环insert单条数据从源头调整写入节奏比天天手动触发COMPACT省心得多。我在实际运维Doris的过程中最大的体会就是很多看起来玄乎的线上问题追到根上都是建表时的几个小决定没做对。分桶键选得好后面所有查询优化都事半功倍分桶键选得糙后续花多少精力擦屁股都找补不回来。所以如果你只记住一句话那就是分区分桶不是建表模板里的填空项而是Doris性能的命门。
返回列表