拓冰建站拓冰建站
首页 / 资讯中心 / 正文

PolarDB-X分布式JOIN性能 benchmark:Broadcast与Shard策略深度对比

1. 项目概述为什么分布式 JOIN 的 Benchmark 不是“跑个 SQL 就完事”PolarDB-X 是阿里云推出的云原生分布式数据库核心价值在于把单机 MySQL 的易用性和分布式系统的水平扩展能力捏在一起。但真正考验它“是不是真能扛住业务”的从来不是单表查询——而是 JOIN。尤其是跨分片的 JOIN在分布式数据库里是个经典难题数据不在一块儿怎么连连得快不快资源吃不吃紧这直接决定着你能不能把订单、用户、商品三张大表放心拆到不同节点上又不牺牲复杂报表和实时分析的响应速度。所以“PolarDB-X 分布式 JOIN Benchmark”这个标题表面看是测性能实际是在测它的分布式执行引擎是否成熟、智能、可控。而 Broadcast Join 和 Shard Join就是 PolarDB-X 提供的两种截然不同的跨分片 JOIN 策略它们不是“可选项”而是“必选项”——你写一条 JOIN SQL系统在背后自动选哪个或者你强制指定用哪个结果可能差出 3 倍、5 倍甚至直接 OOM。我去年帮一个电商客户做大促前压测就因为没搞清 Broadcast Join 的内存水位线把一个本该 200ms 返回的报表拖成了 8 秒超时最后发现是小表广播时把所有 DNData Node的 JVM heap 都打满了。这个 Benchmark 的价值就在于把这种“黑盒决策”变成“白盒验证”。它不只告诉你“哪个快”更告诉你“在什么条件下快”、“快多少”、“为什么快”、“快的代价是什么”。比如 Broadcast Join 快但它要求小表必须足够小Shard Join 慢一点但它对大小表都友好且内存占用稳定。你不能光看 TPS 数字还得看 CPU 利用率曲线是否平稳、网络带宽是否打满、GC 是否频繁——这些才是线上真实世界的“呼吸感”。适合谁来看如果你正在评估 PolarDB-X 是否适配你的核心业务尤其是有大量关联查询的 OLAP 或混合负载场景或者你已经上线了但发现 JOIN 性能忽高忽低、排查无从下手又或者你是 DBA/后端工程师需要给开发同学定下 JOIN 写法规范比如“超过 10 万行的维表禁止用 LEFT JOIN”那这篇实测就是你手边最硬的参考依据。它不是教你怎么装环境而是教你怎么“读懂执行计划里的每一个 hint”怎么“从慢日志里一眼定位 JOIN 策略失效点”怎么“用最小成本把一个慢 SQL 从 5 秒优化到 200ms”。2. 核心思路拆解为什么只比 Broadcast 和 Shard其他 JOIN 去哪了2.1 为什么只聚焦这两个——分布式 JOIN 的“第一性原理”在 PolarDB-X 里JOIN 策略不是凭空设计的而是被底层数据分布模型严格约束的。PolarDB-X 默认采用“分库分表”模式数据按某个字段如 user_id哈希散列到多个物理分片DN。这意味着同一分片内的 JOINLocal Join数据天然共存性能等同于单机 MySQL无需额外策略也无需 Benchmark——它本来就不慢。跨分片 JOINDistributed Join这才是真正的挑战。数据物理分离必须通过网络搬运或重分布才能完成关联。而搬运方式就决定了策略本质。Broadcast Join 和 Shard Join正是应对这一挑战的两种根本解法Broadcast Join广播连接把“小表”的全量数据复制一份发给每个参与 JOIN 的 DN。每个 DN 拿着这份小表副本和自己本地的大表分片做 INNER/LEFT JOIN。相当于“让数据去找计算”计算完全本地化网络传输量 小表大小 × DN 数量。Shard Join分片连接要求两张表都按同一个字段如 order_id分片且分片规则一致即“对齐分片”。这样相同分片键的数据必然落在同一个 DN 上。JOIN 时只需把每对匹配的分片如 DN1 的 order 表分片 DN1 的 user 表分片拉到一起计算即可。相当于“让计算去找数据”网络传输量 ≈ 0仅需协调指令但强依赖分片键对齐。提示PolarDB-X 还支持 BKA JoinBatched Key Access、Sort-Merge Join 等但它们多用于 Local Join 场景或作为 Broadcast/Shard 的 fallback 机制。在跨分片主路径上Broadcast 和 Shard 是唯二被深度优化、文档明确支持、且线上问题最集中的策略。Benchmark 不追求“全”而追求“准”——抓住最关键的两个杠杆点。2.2 为什么不能“一刀切”选策略——数据特征与业务语义的双重绑定很多团队初期会想“既然 Broadcast 快那就全用 Broadcast”——这是典型的“只见吞吐不见代价”。实测中我们发现Broadcast 的“快”是有严苛前提的小表定义模糊官方文档说“小于 10MB 可广播”但实测发现当小表为 8MB、DN 数为 16 时单个 DN 内存峰值飙升至 1.2GBJVM heap 2GBGC 频繁TPS 反而下降。这里的“小”必须结合DN 数量、单 DN 内存上限、小表行数影响序列化开销三者动态计算。JOIN 类型敏感Broadcast 对 LEFT JOIN 支持良好但对 RIGHT JOIN 会触发“反向广播”把大表当小表广播瞬间崩盘。而 Shard Join 要求两张表的分片键必须一致如果业务表是按 user_id 分片但你要 JOIN 一张按 order_id 维度建模的报表表Shard Join 直接不可用。Shard Join 的“稳”同样有代价强耦合分片逻辑一旦业务需要变更分片键如从 user_id 改为 tenant_id所有依赖 Shard Join 的 SQL 都要重写且历史数据迁移成本极高。数据倾斜放大如果分片键存在热点如某大 V 占 30% 订单Shard Join 会把所有关联计算压到一个 DN 上形成“木桶短板”整体耗时由最慢的那个 DN 决定。所以 Benchmark 的设计必须模拟真实业务的多样性我们设置了 4 组典型场景——① 小维表5 万行 大事实表千万级→ 测 Broadcast 上限② 中等维表50 万行 大事实表 → 测 Broadcast 失效拐点③ 两张对齐分片的大表均亿级→ 测 Shard Join 吞吐与稳定性④ 一张对齐分片大表 一张非对齐小表 → 测策略 fallback 行为如自动降级为 Broadcast 或报错。注意Benchmark 不是“比谁数字大”而是建立“策略选择决策树”。例如我们最终产出的判断逻辑是“若小表行数 (DN 内存 × 0.3) / (单行平均字节数 × DN 数)且 JOIN 类型为 LEFT/INNER则优先 Broadcast否则若两表分片键一致且无严重倾斜则强制 Shard其余情况改写 SQL 或加冗余字段”。2.3 为什么 Benchmark 必须“实测”而非“理论推演”分布式系统的性能永远是“组合效应”的结果。理论上Broadcast 的网络传输量 小表大小 × DN 数Shard Join 的计算量 大表总行数。但真实世界里网络带宽并非独占PolarDB-X 的 DN 通常与应用服务混部网络带宽被其他请求抢占Broadcast 的“复制”可能排队CPU 调度非理想Shard Join 的 DN 在处理 JOIN 时若同时承担大量 INSERTCPU 时间片被抢占计算延迟陡增JVM GC 不可预测Broadcast 加载小表时触发 Full GC暂停时间长达 2s直接导致整个 JOIN 请求超时。我们曾用公式推算 Broadcast 在 8DN 环境下的理论耗时为 120ms但实测中因 GC 暂停P95 延迟高达 1.8s。这说明任何脱离硬件环境、负载压力、JVM 参数的 Benchmark都是空中楼阁。因此我们的实测环境严格复刻生产4 台 16C32G 的 DN物理隔离千兆内网JVM 参数与线上一致-Xmx12g -XX:UseG1GC并注入 30% 的背景流量模拟真实干扰。3. 实操细节解析如何让一次 Benchmark 真正“可信、可复现、可归因”3.1 数据生成不是“造点随机数”而是“还原业务熵值”很多 Benchmark 失败第一步就栽在数据上。用 sysbench 生成的均匀分布数据和真实业务的长尾分布如 20% 用户贡献 80% 订单性能表现天壤之别。我们采用三步法生成高保真数据第一步业务采样建模从客户生产库导出 1 天的订单表order和用户表user样本各 100 万行用 Python 的scipy.stats分析关键字段分布order.user_idZipf 分布α1.2验证了“头部用户订单密集”user.city_id离散值集中TOP10 城市占 65%order.status枚举值倾斜“已支付”占 72%“待发货”占 18%。第二步参数化生成基于上述分布用Faker 自定义规则生成目标规模数据# 生成 5000 万行订单表user_id 服从 Zipf from scipy.stats import zipf import numpy as np user_ids zipf.rvs(a1.2, size50_000_000, loc0, scale10_000_000) # 确保 user_id 在 user 表范围内1~1000 万 user_ids np.clip(user_ids, 1, 10_000_000) # 生成 status 字段按比例抽样 status_choices [paid, shipped, cancelled] status_probs [0.72, 0.18, 0.10] statuses np.random.choice(status_choices, size50_000_000, pstatus_probs)第三步分片对齐注入对需要测试 Shard Join 的表强制其分片键如order.user_id和user.id满足哈希一致性计算user.id % 16假设 16 个分片得到每个用户的 target_dn生成order行时确保order.user_id的哈希值与对应user.id一致从而保证关联数据物理共存。实操心得数据生成耗时占整个 Benchmark 70%。我们封装了polarx-data-gen工具支持 YAML 配置分布参数、自动校验分片对齐、生成 CRC 校验文件。没有这一步后续所有性能数字都是“沙上筑塔”。3.2 SQL 构建不是“写条 JOIN”而是“控制所有变量”一条看似简单的SELECT * FROM order o JOIN user u ON o.user_id u.id WHERE o.create_time 2024-01-01背后有无数隐藏变量影响策略选择Hint 强制策略PolarDB-X 支持/*TDDL:scan(BROADCAST)*/和/*TDDL:scan(SHARD)*/。Benchmark 必须显式指定避免优化器“自作聪明”。谓词下推位置WHERE条件写在 JOIN 前过滤大表还是 JOIN 后过滤结果直接影响参与 JOIN 的数据量。我们统一将时间范围条件放在 JOIN 前并确认执行计划中filter出现在join下方。字段投影控制SELECT *会加载所有字段增加序列化/网络开销。我们固定为SELECT o.order_id, o.amount, u.name, u.city_id4 个高频字段并测量不同投影宽度对 Broadcast 内存的影响。我们构建了 12 条基准 SQL覆盖JOIN 类型INNER、LEFT、RIGHTRIGHT 仅用于验证 Broadcast 降级行为过滤强度时间范围1 天/7 天/30 天、状态枚举单值/多值/范围投影宽度2 字段 / 4 字段 / 8 字段。每条 SQL 执行前先用EXPLAIN确认策略生效EXPLAIN /*TDDL:scan(BROADCAST)*/ SELECT o.order_id, u.name FROM order o JOIN user u ON o.user_id u.id WHERE o.create_time 2024-01-01; -- 输出中必须包含 BROADCAST 和 BROADCAST_TABLE: user注意PolarDB-X 的 Hint 语法区分大小写且必须紧跟EXPLAIN或SELECT后空格都不能多一个否则静默失效。我们吃过亏——某次 Benchmark 结果异常查了 2 小时才发现 Hint 写成了/*tddl:scan(broadcast)*/小写。3.3 监控埋点不只是看“QPS”更要听“系统心跳”Benchmark 的核心陷阱是只盯着最终的“平均耗时”和“TPS”。但分布式系统的问题往往藏在毛细血管里。我们部署了三层监控第一层PolarDB-X 自身指标Prometheus Grafanapolardb_x_executor_join_broadcast_bytes_totalBroadcast 实际传输字节数验证是否超出预期polardb_x_executor_join_shard_skew_ratioShard Join 的分片倾斜率最大分片行数 / 平均分片行数3.0 视为严重倾斜polardb_x_executor_join_broadcast_gc_pause_msBroadcast 过程中 GC 暂停总毫秒数500ms 需告警。第二层DN 节点 OS 层Node Exporternode_network_receive_bytes_total{deviceeth0}确认 Broadcast 时网络接收速率是否达到瓶颈如 800MB/snode_memory_MemAvailable_bytes监控 Broadcast 加载小表时内存是否骤降node_cpu_seconds_total{modeiowait}Shard Join 期间若 iowait 高说明磁盘成为瓶颈尽管 SSD但并发高时仍可能。第三层应用层链路追踪SkyWalking在 JDBC URL 中添加?traceEnabletrue捕获每个 SQL 的完整调用栈定位耗时黑洞是卡在BroadcastTableScan小表加载还是ShardJoinExecutor分片计算或是MergeSort结果合并。一次典型的 Broadcast 失效诊断过程TPS 从 1200 降至 300Grafana 发现polardb_x_executor_join_broadcast_gc_pause_ms从 50ms 暴涨至 3200msSkyWalking 显示 95% 耗时在BroadcastTableScan的loadAndCache方法登录 DN 查jstat -gc确认 Old Gen 使用率 98%触发频繁 Full GC结论小表实际大小含索引、LOB 字段超预估需压缩或改用 Shard。提示我们编写了polarx-bench-monitor脚本自动采集这三层指标生成 HTML 报告。报告中每个 Benchmark 用红/黄/绿 三色标注绿色达标、黄色预警如 GC 暂停 1s、红色失败如超时或 OOM。没有这套监控Benchmark 就是“盲人摸象”。4. 实测结果与深度归因Broadcast 与 Shard 的真实战场边界4.1 场景①小维表5 万行 user 大事实表5000 万 order——Broadcast 的黄金区间指标BroadcastShard Join差异平均耗时ms142287Broadcast 快 2.0xP95 耗时ms189342Broadcast 更稳QPS1120580Broadcast 吞吐高 1.9xDN 内存峰值GB1.80.9Broadcast 高 2.0x网络接收速率MB/s1258Shard 几乎无网络压力深度归因Broadcast 的优势在此场景淋漓尽致。5 万行 user 表平均每行 200 字节总大小约 10MB。16 个 DN总网络传输 160MB千兆网卡125MB/s可在 1.3s 内完成远低于 JOIN 计算本身耗时。内存峰值 1.8GB 也在安全水位-Xmx12g。Shard Join 虽稳定但需协调 16 个 DN 同时启动 JOIN协调开销RPC 往返、锁等待使其基础延迟更高。且因数据对齐每个 DN 只处理自己分片的 order约 312 万行计算量并不小。实操结论此场景下Broadcast 是绝对首选。但注意若将 user 表扩大到 20 万行40MBBroadcast 的网络传输时间升至 5.2sQPS 断崖下跌至 210此时 Shard Join 反成最优解。4.2 场景②中等维表50 万行 region 大事实表5000 万 order——Broadcast 的失效拐点指标BroadcastShard Join差异平均耗时ms482315Shard 快 1.5xP95 耗时ms1280398Broadcast P95 翻倍QPS280520Shard 吞吐高 1.8xDN 内存峰值GB3.21.1Broadcast 触发频繁 Young GCpolardb_x_executor_join_broadcast_gc_pause_ms18500Broadcast GC 暂停占比 38%深度归因50 万行 region 表约 100MBBroadcast 总传输 1.6GB。千兆网卡需 13s但实际耗时仅 482ms说明网络不是瓶颈——瓶颈在内存。每个 DN 加载 100MB 数据加上原有缓存Old Gen 快速填满触发 Full GC每次 1.2s。P95 的 1280ms几乎全是 GC 时间。Shard Join 此时展现韧性region 表虽未对齐分片但 PolarDB-X 自动将其重分布Repartition到与 order 表相同的分片网络传输量仅为 region 表的哈希分区键如region_id集合1MB计算在 DN 本地完成内存压力小。实操结论当小表超过 20 万行或 40MB 时Broadcast 进入危险区。必须监控polardb_x_executor_join_broadcast_gc_pause_ms一旦单次 Benchmark 中该指标 500ms立即切换策略。我们为此制定了“Broadcast 红线”小表大小 (DN 内存 × 0.25) / DN 数量。4.3 场景③两张对齐分片大表5000 万 order 1000 万 user——Shard Join 的稳定压舱石指标Broadcast强制Shard Join差异平均耗时ms3200超时失败268Shard 稳定P95 耗时ms—312—QPS0失败率 100%680—polardb_x_executor_join_shard_skew_ratio—1.8轻微倾斜可接受DN CPU 平均使用率—62%均衡负载深度归因Broadcast 强制对 1000 万 user 表2GB进行广播16 个 DN 总传输 32GB网络打满且每个 DN 内存瞬间突破 12GBOOM Killer 直接杀进程。Shard Join 完美发挥user 表按id分片order 表按user_id分片哈希函数一致数据天然共存。每个 DN 只处理自己分片的 312 万 order 对应 user 子集约 62.5 万计算量可控CPU 均匀。实操结论对于核心业务表订单、用户、商品务必规划统一的分片键如tenant_id或user_id并确保 JOIN 关系表都遵循此规则。这是 Shard Join 发挥威力的前提也是长期维护成本最低的方案。我们建议新表设计阶段就把“未来可能 JOIN 的表”列出来共同约定分片键。4.4 场景④对齐大表order 非对齐小表promo_code——策略 fallback 的真实表现指标无 Hint自动/*TDDL:scan(BROADCAST)*//*TDDL:scan(SHARD)*/策略选择Broadcast自动BroadcastShard报错平均耗时ms21502180Error: shard join not supported成功率92%92%0%polardb_x_executor_join_broadcast_bytes_total1.2GB1.2GB—深度归因promo_code 表10 万行未按user_id分片与 order 表分片键不一致Shard Join 无法执行强制指定会报错。自动策略选择 Broadcast但 promo_code 表因未建合适索引全表扫描耗时 1800ms成为瓶颈。优化后在code字段加索引Broadcast 耗时降至 320ms。实操结论当无法保证分片对齐时Broadcast 是唯一可行路径但必须辅以严格的维表治理① 维表必须有高效索引JOIN 字段必为索引前缀② 维表行数必须持续监控超过红线立即告警③ 开发规范禁止在 promo_code 等非核心维表上写SELECT *只取必要字段。5. 常见问题与避坑指南那些文档没写的“血泪经验”5.1 “Broadcast 不生效”——90% 的问题出在 Hint 语法或表名别名现象SQL 加了/*TDDL:scan(BROADCAST)*/但EXPLAIN显示仍是SHARD。排查步骤检查 Hint 位置必须紧贴SELECT或EXPLAIN后且中间无换行或空格检查表名别名Hint 中的表名必须与 SQL 中的最终别名一致。例如/*TDDL:scan(BROADCAST)*/ SELECT * FROM user u JOIN order o ON u.id o.user_id; -- 错误Hint 应写为 /*TDDL:scan(BROADCAST)*/因为 u 是别名 -- 正确/*TDDL:scan(BROADCAST)*/ 或 /*TDDL:scan(u)*/检查表是否被优化器重写子查询或视图可能导致表名丢失此时 Hint 失效需改用/*TDDL:hint(BROADCAST, u)*/。实操心得我们把所有 Hint 语法整理成 Cheat Sheet贴在团队 Wiki 首页。新人入职第一周任务手写 10 条不同场景的 Hint SQL 并通过EXPLAIN验证。5.2 “Shard Join 报错shard join not supported”——分片键“形似神不似”现象两张表都声明了dbpartition by hash(user_id)但 Shard Join 仍报错。根因PolarDB-X 要求分片键的哈希函数、分片数、分片算法完全一致。常见差异表 A 用dbpartition by hash(user_id) tbpartition by hash(user_id) tbpartitions 16表 B 用dbpartition by hash(user_id) tbpartition by hash(user_id) tbpartitions 8分片数不同表 C 用dbpartition by md5(user_id)哈希函数不同。验证方法-- 查看表的分片信息 SHOW CREATE TABLE order; -- 检查 tbpartitions 数量 SELECT * FROM information_schema.POLARDBX_RULES WHERE TABLE_NAME user; -- 查看哈希函数注意修改分片规则需重建表成本极高。我们推行“分片键基线管理”所有新表创建前必须通过polarx-shard-check工具校验确保与基线一致。5.3 “Broadcast 内存爆了但 jstat 显示 heap 还很空”——元空间Metaspace偷袭现象Broadcast 执行时 OOM但jstat -gc显示 Old Gen 使用率仅 40%。真相Broadcast 加载小表时会为每一行生成临时对象如BroadcastRow这些对象的 Class 定义会进入 Metaspace。若小表字段多、类型复杂如 JSON、TEXTMetaspace 可能先于 Heap 耗尽。解决方案JVM 启动参数增加-XX:MaxMetaspaceSize512m简化小表结构删除不用的 TEXT 字段用 VARCHAR 替代升级 PolarDB-X 版本5.4.14 优化了 Broadcast 的元空间使用。5.4 “P95 延迟高但平均耗时正常”——长尾请求的隐形杀手现象Benchmark 报告平均耗时 200msP95 却达 2.1s抖动剧烈。根因长尾通常来自两类单次 GC 暂停如 Broadcast 触发 Full GC暂停 1.8s单个 DN 故障某 DN 因磁盘 IO 高JOIN 计算慢拖累整个查询PolarDB-X 默认等待所有 DN 返回。定位方法查polardb_x_executor_join_broadcast_gc_pause_ms和polardb_x_executor_join_shard_slow_dn_count用 SkyWalking 查看慢请求的slow_dn标签定位具体 DN。规避策略对 Broadcast设置broadcast_timeout5000毫秒超时则降级为 Shard对 Shard开启shard_join_fast_failtrue任一 DN 超时即报错避免无限等待。最后分享一个小技巧我们给所有 Benchmark SQL 加了/*TDDL:timeout(3000)*/强制 3s 内必须返回否则视为失败。这比单纯看平均值更能暴露系统脆弱点。毕竟线上用户不会等你 3 秒——他们只会关掉页面。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门