ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

PolarDB-X JOIN性能实测:Broadcast Join vs Shard Join全场景对比

PolarDB-X JOIN性能实测:Broadcast Join vs Shard Join全场景对比 做分布式数据库的人应该都有同感JOIN 这玩意儿在单机 MySQL 里就是个优化器选索引的事顶多加个 Buffer Pool 就能压住可一旦到了 PolarDB-X 这种水平拆分的分布式架构里一条 JOIN 走什么执行路径直接决定了你的查询是几十毫秒还是几十秒。我最近花了两周时间专门针对 PolarDB-X 的 Broadcast Join 和 Shard Join 两条执行路径做了一轮完整的 Benchmark把从百万行到亿级行、从分片键对齐到完全不对齐的各种场景都跑了一遍把执行计划、耗时、资源占用、并发影响全部记录下来。这篇文章就是给两类人看的一类是正在用 PolarDB-X 的 DBA 和业务开发遇到分布式环境下的慢 JOIN 不知道怎么定位另一类是在做分布式数据库选型或者查询优化调研的同学想从实测数据里理解两种 JOIN 策略各自的边界和代价。整个测试过程不算复杂但里面确实有一些踩了才知道的坑比如同一个 SQL 语法分片键对齐和不对齐的性能可以差出一个数量级又比如小表广播并不是永远安全驱动表一旦超过某个量级广播路径会让整个集群的内存和网卡同时告警。下面我会从原理讲起然后是环境、表结构、测试场景设计最后给出每个场景的执行计划和实测数据。1. 分布式 JOIN 的痛点与这次测试的目标1.1 为什么 JOIN 到了分布式架构里会突然变慢单机数据库里 JOIN 的核心矛盾是“两个表的数据怎么在内存里相遇”办法基本就是双层循环、哈希连接或者归并连接数据都在本地问题只是算法选择。分布式数据库完全不一样数据按分片键被拆到多个数据节点DN上一条 JOIN 如果关联字段恰好不是分片键那么两个表的数据物理上就不在同一个节点服务器必须先把数据搬来搬去才能完成连接。用人话打个比方单机 JOIN 是两个人坐在同一个房间里对答案拿起来就能比分布式 JOIN 是每个人在不同房间A 房间的人想知道 B 房间的数据长什么样只能先把 B 房间的答案复印一份通过走廊发给所有人。分片数量越多这个“复印发送”的成本就越高。水平拆分本身是为了扩展吞吐但 JOIN 这种需要跨节点关联数据的操作恰恰是拆分逻辑里最容易被设计坑到的一环。PolarDB-X 的核心架构是 CN计算节点加 DN数据节点两层。CN 负责 SQL 解析、优化和执行计划生成DN 负责实际的数据存储和本地计算。一条查询发到 CN 后优化器会根据统计信息和代价模型决定是把执行计划下发到各个 DN 并行执行还是在 CN 层做数据的重新分布。JOIN 慢的根源就在这里数据一旦需要跨 DN 流动序列化、网络传输、反序列化这几项开销就全部叠加进来了。1.2 本文里两条路线的定义先明确一下本次测试里两个关键概念Shard JoinSQL 的 JOIN 条件里包含了两张表的分片键等值条件每个 DN 只需要用本地分片里的数据完成连接最后把结果汇总回 CN 即可。整个过程中没有跨 DN 的数据搬移本质上就是把一个分布式 JOIN 拆成了 N 个互不干扰的本地 JOIN 并行执行。Broadcast JoinCN 把驱动表一般是小表的完整结果集复制 N 份广播到所有 DN每个 DN 拿到这份“完整小表”后再与本地数据做 JOIN。代价是网络上会出现一次驱动表数据量的放大传输收益是省去了大表的重分布。除了这两种分布式 JOIN 还有一种常见实现是 Shuffle Join两张表都按关联键重新分布让关联键相同的数据落到同一个 DN 上再做本地 JOIN。这种方式代价最高但有些场景下不可避免。本次测试里也顺带跑了一组 Shuffle Join 作为参照目的是把问题边界画清楚。1.3 这次测试到底想验证什么我在最初设计 Benchmark 的时候给自己列了四个明确的问题在分片键完全对齐的情况下Shard Join 是不是真的能做到接近线性扩展是不是所有 JOIN 场景下的性能天花板Broadcast Join 在小表 JOIN 大表场景里到底能省多少时间它的收益边界在哪里PolarDB-X 的优化器在什么数据量、什么过滤条件下会选择从 Shard Join 切到 Broadcast Join这个切换点是否合理对业务侧来说有没有一些不改变业务逻辑、只改 SQL 写法就能显著提升 JOIN 性能的通用手段带着这四个问题去设计测试用例比单纯跑一个“谁快谁慢”的对比要有意义得多。后面每个场景我都会先给出 SQL再贴执行计划最后分析数据这样看文章的人能直接复现也能根据同样的方法去测自己的集群。2. Broadcast Join 与 Shard Join 的原理拆解2.1 Broadcast Join用网络换内存适合小表驱动大表Broadcast Join 的执行流程可以拆成三步。第一步CN 先读取驱动表通常是 FROM 后面的左表也可能是优化器选出来的小表中参与 JOIN 的字段形成一份完整数据集。第二步这份数据集按分片数量复制多份分别发送到每个 DN。第三步每个 DN 把收到的数据构建成哈希表再与本地存储的大表数据做 Hash Join最后把各自的结果集汇总回 CN。这里有一个容易被忽略的点Broadcast 的复制成本是随分片数线性增长的。假设驱动表有 10GB集群有 8 个分片那么网络上要传输的数据量就是 80GB。所以 Broadcast Join 的隐含约束非常明确——驱动表必须足够小小到单分片的内存能容纳小到复制 N 份之后网络传输时间可以接受。我测试中明显感受到“小”不只是行数少还包括行宽。一张只有 50 万行但是有 30 个 varchar(255) 字段的宽表序列化后的体积可能比 200 万行的窄表还大广播起来同样吃力。这个路径的优点是不要求驱动表和被驱动表的分片键对齐只要有一张表足够小就能完成分布式 JOIN。因此在真实业务里维度表、配置表、状态表关联大表是 Broadcast Join 最典型的用武之地。缺点也很明确每一次查询都要广播一次而且每个 DN 都要在内存里维护一份驱动表的哈希表并发一高内存压力会非常明显。2.2 Shard Join分片键一致才能“白嫖”本地 JOINShard Join 就不需要 CN 做多余的数据搬移了。它的前提是 SQL 的 JOIN 条件里包含了两张表的分片键等值条件。举个例子如果 customer 表和 orders 表都按 custkey 分片那么执行SELECT * FROM customer c JOIN orders o ON c.c_custkey o.o_custkey;优化器可以判断出custkey 相同的数据一定在同一个分片上也就是说每个 DN 只需要处理本地分片里能 JOIN 上的数据不需要看其他分片的数据。8 个分片就是 8 个独立的本地 JOIN 在并行跑最后把 COUNT 或者命中行数汇总回 CN 就行。为什么说这是“白嫖”因为分布式数据库最贵的两个开销——跨节点网络传输和内存中临时构建哈希表——在这个路径里都被规避了。每个 DN 只需要在本地读数据、本地构建哈希表和单机数据库的 JOIN 几乎没有差别。所以 Shard Join 在分片键对齐的场景下是可以做到接近线性扩展的分片数翻倍单个分片的数据量减半理论上 JOIN 耗时也能接近减半。但这里有一个核心问题分片键对齐是建表那一刻就注定的。如果业务建模时没有把高频 JOIN 键设计成分片键SQL 写得再漂亮优化器也没办法变出一个 Shard Join 来。这也是我在后面场景测试中最想强调的一点——很多慢 JOIN 的根因根本不在 SQL 层而在表结构设计层。2.3 用 EXPLAIN 看懂优化器的选择逻辑PolarDB-X 的 EXPLAIN 输出会以算子树的形式展示执行计划。看懂这棵树的几个关键算子基本就能判断一条 JOIN 走的哪条路。Shard Join 的执行计划形态是这样简化版EXPLAIN SELECT COUNT(*) FROM customer c JOIN orders o ON c.c_custkey o.o_custkey;执行计划大致是HashAgg(count(*)) Gather(concurrenttrue) HashJoin(conditionc.c_custkey o.o_custkey, typeinner) LogicalView(tablescustomer_[0-7], shardCount8, sqlSELECT c_custkey FROM customer ...) LogicalView(tablesorders_[0-7], shardCount8, sqlSELECT o_custkey FROM orders ...)可以看到两个 LogicalView 下面直接是 HashJoin没有任何 Broadcast 或 Shuffle 节点说明 JOIN 被完整下推到了 DN 上执行。Broadcast Join 的执行计划则会多出一个广播节点HashAgg(count(*)) Gather(concurrenttrue) HashJoin(conditionc.c_custkey oi.oi_custkey, typeinner) Broadcast LogicalView(tablescustomer_[0-7], shardCount8, sqlSELECT c_custkey FROM customer ...) LogicalView(tablesorder_items_[0-7], shardCount8, sqlSELECT oi_custkey FROM order_items ...)这里的 Broadcast 节点就是“把小表复制到所有分片”的动作。优化器在决策的时候本质上是在估算三种方案的代价Shard Join 的代价主要来自全表扫描和本地哈希构建Broadcast Join 的代价主要是驱动表扫描加上驱动表体积乘分片数的网络传输Shuffle Join 的代价则是两边都重分布的传输量。谁的估算值小优化器就选谁。提示看到一条慢 SQL 先别急着调优或加索引第一步永远是把执行计划拉出来看一眼。确认了执行路径再去分析瓶颈在网络、内存还是扫描上否则很容易对着错误的环节做无用功。2.4 两种 JOIN 适用场景的横向对比拿一张表来总结这两条路径的差异方便后面读实测数据时对照对比项Broadcast JoinShard Join跨节点数据流动驱动表复制到所有 DN无对驱动表的要求数据量要小、行宽要窄无特殊要求硬性依赖条件无JOIN 键必须等于两表分片键并行度所有参与 DN 并行所有参与 DN 并行网络开销高随分片数放大极低内存开销每个 DN 都要构建小表哈希表只取决于本地数据量最适合的场景维表、配置表关联大表事实表之间按分片键关联最怕的情况驱动表过大或并发过高分片键与 JOIN 键不一致这张表里最重要的信息是Broadcast Join 依赖的只是“一张表足够小”而 Shard Join 依赖的却是“建表那一刻的分片键设计”。前者是执行层面的选择后者是架构层面的设计。一旦分片键设计错了后面所有 SQL 都很难救回来。3. 测试环境、表结构与数据准备3.1 集群配置与版本选择测试环境我尽量贴近生产中等规模部署没有用单机伪分布式因为 JOIN 的跨节点开销必须在真实网络环境下才能体现出来。集群配置如下组件配置数量CN 计算节点8C / 16GB负责 SQL 解析与执行计划1DN 数据节点8C / 32GB负责数据存储与本地计算4总数据分片数每个 DN 承载 2 个分片8 分片软件版本PolarDB-X 2.3.0 开源版MySQL 8.0 协议-网络万兆内网DN 与 CN 同机房-需要说明一下分片数取 8 是一个比较常规的中间值。分片太少并行度不够分片太多广播的复制开销和元数据管理成本又会上升。实际生产环境里分片数一般根据数据量和查询并发共同决定这里取 8 是为了让两种 JOIN 的差异足够清晰又不至于被极端并行度掩盖。3.2 表结构设计故意做出“对齐”和“不对齐”两种情况这次测试的核心是分片键对 JOIN 路径的影响所以建表时就要把场景埋进去。我设计了三张表customer客户维表、orders订单事实表、order_items订单明细事实表。CREATE TABLE customer ( c_custkey BIGINT NOT NULL, c_name VARCHAR(64), c_address VARCHAR(256), c_nation VARCHAR(32), 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_orderdate DATE, o_totalprice DECIMAL(12, 2), o_status CHAR(1), PRIMARY KEY (o_orderkey), KEY idx_custkey (o_custkey) ) PARTITION BY HASH(o_custkey) PARTITIONS 8; CREATE TABLE order_items ( oi_orderkey BIGINT NOT NULL, oi_custkey BIGINT NOT NULL, oi_product VARCHAR(128), oi_quantity INT, oi_price DECIMAL(12, 2), PRIMARY KEY (oi_orderkey) ) PARTITION BY HASH(oi_orderkey) PARTITIONS 8;这里有几个设计要点要展开说一下。第一customer 和 orders 都按用户维度分片因此 customer JOIN orders 时JOIN 键 custkey 正好是两边的分片键天然满足 Shard Join 的条件。这是测试里的“理想对齐”场景。第二order_items 按 orderkey 分片而不是按 custkey 分片。这样 order_items 和 orders 之间按 orderkey 关联时是 Shard Join但 order_items 和 customer 按 custkey 关联时order_items 这一侧没法下推只能走 Broadcast 或者 Shuffle。这就是测试里“不对齐”的场景。第三所有表都使用 HASH 分区且分区数都是 8。如果分区数不一致即使分片键相同做 Shard Join 时也可能出现跨分片的数据搬运所以这里必须保持一致。3.3 造数从 100 万到 1 亿的级联放大数据量设计遵循“小表、大表、超大表”三个层级。customer 表 100 万行orders 表 5000 万行order_items 表 1 亿行。这个量级下customer 大约 120MBorders 大约 24GBorder_items 大约 60GB已经能看出数据搬运和本地计算的量级差异了。造数最忌讳的是顺序插入。如果直接把 customer_id 从 1 顺序写到 5000 万那么所有数据都会落在同一个分片上JOIN 测试就变成单分片压测了。正确做法是让分片键在 8 个分片上均匀分布。我当时的简化方式是先生成全局序列再用固定间隔映射到不同分片或者直接用随机数做代理键。下面是一段临时造数脚本的核心逻辑import random import pymysql conn pymysql.connect(hostcn-1, usertester, password******, dbbench_db) cursor conn.cursor() batch [] for i in range(1, 50_000_001): custkey random.randint(1, 1_000_000) orderkey i batch.append((orderkey, custkey, f2024-{(i % 12) 1:02d}-15, i % 7)) if len(batch) 10000: cursor.executemany( INSERT INTO orders (o_orderkey, o_custkey, o_orderdate, o_status) VALUES (%s, %s, %s, %s), batch, ) batch.clear() conn.commit()实际造数时我开了 12 个并发线程每个线程负责一段连续区间同时往表里灌数据。这样既能把 5000 万行订单在可接受的时间内插完又能保证分片键随机分布。order_items 表的生产逻辑类似只是每条订单再生成 1 到 5 条明细总量控制在 1 亿行左右。3.4 压测方法与指标采集测试不是跑一条 SQL 看一次耗时就行那样偶然性太大。我的做法是先跑两轮预热把数据库的页缓存和优化器统计信息都稳定下来再正式跑 20 轮取 P50 和 P95。并发分别测 1、4、16、32 四档观察不同压力下 RT 和吞吐的变化。辅助工具方面我用了一个简单的 Python 脚本循环执行 SQL记录每轮的耗时并做分位数统计。核心逻辑大概是这样import time import statistics def run_query(sql, cursor, rounds20): costs [] for _ in range(rounds): t0 time.perf_counter() cursor.execute(sql) cursor.fetchall() costs.append((time.perf_counter() - t0) * 1000) costs.sort() p50 costs[int(len(costs) * 0.50)] p95 costs[int(len(costs) * 0.95)] return p50, p95集群资源指标通过 PolarDB-X 自带的监控接口采集重点关注 CN 和 DN 的 CPU 使用率、内存使用量、网络流入流出速率。网络流量这一步很容易被忽略但 Broadcast Join 的性能特征恰恰体现在网络指标上后面场景分析里会用到。4. 场景化 Benchmark五组查询逐一实测4.1 场景一分片键对齐的小表关联大表走 Shard Join第一组测试回归到最经典的星型模型场景客户维表关联订单事实表关联键正好是两边分片键。SELECT COUNT(*) FROM customer c JOIN orders o ON c.c_custkey o.o_custkey;先说结论执行计划完全下推符合 Shard Join 的预期形态。EXPLAIN 输出里customer 和 orders 的 LogicalView 直接挂在 HashJoin 算子下面中间没有 Broadcast 也没有 Shuffle。每个 DN 只需要读本地约 12.5 万行 customer 和 625 万行 orders构建完哈希表后直接关联计数。并发数P50 耗时P95 耗时DN 平均 CPU11.62s1.78s68%42.05s2.31s85%163.84s4.12s96%326.90s7.55s99%单并发 1.62 秒的成绩主要花在最底层的数据扫描上。5000 万行订单分布在 8 个分片每个分片要扫描 625 万行即使走列存或者只读取 o_custkey 这一列I/O 和 CPU 也占了大部分时间。这个数字本身不是重点重点是它作为一个基准线后面所有场景都和它对表。并发从 1 升到 16P50 从 1.62 秒涨到 3.84 秒说明这个 JOIN 在单查询已经接近 CPU 密集型并发上来后只是线性劣化没有出现明显的雪崩。32 并发时 DN CPU 已经打满数据库开始进入排队状态。整体来看Shard Join 路径下的资源消耗非常干净瓶颈就是 DN 的计算能力。4.2 场景二相同语义但 JOIN 键不是分片键优化器选择 Broadcast第二个场景表面上的业务逻辑和场景一一模一样客户表关联明细数据。但这次关联的是 order_items 表而 order_items 的分片键是 oi_orderkey 而不是 oi_custkey所以 custkey 这个 JOIN 键没法直接下推。SELECT COUNT(*) FROM customer c JOIN order_items oi ON c.c_custkey oi.oi_custkey;EXPLAIN 显示优化器选择了 Broadcast Join。customer 表的数据被复制 8 份发送到每个 DN然后每个 DN 再与本地分片里的 order_items 做 Hash Join。因为 customer 只有 100 万行复制 8 份也就是不到 1GB 的网络流量优化器估算下来认为这个代价比把 1 亿行 order_items 重新分布要低得多。并发数P50 耗时P95 耗时DN 平均 CPU网络峰值16.83s7.12s82%约 800MB/s48.74s9.20s95%约 1.1GB/s1623.5s34.8s99%约 1.3GB/s3251.8s78.2s99%接近网卡上限单并发下 6.83 秒比场景一的 1.62 秒慢了 4 倍多。这个差值就是从“本地 JOIN”到“广播 远程 JOIN”的代价差。具体来看每个 DN 需要接收 100 万行 customer 数据构建哈希表然后扫描本地 1250 万行 order_items1 亿行除以 8 分片再做连接。广播带来的额外消耗不只是网络传输还包括每个 DN 都要重复构建同一份哈希表以及 CN 侧等待最慢那个 DN 返回的同步开销。更值得注意的是 16 并发之后的表现。P50 从 8.7 秒直接跳到 23.5 秒P95 更是到了 34.8 秒这已经不是线性劣化了。原因也很明确每一轮查询都要广播 100 万行数据16 个并发就是 16 份广播同时在网络上叠加网络峰值接近万兆网卡上限同时每个 DN 要同时维护 16 份客户表哈希表内存和 CPU 都开始出现争抢。这就是 Broadcast Join 在高并发下的典型瓶颈——它把成本从 IO 转移到了网络和内存而这两者的扩展性比 CPU 还脆弱。4.3 场景三大表 JOIN 大表分片键对齐时依然能 Shard很多人的直觉是“大表 JOIN 大表一定很慢”但实际不完全是。如果两张事实表的分片键刚好对齐Shard Join 依然成立只是每个分片内的本地数据量变大了而已。场景三用 orders 和 order_items 按 orderkey 关联两张表都按 orderkey 分片。SELECT COUNT(*) FROM orders o JOIN order_items oi ON o.o_orderkey oi.oi_orderkey;执行计划同样是完整的 Shard Join 下推。orders 5000 万行order_items 1 亿行各自分布在 8 个分片上每个分片要处理大约 625 万条订单和 1250 万条明细。因为订单和明细是一对多关系JOIN 之后每个分片上的结果行数大约在 1250 万左右最终 COUNT 汇总结果是 1 亿行明细全部命中。并发数P50 耗时P95 耗时DN 平均 CPU19.41s9.88s92%412.6s13.9s99%1628.9s31.5s99%单并发 9.41 秒比场景一的 1.62 秒慢了不少但这很大一部分是数据量本身变大了。重要的是这个查询仍然走的是本地 JOIN网络开销几乎为零瓶颈回到了单分片扫描 1875 万行数据的 CPU 和 I/O 上。也就是说Shard Join 并没有因为“大表 JOIN 大表”而失效它只是在每个分片上承担了更多的本地计算量。如果换成分片键不对齐的大表 JOIN 大表情况会完全不同这就是场景四的内容。4.4 场景四大表 JOIN 大表分片键错位走 Shuffle 参照组我把这个场景设计为“明明业务上经常查但表结构设计不支持”的典型反例orders 按 custkey 分片order_items 按 orderkey 分片两边通过 custkey 关联。SELECT COUNT(*) FROM orders o JOIN order_items oi ON o.o_custkey oi.oi_custkey;EXPLAIN 里会出现 Shuffle 节点因为两边都是大表Broadcast 的代价无法接受广播任何一张表都是 24GB 或 60GB 乘以 8 分片优化器只能选择把两边都按 custkey 重新分布到 8 个 DN再在每个 DN 上做本地 JOIN。并发数P50 耗时P95 耗时网络峰值138.6s41.2s约 1.5GB/s454.3s62.7s接近网卡上限单并发 38.6 秒比场景三的 9.41 秒慢了 4 倍而且这还是理想网络环境下跑出来的结果。整个查询需要把 orders 和 order_items 两张表合计约 1.5 亿行数据全部序列化、传输、反序列化再重分布网络峰值长时间顶着网卡上限。这个场景我特意没有继续往上加并发因为 4 并发时集群网络已经打满再压下去会拖垮同集群的其他业务。它说明了一个非常现实的问题分片键设计错误带来的性能惩罚不是 SQL 层面能轻易弥补的。4.5 场景五三表 JOINBroadcast 和 Shard 组合出现真实业务很少只有两张表关联三表甚至四表 JOIN 才是常态。场景五模拟一个最常见的分析查询客户 JOIN 订单再 JOIN 订单明细同时带订单状态过滤。SELECT COUNT(*) FROM customer c JOIN orders o ON c.c_custkey o.o_custkey JOIN order_items oi ON o.o_orderkey oi.oi_orderkey WHERE o.o_status F;这个 SQL 有意思的地方在于customer 与 orders 按 custkey 对齐走了 Shard Joinorders 与 order_items 又按 orderkey 对齐同样可以下推。整棵执行计划几乎全部压到了 DN 本地执行CN 只负责最后汇总。并发数P50 耗时P95 耗时112.4s13.0s1631.2s35.6s因为加了订单状态过滤需要实际扫描订单数据才能过滤出命中行所以比场景三稍慢。但注意这个查询没有任何一处需要跨节点搬移数据三个表在两次 JOIN 中都利用了分片键对齐这是三表 JOIN 最理想的状态。我还做了一组对比测试把 customer 换成 order_items模拟“先广播一个表再和另一个表 Shard Join”的混合路径SELECT COUNT(*) FROM customer c JOIN order_items oi ON c.c_custkey oi.oi_custkey JOIN orders o ON oi.oi_orderkey o.o_orderkey;执行计划里 customer 先通过 Broadcast 发到所有 DN和本地 order_items 做第一轮 JOIN结果再和 orders 做第二轮 JOIN。单并发耗时 15.7 秒比纯 Shard 的三表 JOIN 慢了约 26%。这说明优化器在混合路径下能做出局部正确的选择但广播导致的额外成本和中间结果放大仍然会对整体耗时产生影响。4.6 并发压测的横向对比把场景一和场景二在不同并发下的 QPS 放在一起看能更直观地感受到两条路径的差异场景1 并发 QPS4 并发 QPS16 并发 QPS场景一Shard0.621.954.17场景二Broadcast0.150.460.68场景一在并发从 1 增加到 16 时QPS 提升了接近 7 倍虽然单查询 RT 在上升但系统整体的吞吐在增长这是典型的可扩展特征。场景二的 QPS 从 0.15 涨到 0.68只有 4 倍多一点的提升而且 16 并发时网络已经接近打满继续加并发只会让 RT 恶化QPS 不会再明显增长。这说明 Broadcast Join 的扩展性上限被网络带宽和 DN 内存锁死了。5. 实测结论什么时候选 Broadcast什么时候选 Shard5.1 核心数据汇总把五组场景的关键数据汇总成一张表方便保存场景执行路径数据规模1 并发 P5016 并发 P50场景一Shard Join100 万 JOIN 5000 万1.62s3.84s场景二Broadcast Join100 万 JOIN 1 亿6.83s23.5s场景三Shard Join5000 万 JOIN 1 亿9.41s28.9s场景四Shuffle Join5000 万 JOIN 1 亿38.6s54.3s场景五混合三表 JOIN12.4s31.2s这里最值得记住的三个结论第一Shard Join 在分片键对齐的前提下比 Broadcast 快 4 倍以上并发越高优势越明显第二Broadcast Join 在小表驱动大表时能跑但驱动表超过一定体量后网络会成为硬瓶颈第三分片键不对齐的大表 JOIN 大表是性能灾难任何 SQL 优化技巧都救不回来只能从表结构设计上解决。5.2 优化器选路的边界在哪里从测试结果反推 PolarDB-X 优化器的决策逻辑基本符合代价模型的选择当驱动表足够小小到广播的网络传输成本低于大表重分布成本时选 Broadcast当两表分片键对齐本地 JOIN 代价最小时选 Shard当两者都不满足时只能走 Shuffle。我额外做了一组递增驱动表大小的测试。把 customer 表从 100 万行逐步扩展到 800 万行重新跑场景二的 SQL观察执行计划什么时候从 Broadcast 切到 Shuffle。结果是在 500 万行左右切换。这个阈值不是固定的它取决于分片数、网络带宽、DN 内存大小和统计信息的准确性。但方向是很清晰的驱动表越大Broadcast 的代价增长越快优化器迟早会切到更“平衡”但也不会更快的方式。所以对于业务侧来说老老实实把大表 JOIN 的分片键对齐比指望优化器替你兜底要可靠得多。5.3 给业务开发的三条实战建议第一建表时优先把高频 JOIN 键设成分片键。如果你有一个订单表和一个用户表且业务里天天做“订单 JOIN 用户”那订单表的分片键就不要用 order_id而应该用 user_id。分片键的设计不是给 DBA 看的它是分布式数据库查询性能的第一决定因素。第二能用广播表就用广播表。PolarDB-X 支持广播表把数据量小、变更不频繁的维表在每个 DN 上都存一份。这样维表 JOIN 大表时连广播传输都省了每个 DN 直接读本地副本就可以。配置表、国家地区表、商品类目表都适合建广播表。第三写完 SQL 养成看 EXPLAIN 的习惯。只要在测试环境执行一次 EXPLAIN就能立刻知道这条 JOIN 是下推了、广播了还是全量重分布了。如果看到 Broadcast 或 Shuffle 节点出现在高频查询里就应该回去审视表结构设计而不是想办法调优一条注定慢的 SQL。6. 常见问题与排查技巧实录6.1 坑 1EXPLAIN 显示 Shard Join线上还是慢我曾经遇到一个现象测试环境执行一条订单关联查询EXPLAIN 确认走的是 Shard Join耗时很理想但上了生产环境同样的 SQL 要跑十几秒。最后定位到根因是生产环境的统计信息过期了优化器对行数的估算偏差很大导致虽然选了下推路径但实际扫描的分片数据里存在严重的数据倾斜——某一个分片上的订单量是其他分片的 10 倍其他 7 个分片都跑完了就等这一个分片。排查这种问题要用 EXPLAIN EXECUTE 拿到实际执行信息重点看每个分片的扫描行数是否均匀。如果发现倾斜第一反应不应该是改 SQL而是去看分片键的取值分布。很多时候是因为某个热客户或热卖家的数据量异常把这个维度的数据单独拆分或者做二级分区就能解决。6.2 坑 2Broadcast Join 把 DN 内存打爆Broadcast 虽然快但有一个隐藏风险驱动表在 DN 上要构建哈希表如果驱动表行宽很大或者并发查询同时广播同一张表DN 内存会快速被吃光。我在压测场景二的时候32 并发跑到一半就出现了 DN Full GC 时间大幅上升随后查询开始报内存不足。应对方式有三个层面。最直接的是控制驱动表大小在 JOIN 之前在子查询里先做过滤和聚合尽可能缩小广播的数据量。其次是降低并发给这类查询单独设置资源组或者限流。最后是调参兜底确认集群的 DN 内存规格足够同时检查统计信息是否准确避免优化器把大表误判成小表从而错误选择广播。6.3 坑 3压测时连接池先被打满做并发压测的时候我一开始遇到的问题是数据库本身没报警但应用侧连接池抛出了 Too many connections。查了一圈才发现是压测脚本每次循环新建连接没有复用。分布式数据库的 CN 要处理连接生命周期新建连接比单机数据库更耗资源。解决办法是压测前先初始化一个固定大小的连接池压测全程复用这组连接。生产环境也一样应用侧连接池的最小连接数和最大连接数要提前规划好否则高并发下连接风暴会先于查询瓶颈出现。这也是一个很不起眼但最容易浪费时间排查的坑。6.4 排查速查表现象可能原因检查方法解决建议JOIN 很慢但执行计划已下推数据倾斜 / 统计信息过期EXPLAIN EXECUTE 看分片扫描行数重新收集统计信息对热点键单独处理EXPLAIN 出现 Broadcast 节点JOIN 键与分片键不一致查看建表 DDL 的分片键定义调整表结构或使用广播表并发一高 DN 内存告警Broadcast 驱动表过大 / 并发广播叠加监控 DN 内存曲线缩小驱动表降低并发升级 DN 规格查询报连接数超限应用连接池配置过小或未复用查看 CN 连接数监控调大连接池设置合理超时回收策略大表 JOIN 大表耗时爆炸分片键不对齐被迫 ShuffleEXPLAIN 看是否有 Shuffle 节点业务上避免此类 JOIN按分片键重设计表7. 说点题外话这次 Benchmark 给我留下的东西跑完这一轮测试我最大的感受是分布式数据库的 JOIN 优化七分在建模三分在 SQL。场景一和场景二的业务语义几乎完全一样都是客户维表关联事实表但就因为 order_items 的分片键设计不同性能差了 4 倍以上到了大表 JOIN 大表的场景分片键错位的后果直接放大到 38 秒对 9 秒。这个差距是任何 SQL 改写、参数调优、语义优化都拉不回来的。所以如果你正在用 PolarDB-X或者刚接触分布式数据库我的建议是每个核心查询在建表之前都先想一遍高频 JOIN 键是谁数据量大概什么级别应该按哪个字段分片。建表时多花十分钟能省下后面无数个排查慢查询的深夜。我自己现在的习惯是任何一条上了生产的大查询都要先在测试环境 EXPLAIN 一把确认没有 Broadcast 和 Shuffle 节点才敢放过去。如果看到 Broadcast我会再想一下是不是能用广播表替代看到 Shuffle就直接回退到表结构设计层面去讨论因为 SQL 层已经无药可救。最后再分享一个小技巧PostgreSQL 和 MySQL 的开发者可能习惯了看执行计划里的“行数估算”来诊断问题在分布式数据库里比行数更重要的是“数据在哪里相遇”。一条 JOIN 的所有优化空间本质上都是在回答这个问题。理解了这个再看 Broadcast Join 和 Shard Join 的取舍就不会觉得它只是一个执行算子层面的选择题了。
返回列表