
搞分布式数据库这些年最常被问到的一句话就是分布式下到底能不能愉快地JOIN我最近在PolarDB-X上专门做了一轮Broadcast Join和Shard Join的对比实测用TPC-H风格的数据跑了几个典型场景把执行计划、响应时间、网络开销、踩坑点全部记录下来。结论先放在前面同分片键下的Shard Join可以把分布式JOIN的耗时拉回接近单机水平而Broadcast Join在做大小表关联时依然是性价比极高的默认策略但前提是你得搞清楚优化器为什么这么选以及什么时候它会选错。这篇文章面向正在做分布式数据库选型、或者已经把PolarDB-X原DRDS体系用起来的开发和DBA同学我会把测试设计、执行计划识别、关键参数和问题排查方法都摊开讲。1. 为什么分布式JOIN是绕不开的坎1.1 从单机JOIN到分布式JOIN单机数据库里写JOIN是一件天经地义的事执行引擎打开两张表选一个哈希连接还是嵌套循环连接索引有没有用上内存够不够排序基本上都围绕CPU、内存、磁盘这三样东西在转。到了分布式数据库这里数据被水平拆到了多个数据节点DN上两张表的关联数据很可能根本不在同一个节点里JOIN就不再是“读出来、连起来”这么简单了它变成了一个分布式计算问题。我打个比方。单机JOIN就像大家都在同一个自习室找人直接喊一嗓子就行。分布式JOIN等于人分散在不同教学楼要么把其中一栋楼的人全部集中到另一栋楼再找广播要么提前按教室把座位排好让相关联的人本来就在同一间教室里分片对齐。前者灵活但路上耗时后者需要提前设计好座位表但找起来最快。PolarDB-X的分布式JOIN核心就是围绕这两种思路展开的。1.2 一次PolarDB-X查询背后发生了什么要理解两种JOIN策略的差距先要清楚PolarDB-X的执行架构。它大致分三层GMS负责全局元数据和分片信息CN是计算节点负责接收SQL、做语法解析、逻辑优化、物理优化生成分布式的执行计划并调度执行DN是数据节点真正存储数据并执行下推过来的SQL片段。比如一条简单的两表JOIN查询CN会先判断JOIN条件里的列和两张表的分片键能不能对上如果能对上就把整个JOIN下推到每个DN去并行执行CN只负责汇总结果。如果对不上CN就得想办法把数据搬到一起来要么把其中一张表广播到所有DN要么把两张表都按JOIN列重新打散再连接。这一搬网络传输就来了性能的差距也来了。所以分布式JOIN的性能很大程度上在表结构设计那一刻就已经被决定了。这也是为什么很多业务从单机迁到分布式之后明明单表查询都快一遇到JOIN就慢得离谱——不是数据库不会做JOIN而是分片键从一开始就没设计对。2. Broadcast Join与Shard Join的原理和适用场景2.1 Broadcast Join把“小字典”复制给所有人Broadcast Join的思路非常直白既然两张表的数据不在一起那我把其中一张表通常是小表完整发给所有数据节点让每个节点在本地完成连接。假设orders表有1000万行分布在8个DN上customer表有100万行也分布在8个DN上两边分片键不一致。执行JOIN时CN可以把customer表的数据广播到所有涉及orders的DN节点每个DN拿到完整customer副本后和本地orders切片做连接。这样orders仍然保持分布式并行扫描不会因为JOIN而变成单点处理。这个策略的代价主要是网络传输和DN内存。传输量大概等于小表数据量乘以DN数量比如100MB的小表广播到8个DN网络层面就要搬800MB。所以Broadcast Join特别适合“大表 JOIN 小表”的场景小表控制在百万行级别以内比较舒服。如果所谓的小表本身就有几千万行广播代价就非常难看了。我用一个成本公式帮助理解Broadcast成本 ≈ 大表扫描成本 小表全量传输成本 × DN数 各DN对小表副本扫描/哈希成本好处是大表只扫一遍小表虽然被扫了N遍但因为它小总体可控。这也是绝大多数分布式数据库在优化器里默认优先考虑的策略之一。2.2 Shard Join让数据在分片内“门当户对”Shard Join有些资料也叫Partition-Wise Join或者Colocate Join思路就高级一些。如果两张表用同一个分片键、相同的分片数量、相同的分区方式并且JOIN条件里包含这个分片键的等值条件那么可以说数据在写入时就已经按照关联关系提前对齐了。还是上面那个例子如果orders和customer都按客户ID分片orders的o_custkey和customer的c_custkey值域是一致的。那么某个客户的全部订单和这个客户的档案本来就在同一个DN上。执行JOIN时每个DN只需要把自己手上的orders分片和customer分片做本地连接然后CN把8个DN的结果汇总就行全程没有一行数据需要跨节点搬迁。这种 JOIN 的性能最接近单机数据库也是我在这次测试里最推荐的设计目标。它有几个硬性前提JOIN条件是分片键的等值条件两张表的分片键类型、分区方式、分区数量完全一致查询条件没有破坏分片键的下推能力比如对分片键套了函数就不会被识别。如果分片键不一致或者JOIN条件里不包含分片键那么即使能下推也只是把单张表的扫描下推下去连接动作还是会被提到CN层来做性能会明显下降。2.3 从执行计划快速识别JOIN策略判断一条SQL在PolarDB-X上走的到底是Broadcast Join还是Shard Join最直接的办法就是看EXPLAIN输出。我测试时常用的判断口诀看JOIN算子在Gather上面还是下面。看一个同分片键下推成功的执行计划Shard JoinEXPLAIN SELECT count(*) FROM orders o JOIN customer c ON o.o_custkey c.c_custkey WHERE c.c_nationkey 1 AND o.o_orderdate 2020-01-01 AND o.o_orderdate 2020-04-01;Gather(concurrenttrue) LogicalView(tablesorders_[0-7],customer_[0-7], shardCount8, sqlSELECT count(*) FROM orders o JOIN customer c ON o.o_custkey c.c_custkey WHERE c.c_nationkey ? AND o.o_orderdate ? AND ...)这里的核心特征是JOIN在LogicalView内部也就是在DN上完成的CN层的Gather只负责把各分片算出来的count汇总没有跨节点搬运数据。再看分片键不对齐时可能出现的执行计划HashJoin(conditiono.o_custkey c.c_custkey, typeinner) Gather(concurrenttrue) LogicalView(tablesorders_[0-7], shardCount8, sqlSELECT ...) Gather(concurrenttrue) LogicalView(tablescustomer_[0-7], shardCount8, sqlSELECT ...)这种形态下JOIN算子跑在Gather之上意味着两张表的数据都被拉到了CN层由CN完成最终的连接。这个计划不一定等于Broadcast Join也可能是Shuffle重分布Join但无论如何数据跨节点搬运已经发生了代价必然比Shard Join高。我建议每个使用PolarDB-X的团队把核心JOIN语句的EXPLAIN都收集一份存起来作为上线评审的一部分。执行计划好不好一眼就能看出来等线上慢查询出现再去分析就晚了。3. 测试环境与数据集设计3.1 集群拓扑与配置这次测试用的是PolarDB-X标准分布式部署拓扑是1个GMS节点、2个CN节点、8个DN节点。CN和DN规格均为8核32G操作系统和数据库版本就不贴细节了避免不同小版本行为差异带来误导。我特别提一下PolarDB-X不同版本在优化器细节和Hint语法上可能会有差异但核心的Broadcast/Shard逻辑是大体一致的。存储用的是ESSD云盘网络环境是万兆内网。这个配置不算高配但对验证JOIN策略的差距足够了。因为我们要对比的是同一条SQL在相同硬件下不同分片设计带来的性能差异重点在相对差距而不是绝对耗时。3.2 表结构与分片键设计我设计了两套Schema用来模拟“分片键对齐”和“分片键不对齐”两种场景。第一套是对齐场景表结构如下CREATE TABLE customer ( c_custkey BIGINT NOT NULL, c_name VARCHAR(25) NOT NULL, c_nationkey INT NOT NULL, c_phone VARCHAR(15), c_mktsegment VARCHAR(10), PRIMARY KEY (c_custkey) ) PARTITION BY HASH(c_custkey) PARTITIONS 8; CREATE TABLE orders ( o_orderkey BIGINT NOT NULL, o_custkey BIGINT NOT NULL, o_orderstatus CHAR(1), o_totalprice DECIMAL(12,2), o_orderdate DATE NOT NULL, PRIMARY KEY (o_orderkey) ) PARTITION BY HASH(o_custkey) PARTITIONS 8;orders表按o_custkey分片customer表按c_custkey分片两边在“客户ID”这个维度上天然对齐JOIN条件写o_custkey c_custkey就能走Shard Join。第二套是把customer表改成按c_nationkey分片CREATE TABLE customer ( c_custkey BIGINT NOT NULL, c_name VARCHAR(25) NOT NULL, c_nationkey INT NOT NULL, c_phone VARCHAR(15), c_mktsegment VARCHAR(10), PRIMARY KEY (c_custkey) ) PARTITION BY HASH(c_nationkey) PARTITIONS 8;这种情况下即使JOIN条件还是o_custkey c_custkey但因为customer在本地DN上的数据不是按c_custkey组织的无法直接和orders的本地分片对齐优化器只能选择Broadcast或者重分布。3.3 数据生成与装载要点数据规模我定在orders 1000万行、customer 100万行参考TPC-H生成逻辑造数。生成工具可以自己写脚本也可以直接用TPC-H的dbgen然后改列名。导入方式我建议用DataX或者多线程分批INSERTLOAD DATA虽然也可以用但分布式下要注意分片路由和主键冲突的问题。数据装载完成后我强烈建议立刻执行统计信息收集ANALYZE TABLE orders; ANALYZE TABLE customer;这一步非常关键。PolarDB-X的优化器依赖统计信息来估算行数和选择JOIN策略如果统计信息缺失或者严重过期优化器可能把一个本该走Shard Join的查询错误地变成了Broadcast Join或者把广播小表选成广播大表。我在测试中就遇到过因为没收集统计信息导致同一条SQL执行计划不稳定的情况。另外还要准备一份倾斜数据。我在orders里额外造了一个特殊客户让这个客户的订单量占总订单量的20%也就是200万行用来测试数据倾斜对两种JOIN策略的影响。4. 性能实测与结果对比4.1 场景一大小表关联的Broadcast Join先看最典型的场景大表JOIN小表。测试SQL带上了过滤条件让结果集更接近真实业务SELECT count(*) FROM orders o JOIN customer c ON o.o_custkey c.c_custkey WHERE c.c_nationkey 1 AND o.o_orderdate 2020-01-01 AND o.o_orderdate 2020-04-01;customer经过nationkey1过滤之后大概剩10万行orders过滤后大约250万行这是非常舒服的“大表 JOIN 小表”场景。在分片键不对齐的Schema B下优化器选择了Broadcast Join。执行计划里能看到customer表的所有分片数据被拉起来广播到orders涉及的8个分片最终在CN层或者DN本地完成连接。由于广播数据量不大整体耗时还可以接受测试三次取中位数耗时在2.3秒左右。我专门看了一下网络监控broadcast阶段把约20MB的customer结果集复制到8个DN实际网络传输大约160MB。按万兆网络理论带宽来算这部分应该只有0.2秒左右的开销但实际表现会受序列化、DN内存拷贝和哈希构建影响整体上广播阶段占查询总耗时的占比还不小。4.2 场景二同分片键下的Shard Join同一套SQL换到Schema A也就是orders和customer都按客户ID分片的场景这次走了Shard Join。执行计划中JOIN被完整下推到DNCN层只汇总8个分片的count结果。实测耗时0.9秒左右比Broadcast Join快了大约2.5倍。注意这个差距是在小表只有10万过滤行、广播数据量不大的情况下出现的。如果小表再大一些比如过滤后有几百万行Broadcast Join的传输量翻几倍差距会更夸张。为什么Shard Join能这么快因为每个DN只需要处理自己本地大约125万行orders和12.5万行customer的连接8个分片并行跑热点不重、网络不搬数据、CN不参与连接计算。整个查询的资源消耗非常干净DN的CPU是主要开销网络几乎可以忽略。4.3 场景三分片键不一致引发的重分布我很好奇如果分片键不一致优化器在什么时候不会选Broadcast而选重分布。于是我把customer表调整到1000万行模拟“两个大表JOIN”的情况。这种情况下如果还走Broadcast广播的数据量是1000万行每个DN都要接收完整拷贝网络和内存都会爆炸。优化器显然也明白这一点它选择了把orders和customer都按o_custkey/c_custkey重新打散再在CN层做HashJoin。这个过程就是典型的Shuffle Join重分布需要把orders的1000万行和customer的1000万行通过网络按JOIN列重新分配。按每行平均100字节估算网络传输量轻松超过1GB整体耗时直接飙到4.7秒左右。这个结果说明一个很现实的问题当业务里存在两个体量都不小的表需要JOIN时分片键不对齐的代价是成百上千倍的网络开销放大。不仅仅是慢还会让集群的网络和CN内存水位在查询期间明显上涨影响同一时间跑在上面的其他查询。4.4 数据倾斜场景下的表现我把倾斜数据加进来之后两种JOIN策略的表现又开始分化。某个客户的200万订单经过HASH分片后全部落在同一个DN上Shard Join执行时这个热点DN要比其他DN多处理很多数据整个查询的耗时被这个慢分片拖住。实测从正常的0.9秒涨到了1.6秒左右但整体仍然在可控范围。Broadcast Join在倾斜场景下更惨一点因为广播过来的customer副本在每个DN都要参与连接热点DN上不仅要处理倾斜的订单数据还要构建和探测哈希表加上网络传输开销实测耗时到了3.5秒。这里我想强调一个经验数据倾斜对任何分布式JOIN策略都不友好但相对而言Shard Join的短板更可控。热点DN的负载可以通过拆分大客户、调整分片键、给热点键加盐等方式缓解而Broadcast Join一旦小表变大网络瓶颈会变成全局问题影响的不只是这一条SQL。三类场景的对比汇总如下场景JOIN策略耗时中位数网络传输关键瓶颈大表 JOIN 小表10万维表Broadcast Join2.3s约160MB广播传输 DN内存同分片键大表 JOIN维表10万Shard Join0.9s接近0热点DN计算两个大表 JOIN各1000万Shuffle重分布4.7s超过1GBCN层计算 网络同分片键 数据倾斜Shard Join倾斜1.6s接近0热点分片Broadcast 数据倾斜Broadcast Join3.5s约160MB网络 热点分片以上数据是基于我这次测试环境得到的不同规格、不同版本的集群会有差异。但相对趋势是稳定的Shard Join优于Broadcast JoinBroadcast Join优于重分布Shuffle数据倾斜会放大所有策略的缺点。5. 这些性能差距对业务场景意味着什么5.1 交易类系统的订单-客户关联在电商、金融、订单中心这类系统里订单表和用户表关联是非常高频的查询。用户表的体量通常会增长到千万级甚至亿级订单表更是动辄上亿。这种场景如果用户表没有和订单表按同一个维度分片每次关联都走广播或者重分布查询延迟和数据传输量都会非常难看。我见过一个实际案例订单表按买家ID分片用户表却按用户注册时间分片两边做关联查询时PolarDB-X把所有用户数据都拉到CN层做连接用户量一上来查询直接超时。后来把用户表改成按用户ID分片同样一条SQL从几十秒降到了1秒以内这就是Shard Join带来的直观收益。5.2 实时报表与宽表构建实时报表场景经常要把大明细表和多个维度表JOIN再汇总出指标。如果维度表数量多、体量小比如地区表、渠道表、品类表这类数据适合建成广播表让每个DN都保留一份完整副本和明细表的JOIN完全本地化。广播表在PolarDB-X里可以直接建CREATE TABLE nation ( n_nationkey INT NOT NULL, n_name VARCHAR(25) NOT NULL, PRIMARY KEY (n_nationkey) ) BROADCAST;但要守住一个底线广播表一定是稳定的、体量小的表。如果维度表持续膨胀到千万行以上广播的收益就变成负担了。宽表构建时要定期梳理维度表的体量该转成普通分片表就及时转。5.3 峰值流量下的资源规划大促和业务高峰期间数据库最怕的不是CPU高而是网络和内存被打爆。Broadcast Join和Shuffle Join都会在某一瞬间产生较大的网络脉冲特别是多个大查询同时并发时带宽争抢会拖慢所有查询。如果你核心链路的JOIN都设计成了Shard Join那么峰值期间的容量规划就简单很多主要关注每个DN的数据倾斜和CPU水位就够了。反过来说如果核心链路大量依赖Broadcast Join网络带宽和CN内存就要留出足够余量最好在压测阶段就把这些查询跑一遍直接观察网络吞吐峰值。6. 常见问题与调优建议6.1 问题速查表我把这次测试和平时运维中遇到的典型问题整理成一个速查表方便直接对号入座。现象可能原因排查手段处理建议分片键一致却没走Shard JoinJOIN条件没包含分片键等值条件EXPLAIN看JOIN算子位置改写SQL让分片键参与等值JOIN分片键一致但执行计划仍广播分片键列被函数包裹无法下推查看SQL写法检查是否有函数去掉函数或者用冗余列存储加工后的值同一条SQL执行计划不稳定统计信息缺失或过期ANALYZE TABLE后重新EXPLAIN建立定期统计信息收集任务广播小表太大导致DN内存上涨维度表膨胀广播传输量成倍增加监控DN内存和网络转普通分片表或拆分维度某个DN明显比其他DN慢分片键分布不均数据倾斜查看各分片行数和DN慢日志调整分片键、热点键加盐、拆大客户JOIN列字符集或类型不一致隐式类型转换导致无法下推检查两张表DDL字符集和字段类型统一为相同类型和字符集6.2 分片键与表模型设计建议经过这轮测试我对PolarDB-X表模型设计有几点很实在的建议。第一核心业务表的关联关系在设计之初就要梳理清楚。主维度和事实表之间尽量选择同一个业务主体作为分片键。订单表按买家ID分片用户表也按用户ID分片流水表按商户ID分片商户表也按商户ID分片。这是Shard Join能成立的根本前提。第二维表分三类处理。体量小且稳定的维表建广播表所有节点都有一份副本最省心体量中等、增长可控的维表和大表同分片键对齐体量巨大而且无法对齐的维表要做专门设计比如只把JOIN需要的字段冗余到事实表从源头减少JOIN。第三不要迷信分片数越多越快。分片数量必须和DN规模匹配如果分片数量不一致即使分片键对齐也无法形成真正意义上的Shard Join。我见过有人把分区数从8改成16之后执行计划突然退化成了Broadcast就是因为另一张表还是8个分区两边分区数对不上。6.3 用Hint干预JOIN策略有时候优化器确实会选错比如它低估了某张表的行数选了Broadcast Join导致小表实际不“小”。这时候可以通过Hint或参数临时干预。PolarDB-X上我比较常用的是会话级参数和Hint两种方式。会话级参数示例SET ENABLE_BROADCAST_JOIN TRUE;Hint示例/*TDDL:CMD_EXTRA(ENABLE_BROADCAST_JOINTRUE)*/ SELECT count(*) FROM orders o JOIN customer c ON o.o_custkey c.c_custkey WHERE c.c_nationkey 1;还有一些环境里可以设置ENABLE_SHUFFLE_JOIN和BROADCAST_JOIN_SIZE来控制重分布和广播的数据量阈值。我的经验是这类参数在测试环境验证过之后才能上生产而且一定要配合EXPLAIN确认执行计划真的变了。不同版本对Hint的支持细节有差异用之前先看官方文档或者直接在测试库跑一下。还有一个不太起眼但很关键的技巧在业务SQL里主动把分片键条件写清楚。比如查询某个客户的订单和档案加上o_custkey 12345这种条件PolarDB-X可以做分片裁剪只路由到1个DN执行配合Shard Join性能直接拉满。我在实际测试里还有一个体会值得单独说一下。PolarDB-X这类分布式数据库的JOIN性能七分靠表结构设计三分靠SQL写法真正留给优化器临场发挥的空间没有想象中那么大。很多线上的慢JOIN根因不是数据库能力不够而是分片键从一开始就没对齐。所以我建议所有用分布式数据库的团队把核心JOIN语句的EXPLAIN纳入日常上线流程每次发布前都看一眼JOIN算子是不是在Gather下面如果不在那就要想清楚为什么不在以及你愿不愿意为这个“不在”承担网络开销。这个习惯坚持半年能帮你避开绝大多数JOIN性能事故。