PolarDB-X分布式JOIN性能优化实战:Broadcast与Shard选型指南
1. 项目概述为什么分布式 JOIN 的 Benchmark 不是“跑个 SQL 就完事”PolarDB-X 是我过去三年在金融和电商类中大型客户项目里用得最频繁的分布式数据库之一。它不是简单把 MySQL 拆开扔到多台机器上而是真正把计算层和存储层做了逻辑分离让 JOIN 这种原本在单机数据库里“习以为常”的操作在分布式环境下变成了一个需要精心设计、反复权衡的系统级工程问题。很多人一上来就问“Broadcast Join 和 Shard Join 到底哪个快”——这个问题本身就有陷阱。就像问“螺丝刀和电钻哪个更好用”答案永远取决于你要拧的是木板上的自攻钉还是混凝土墙里的膨胀螺栓。这次实测的核心目的从来不是为了贴出一张“XX 场景下 Broadcast 快 3.2 倍”的截图而是要搞清楚在真实业务数据模型下当你的订单表shard key 是 order_id和用户表shard key 是 user_id必须关联时PolarDB-X 底层到底会怎么选它选对了没有如果选错了你作为 DBA 或后端工程师有没有办法干预、引导甚至强制它走更优路径这才是 Benchmark 的价值所在。我这次测试覆盖了从 100 万行小表广播到 5000 万行大表分片关联的全量场景所有数据都来自脱敏后的某电商平台真实订单日志与用户画像快照不是 synthetically generated 的理想化数据。测试环境也完全复刻了客户生产集群的典型配置3 个 CN计算节点8 个 DN数据节点每个 DN 配置 16 核 64G网络带宽限制在 10Gbps这是很多混合云环境的真实瓶颈。所以这篇内容适合正在做 PolarDB-X 选型评估的架构师、刚接手线上慢 JOIN 问题的 DBA以及写业务代码时发现“明明加了索引却还是 10 秒才出结果”的后端同学。它不讲抽象理论只讲你在命令行里敲下EXPLAIN后看到的那几行 Plan Node 到底意味着什么以及你下一步该改哪一行 SQL、调哪个 Hint、动哪个系统变量。2. 核心思路拆解Benchmark 不是比速度是比“可控性”与“可解释性”2.1 为什么不能直接套用 TPC-H 的测试方法TPC-H 是一个被广泛引用的基准测试集但它本质上是一个“为测试而生”的封闭模型。它的 schema 固定8 张表、数据生成规则严格scale factor 决定数据量、查询语句预设22 条标准 SQL。这在学术对比或厂商宣传稿里很有用但在真实世界里它几乎毫无指导意义。举个最典型的例子TPC-H 的lineitem表和orders表天然就按orderkey关联而这个orderkey在 PolarDB-X 里恰恰就是orders表的分片键shard key。这意味着绝大多数 JOIN 查询PolarDB-X 可以直接下推到单个 DN 上执行根本不需要跨节点数据传输。这种“天然是最优”的场景完全掩盖了分布式 JOIN 最核心的痛点当关联字段不是分片键时系统如何决策决策依据是否可靠你能否信任它所以我的整个 Benchmark 设计第一原则就是“制造冲突”。我把orders表按order_id分片但故意让user_info表按user_id分片然后构造一条SELECT * FROM orders o JOIN user_info u ON o.user_id u.user_id的查询。这个user_id对orders表来说只是一个普通二级索引字段对user_info表才是分片键。这就逼出了 Broadcast Join 和 Shard Join 的真实博弈场。2.2 Broadcast Join 的本质不是“广播”是“复制 本地 JOIN”很多人听到 Broadcast Join第一反应是“把小表发给所有节点”这没错但只说对了一半。更准确地说Broadcast Join 的执行流程是CN 节点先将小表比如user_info表假设只有 50 万行的全量数据通过网络发送给每一个参与本次查询的 DN 节点每个 DN 节点收到这份“副本”后再用自己的本地orders分片数据与这份副本做一次标准的 MySQL 内存 JOIN通常是 Hash Join最后CN 汇总所有 DN 的结果去重并返回给客户端。这里的关键在于“副本”的生命周期。它不是永久存在 DN 上的而是本次查询专用的临时缓存。这就带来两个硬性约束第一小表的数据量必须足够小否则网络传输时间会成为瓶颈。我实测过当user_info表超过 200 万行约 1.2GB 原始数据时光是广播阶段就耗时 800ms此时即使 JOIN 本身很快整体延迟也已经失控。第二CN 节点必须能准确预估小表大小。PolarDB-X 默认通过information_schema中的TABLE_ROWS和AVG_ROW_LENGTH来估算但这个值在大表上往往严重不准MySQL 的统计信息本身就是采样估算。所以如果你的user_info表有 1000 万行但TABLE_ROWS显示只有 200 万PolarDB-X 就会误判为“小表”强行触发 Broadcast结果就是 CN 疯狂往 DN 发送 1000 万行数据网络打满DN OOM整个集群雪崩。这就是为什么实测中我必须手动ANALYZE TABLE并验证TABLE_ROWS的准确性而不是依赖默认值。2.3 Shard Join 的本质不是“分片”是“重分布 协同 JOIN”Shard Join 的目标是让两个大表都在各自的数据节点上完成 JOIN避免大规模数据移动。它的实现原理是CN 节点分析两个表的分片键和 JOIN 条件如果发现无法利用现有分片进行本地 JOIN比如orders.order_idvsuser_info.user_id就会启动一个“重分布”Repartition过程。具体来说CN 会下发指令让orders表的所有分片根据user_id字段的哈希值重新路由shuffle到一组新的 DN 上同时也让user_info表的所有分片根据user_id字段的哈希值路由到同一组 DN 上。最终每个目标 DN 上都拥有orders和user_info的一部分数据且这部分数据的user_id哈希值是匹配的于是就可以在本地完成 JOIN。这个过程听起来很美但代价巨大。首先重分布本身就是一个高网络、高 CPU、高磁盘 IO 的操作。我用tcpdump抓包观察过一次 5000 万行orders表的重分布会产生超过 120GB 的网络流量因为每行数据都要序列化、哈希、发送、反序列化。其次它要求两个表的user_id字段数据分布必须相对均匀。如果user_info表里 80% 的用户都来自同一个城市比如city_id1那么重分布后承载city_id1数据的 DN 就会成为热点CPU 持续 100%其他 DN 却在空转。我在测试中特意构造了这种“长尾分布”的user_id结果 Shard Join 的 P99 延迟直接从 1.2s 拉长到 8.7s而 Broadcast Join 因为只广播小表反而更稳定。所以Shard Join 的性能高度依赖于数据的“可分片性”而不是简单的“数据量大小”。2.4 Benchmark 的核心指标不只是 QPS 和 Latency在分布式系统里只看平均响应时间Latency和每秒查询数QPS是危险的。我定义了四个必须监控的核心指标网络吞吐Network Throughput用iftop -P 3306实时监控 CN 和 DN 之间的流量。Broadcast Join 的峰值流量应该集中在 CN→DN 方向且总量 ≈ 小表大小 × DN 数量Shard Join 的流量则应该是双向、均衡、且总量远超前者的。DN CPU 利用率分布CPU Skew用top -H -p dn_pid查看每个 DN 的线程 CPU 占用。理想的 Shard Join所有 DN 的 CPU 应该在 60%-80% 之间波动如果出现一个 DN 100%、其他 DN 20% 的情况说明数据倾斜严重。内存峰值Memory Peak用pmap -x dn_pid在查询执行中抓取。Broadcast Join 的 DN 内存峰值应该等于其本地orders分片大小 广播来的user_info全量大小Shard Join 的内存峰值则是重分布缓冲区 本地 JOIN 的 Hash Table 大小通常高出 3-5 倍。Plan Stability执行计划稳定性连续执行同一条 SQL 100 次用EXPLAIN检查 Plan Node 是否每次都一样。如果出现 30% 的概率走 Broadcast70% 的概率走 Shard那就说明优化器的代价模型在摇摆这种不确定性在生产环境里比慢 1 秒更可怕。这四个指标任何一个异常都比单纯的“平均耗时 200ms”更能揭示问题的本质。这也是为什么我的测试脚本里sysbench只负责发压真正的数据采集全部由自研的polardb-x-profiler工具完成它能精确到毫秒级地记录每一次网络包、每一个线程栈、每一MB内存分配。3. 实操细节与关键参数从建表到压测每一步都是坑3.1 建表语句分片策略是性能的起点不是终点很多人以为只要在建表时指定了DBPARTITION BY和TBPARTITION BYPolarDB-X 就会自动搞定一切。这是最大的误解。分片策略的选择直接决定了后续 JOIN 的成本天花板。以下是我为本次测试编写的、经过反复验证的建表语句-- orders 表按 order_id 分片这是业务主键无法更改 CREATE TABLE orders ( order_id bigint(20) NOT NULL COMMENT 订单ID, user_id bigint(20) NOT NULL COMMENT 用户ID, amount decimal(10,2) DEFAULT 0.00 COMMENT 订单金额, create_time datetime DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (order_id), KEY idx_user_id (user_id) ) DBPARTITION BY HASH(order_id) TBPARTITION BY HASH(order_id) TBPARTITIONS 8; -- user_info 表这里有个关键选择——按 user_id 分片还是按 city_id 分片 -- 我最终选择了按 user_id因为 JOIN 条件是 o.user_id u.user_id -- 如果按 city_id 分片那么 JOIN 就必须走 Broadcast没有其他路可选 CREATE TABLE user_info ( user_id bigint(20) NOT NULL COMMENT 用户ID, city_id int(11) DEFAULT 0 COMMENT 城市ID, age tinyint(4) DEFAULT 0 COMMENT 年龄, gender char(1) DEFAULT U COMMENT 性别, PRIMARY KEY (user_id) ) DBPARTITION BY HASH(user_id) TBPARTITION BY HASH(user_id) TBPARTITIONS 8;这里有一个极易被忽略的细节TBPARTITIONS 8。这个数字不是随便写的。它代表每个 DN 上这张表会被再切分成 8 个物理分片Table Partition。PolarDB-X 的 JOIN 优化器在做 Broadcast 决策时会先估算“小表的总行数”然后再除以TBPARTITIONS得到每个物理分片的平均行数再乘以 DN 数量得出总的广播数据量。如果我把TBPARTITIONS设为 1那么优化器会认为每个 DN 上的user_info分片就是整张表从而大概率放弃 Broadcast。而设为 8它就会认为每个 DN 只需要处理 1/8 的数据大大提高了 Broadcast 的触发概率。我在测试中对比过TBPARTITIONS 1和TBPARTITIONS 8前者在 100 万行user_info时优化器始终选择 Shard Join后者在同样数据量下95% 的概率选择 Broadcast。所以TBPARTITIONS是一个隐性的、强大的性能杠杆它不改变数据分布却能显著影响优化器的决策路径。3.2 系统变量调优broadcast_row_count不是唯一开关PolarDB-X 提供了一个非常直观的变量broadcast_row_count用来设置 Broadcast Join 的行数阈值。默认值是 100 万。很多人的做法是把user_info表的行数一查发现是 80 万就放心大胆地写 SQL等着它自动广播。但现实远比这复杂。这个变量只是优化器决策链中的一个环节。它前面还有broadcast_memory_limit广播内存上限默认 1GB后面还有broadcast_timeout广播超时默认 30s。这三个变量是串联关系优化器先检查user_info表的估算行数是否 ≤broadcast_row_count如果是再检查估算的广播数据大小行数 × 平均行宽是否 ≤broadcast_memory_limit如果都满足最后还要确保在broadcast_timeout内能把数据发完。我在一次测试中user_info表估算行数是 95 万小于 100 万但AVG_ROW_LENGTH被低估为 50 字节实际是 120 字节导致估算数据大小为 47.5MB远小于 1GB看起来没问题。结果执行时CN 发现实际要发送 114MB超过了broadcast_memory_limit的软限制于是立刻回退到 Shard Join。所以正确的调优顺序是先ANALYZE TABLE user_info确保统计信息准确再用SELECT AVG(LENGTH(CONCAT(user_id, city_id, age, gender))) FROM user_info计算真实平均行宽最后把broadcast_memory_limit设置为估算行数 × 真实平均行宽 × 1.2留 20% 余量。我最终的配置是SET GLOBAL broadcast_row_count 2000000; SET GLOBAL broadcast_memory_limit 256*1024*1024; -- 256MB SET GLOBAL broadcast_timeout 60;这个配置让我在user_info表达到 180 万行时依然能稳定触发 Broadcast Join而不会在 100 万行时就因内存估算偏差而失败。3.3 SQL 编写技巧Hint 是你的最后一道防线当优化器的自动决策不可信时Hint 就是你手里的“手术刀”。PolarDB-X 支持两种核心 Hint/* BROADCAST(table_name) */和/* SHARD(table_name) */。但它们的使用有严格的语法和时机要求。首先Hint 必须写在SELECT关键字之后、第一个表名之前且不能换行。错误的写法-- ❌ 错误Hint 写在了 WHERE 后面无效 SELECT * FROM orders o JOIN user_info u ON o.user_id u.user_id /* BROADCAST(u) */; -- ❌ 错误Hint 和 SELECT 之间有空行解析器会跳过 SELECT /* BROADCAST(u) */ * FROM orders o JOIN user_info u ON o.user_id u.user_id;正确的写法-- ✅ 正确紧贴 SELECT无空行 SELECT /* BROADCAST(u) */ * FROM orders o JOIN user_info u ON o.user_id u.user_id; -- ✅ 正确可以同时指定多个表但必须是 JOIN 中的小表 SELECT /* BROADCAST(u) BROADCAST(p) */ * FROM orders o JOIN user_info u ON o.user_id u.user_id JOIN product p ON o.product_id p.product_id;更重要的是Hint 的作用域。BROADCAST(u)只强制user_info表被广播但orders表的分片方式不变。如果orders表本身数据量极大比如 5 亿行而user_info只有 50 万行那么BROADCAST(u)是安全的但如果orders表只有 10 万行user_info有 300 万行你强制BROADCAST(u)就会导致 CN 往每个 DN 发送 300 万行数据网络瞬间打满。所以Hint 不是万能的它必须配合你对两张表数据量的精确掌握。我在生产环境里只在两种情况下用 Hint第一EXPLAIN显示优化器错误地选择了 Shard Join而我知道小表绝对够小第二业务上这个 JOIN 是高频核心路径必须保证 Plan 稳定不能有任何摇摆。除此之外我一律禁用 Hint优先通过调优统计信息和系统变量来解决问题因为 Hint 是“硬编码”一旦数据量增长它就会变成技术债。3.4 压测脚本sysbench的魔改与EXPLAIN的自动化解析标准的sysbench oltp_read_only只能测试单表查询对 JOIN 完全无能为力。所以我基于sysbench的 Lua 脚本框架魔改了一个专门用于分布式 JOIN 的压测工具。核心改动有三点第一SQL 模板化。我定义了一个join_template.lua文件里面包含 5 个不同复杂度的 JOIN 场景-- 场景1简单两表关联无 WHERE 过滤 local query1 SELECT /* BROADCAST(u) */ COUNT(*) FROM orders o JOIN user_info u ON o.user_id u.user_id -- 场景2三表关联其中一张是维度表必须广播 local query2 SELECT /* BROADCAST(u) BROADCAST(c) */ o.order_id, u.name, c.city_name FROM orders o JOIN user_info u ON o.user_id u.user_id JOIN city c ON u.city_id c.city_id -- 场景3带范围过滤的关联考验优化器的谓词下推能力 local query3 SELECT /* BROADCAST(u) */ * FROM orders o JOIN user_info u ON o.user_id u.user_id WHERE o.create_time 2023-01-01第二EXPLAIN自动注入。在每次mysql:query()执行前我插入一行mysql:query(EXPLAIN FORMATJSON .. sql)并将返回的 JSON Plan 解析成结构化数据提取出plan_node、typeBROADCAST_HASH_JOIN 或 SHARD_HASH_JOIN、est_rows估算行数、memory_used估算内存等关键字段写入到独立的plan_log.csv文件中。这样一次压测跑完我不仅有 QPS 和 Latency 曲线还有一份完整的、按时间戳排序的执行计划变迁日志。第三动态数据生成。标准sysbench的oltp_common无法生成符合user_id和order_id关联关系的数据。所以我写了一个 Python 脚本gen_join_data.py它先生成 1000 万个唯一的user_id再为每个user_id生成 1-50 个随机的order_id确保orders表里的user_id字段100% 存在于user_info表中。这样测试数据就具备了真实的业务关联性而不是随机生成的“脏数据”。这套魔改脚本让我能在 2 小时内完成从 10 万行到 5000 万行的全量性能扫描并自动生成一份包含 200 个数据点的 Excel 报告其中 X 轴是user_info表行数Y 轴是 P95 Latency每条曲线代表一种配置组合如broadcast_row_count100wvsbroadcast_row_count200w。没有这套自动化靠人肉EXPLAIN一百次早就崩溃了。4. 性能实测结果与深度分析数据不会说谎但需要你读懂它4.1 基准线测试100 万行user_info表下的表现我们先建立一个清晰的基线。测试环境orders表固定为 1000 万行user_info表从 10 万行逐步增加到 100 万行所有系统变量保持默认broadcast_row_count1000000。以下是关键数据user_info行数优化器选择P95 Latency (ms)CN→DN 网络峰值 (MB/s)DN CPU 最大值 (%)Plan 稳定性100,000Broadcast1824278100%300,000Broadcast21512582100%500,000Broadcast26820885100%800,000Broadcast3423338895%1,000,000Shard1256185099 (1个DN)80%这个表格揭示了第一个残酷真相Broadcast Join 的性能衰减是平滑的而 Shard Join 的性能崩塌是陡峭的。当user_info达到 80 万行时Broadcast 的延迟只比 10 万行时增加了 87%看起来还能接受但一旦突破 100 万行的阈值优化器切换到 Shard Join延迟直接飙升至 1256ms是 Broadcast 的 6.7 倍。更致命的是 Plan 稳定性从 100% 降到 80%意味着每执行 5 次查询就有 1 次会意外走 Shard导致用户体验断崖式下跌。网络峰值也印证了这一点Broadcast 时CN→DN 流量是可控的、线性的而 Shard 时双向流量暴增且集中在少数 DN 上CPU 打满。这说明对于这张user_info表100 万行就是 Broadcast 的“甜蜜点”上限。超过这个点就必须考虑其他方案比如提前物化视图或者重构分片键。4.2 数据倾斜测试当user_id不再均匀分布真实世界的用户数据从来不是均匀的。我用gen_join_data.py构造了一个极端案例user_info表共 100 万行但其中 80 万行的city_id1代表北上广深剩下 20 万行分散在其他 100 个城市。然后我再次运行相同的 JOIN 查询。结果如下优化器选择P95 Latency (ms)DN CPU 分布 (Top 3)内存峰值 (GB)备注Broadcast29582%, 79%, 76%1.8稳定无倾斜Shard482099%, 42%, 38%12.4一个 DN 成为热点内存暴涨Shard Join 的延迟从 1256ms 暴涨到 4820ms增长了近 4 倍EXPLAIN的 JSON Plan 显示SHARD_HASH_JOIN节点的est_rows估算为 100 万但实际执行时承载city_id1数据的 DN收到了 80 万行user_info和对应 400 万行orders而其他 DN 只收到了不到 1 万行。这就是数据倾斜的典型症状优化器的代价模型是基于“均匀分布”这个完美假设的一旦现实打破这个假设它的所有估算都会失效。而 Broadcast Join 完全不受此影响因为它把全量小表发给每个 DN每个 DN 处理的orders分片数据量是固定的与user_info的分布无关。这个测试给我一个明确的结论如果你的关联字段如user_id存在明显的长尾分布Broadcast Join 是更鲁棒的选择哪怕它的理论数据量稍大。稳定性有时候比绝对的峰值性能更重要。4.3 多表 JOIN 测试维度表的“广播链”效应在真实报表场景中一个查询往往涉及 4-5 张表。比如orders→user_info用户维度→product商品维度→category类目维度。product和category表通常都很小10 万行但user_info是中等规模100-500 万行。我测试了三种策略策略 A全 Broadcast/* BROADCAST(u) BROADCAST(p) BROADCAST(c) */策略 B混合/* BROADCAST(p) BROADCAST(c) */让user_info由优化器自动选择策略 C全 Shard不加任何 Hint完全依赖优化器。结果令人惊讶策略 AP95 Latency 312ms但 CN 内存峰值达 4.2GB且EXPLAIN显示 Plan Node 非常臃肿有 12 个BROADCAST_HASH_JOIN节点。策略 BP95 Latency 285msCN 内存峰值 2.1GBPlan Node 清晰只有 3 个BROADCAST_HASH_JOIN针对p和cu表走 Shard。策略 CP95 Latency 1890msPlan 不稳定有时走全 Shard有时在中间某步错误地广播了user_info。这说明在多表 JOIN 中盲目地给所有小表加BROADCASTHint反而会拖累整体性能。因为 CN 节点需要同时管理多份广播数据的内存和网络调度开销呈指数级增长。最优解是“精准打击”只对那些绝对小、绝对稳定、且与主事实表关联紧密的维度表如product、category使用BROADCAST而对于像user_info这样的“中等表”应该交给优化器或者通过broadcast_row_count精细调控。这背后是一个经典的“分治”思想把最复杂的、最容易出错的决策user_info的广播与否交给系统把最简单、最确定的决策category表只有 1000 行留给自己。4.4 生产环境迁移建议从“能跑”到“稳跑”的三步走基于以上所有实测我给正在规划 PolarDB-X 迁移的团队总结了一套落地性极强的三步走建议第一步建模期——用EXPLAIN做静态审查。在应用上线前把所有核心 SQL尤其是报表、导出、后台任务相关的 JOIN全部用EXPLAIN FORMATJSON跑一遍。重点关注plan_node字段是否为BROADCAST_HASH_JOIN或SHARD_HASH_JOIN以及est_rows是否在合理范围内比如user_info表显示est_rows5000但实际有 500 万行这就是统计信息严重失真必须ANALYZE。这一步能发现 80% 的潜在性能雷区。第二步灰度期——用slow_query_log做动态监控。上线后开启 PolarDB-X 的慢查询日志并配置long_query_time11 秒。每天凌晨用脚本自动解析日志提取出所有执行时间 1s 的 SQL再对它们执行EXPLAIN检查 Plan 是否发生了变更。如果发现某条昨天还是 Broadcast 的 SQL今天变成了 Shard就要立刻排查是user_info表数据量增长了还是ANALYZE没有定期执行还是某个UPDATE导致了数据分布变化第三步稳态期——用hint做兜底保障。对于那些已经被验证为高频、核心、且 Plan 必须稳定的 SQL果断加上/* BROADCAST(table_name) */Hint。但要记住Hint 是“止血贴”不是“创可贴”。它解决的是当下的稳定性问题长期来看还是要回到第一步持续优化统计信息和分片策略让系统自己就能做出正确决策。我在一个客户的生产环境里就是用这套方法把一个原来 P95 延迟 3.2s 的核心报表查询优化到了稳定的 220ms且连续 30 天 Plan 零变更。这三步不是教科书式的理论而是我在三个不同客户现场踩着坑、流着汗、熬着夜一点点摸索出来的实战路径。它不追求“一步登天”而是强调“步步为营”每一步都有明确的检查点和交付物让团队能真正掌控住分布式 JOIN 这个最棘手的环节。5. 常见问题与独家避坑指南那些文档里不会写的细节5.1 “为什么我加了BROADCAST(u)EXPLAIN却显示SHARD_HASH_JOIN”这是最常被问到的问题。原因几乎总是同一个Hint 的语法位置错误或者表别名没写对。PolarDB-X 的 Hint 解析器非常严格。假设你的 SQL 是SELECT /* BROADCAST(u) */ * FROM orders o JOIN user_info u ON o.user_id u.user_id;这里u是user_info的别名Hint 里写的BROADCAST(u)是正确的。但如果写成BROADCAST(user_info)就会失效因为优化器找不到名为user_info的表它只认识别名u。更隐蔽的错误是空格和换行。/* BROADCAST(u) */和SELECT之间如果有任何空白字符包括空行、tab、空格解析器就会跳过这个 Hint。我曾经在一个客户的环境里调试了整整一天最后发现是开发同学在SELECT和/*之间不小心按了一个空格键。解决方案很简单把 SQL 复制到一个纯文本编辑器如 Notepad打开“显示所有字符”功能确保SELECT和/*是紧挨着的中间没有任何不可见字符。另外可以用SHOW WARNINGS命令如果 Hint 被忽略这里通常会有一条Query cache is disabled之类的提示这是个误导性提示实际是 Hint 解析失败。5.2 “ANALYZE TABLE后TABLE_ROWS变了但EXPLAIN的est_rows没变为什么”这是因为 PolarDB-X 的优化器会缓存EXPLAIN的结果。ANALYZE TABLE更新的是底层 MySQL 的统计信息但 CN 节点的 Plan Cache 里可能还存着旧的、基于错误统计信息生成的执行计划。解决方案有两个第一执行FLUSH PLAN CACHE;命令强制清空 CN 的 Plan Cache第二更推荐的做法是在ANALYZE TABLE后立即执行一次EXPLAIN让它重新生成并缓存新计划。我写了一个一键脚本#!/bin/bash mysql -h cn_host -P 3306 -u user -ppass -e ANALYZE TABLE user_info; mysql -h cn_host -P 3306 -u user -ppass -e FLUSH PLAN CACHE; mysql -h cn_host -P 3306 -u user -ppass -e EXPLAIN SELECT * FROM orders o JOIN user_info u ON o.user_id u.user_id\G这个脚本我放在了每个客户的定时任务里每天凌晨 2 点自动执行确保统计信息永远是最新的。5.3 “Broadcast Join 时DN 内存 OOM 了但broadcast_memory_limit明明设得很大”这是一个非常隐蔽的坑。broadcast_memory_limit控制的是 CN 节点广播数据时的内存上限但它不控制 DN 节点接收并构建 Hash Table 的内存。DN 节点的内存是由 MySQL 自身的join_buffer_size和tmp_table_size参数决定的。当