
上周三晚上快十一点的时候我收到一条线上告警一个报表查询超时了。打开那条SQL一看第一反应是这不就是个简单的分组统计吗数据量也就几千万行怎么就跑不动了。后来用EXPLAIN一查问题根本不在数据量而是整个聚合都在计算节点上完成的数据从存储节点来回拉了一大圈。那一刻我意识到在PolarDB-X这种分布式数据库里SQL性能的关键很多时候不在你写了什么SQL而在优化器把你的算子安排到哪一层去执行——这也就是算子下推的价值所在。这篇文章我想把PolarDB-X的算子与算子下推机制掰开揉碎讲清楚。适合谁看正在用PolarDB-X但查询偶尔变慢的同学做分布式数据库技术选型想了解执行引擎底细的朋友以及面试前突击分布式数据库核心原理的人。我尽量少讲空泛的概念多讲能直接用来排查问题、优化SQL的干货。1. 为什么查个总数都能慢到让人手抖分布式SQL性能的瓶颈根源1.1 算子到底是什么SQL执行计划里的一个个工位先说个基础但很关键的概念算子。很多刚接触PolarDB-X的朋友听到算子这个词总觉得很高深其实它就是一个SQL在执行时被拆分成的每一个具体操作单元。我习惯用快递分拣流水线来类比。一条SQL到了数据库里不是一口气算完的而是被拆成很多小步骤。去找某个表的数据这是TableScan算子相当于从仓库货架拿包裹按条件筛掉不要的行这是Filter算子相当于按重量或目的地初筛把相同维度的数据汇总这是Aggregate算子相当于把发往同一个城市的包裹归堆按某个字段排序这是Sort算子相当于按区域码排好顺序把两张表的数据按条件拼在一起这是Join算子相当于把订单信息和物流信息合并成一张完整单据。这些工位在单机MySQL里也存在只不过单机场景你不太感知得到它们的存在。但在PolarDB-X这种分布式架构里算子的位置变成了核心问题每个工位到底放在哪台机器上干活全部放在计算节点CN统一处理还是下沉到每个存储节点DN本地处理这就是下推讨论的范畴。1.2 一条SQL在PolarDB-X里的旅程CN、DN与GMS各管什么要理解下推先得知道PolarDB-X的执行架构。它有三类核心组件GMS负责全局元数据管理你可以理解成档案室CN是计算节点负责接收SQL、做语法解析和优化、生成执行计划并调度执行DN是存储节点负责真实数据的存储和本地计算本质上每个DN就是一个MySQL实例。一条SQL到达CN之后大致走这么几步先做语法解析生成逻辑执行计划然后优化器基于统计信息、数据分布等生成物理执行计划这个物理执行计划就是一棵由算子组成的树最后执行器按照这棵树去调度。整个过程里有个关键决策贯穿始终树上的每个算子它的输入数据在哪里它自己又该在哪里执行。这就要说到下推这两个字的字面意思了。原本算子的默认执行位置是CN但CN手里并没有数据它得从DN把数据拉过来。如果把算子从CN往下挪到DN让数据在自己所在的地方先完成加工只把加工后的结果返回给CN这就叫下推。所以算子决定做什么下推决定在哪儿做这两件事合起来基本决定了一条分布式SQL的性能上限。1.3 先算一笔账把数据全拉回中心处理到底有多贵为什么说执行位置这么重要核心原因是分布式环境里网络传输是最贵的资源。我经常跟团队的人说一句话CPU跑得快、内存也不慢但跨越节点搬数据这件事永远是分布式系统的阿喀琉斯之踵。来算一笔最直观的账。假设订单表有1亿行分布在32个分片上每个分片大概312万行。现在执行SELECT COUNT(*) FROM orders。如果聚合算子不下推CN得把1亿行原始数据全部通过网络拉到计算节点内存里再逐行计数。按照平均一行100字节估算这就是大约10GB的网络传输量在万兆网络下也得花十几秒起步而且CN的内存压力会非常大。但如果聚合下推到DN每个分片只需要在本地MySQL里把COUNT数出来向CN返回一个数字。32个分片返回32个数字加起来不超过1KB。同一个查询一个可能是秒级返回另一个直接干到分钟级甚至超时。这就是为什么网上一聊到PolarDB-X性能绕不开算子下推。它不是某个SQL的优化技巧而是整个执行引擎设计的核心逻辑之一。2. 算子下推的本质让数据在自己家门口完成加工而不是把原料全运回中央工厂2.1 数据本地性原理区域仓先加工总仓只做总装下推背后的理论根基是数据本地性Data Locality意思是计算应当尽量靠近数据所在的位置。用物流行业来类比最合适如果全国各地的包裹都要先运到北京总部再分拣、再发出去那整个网络的吞吐都会被北京总部卡死。真实物流体系的做法是每个区域分拨中心先把当地的包裹做初筛、归类、打包只把需要跨区流转的包裹发到总部。PolarDB-X的DN就是这样的区域分拨中心。数据物理上就存放在DN机器上DN自身的MySQL实例完全有能力独立完成过滤、聚合、排序、甚至一部分Join操作。优化器只要判断某个算子在DN上执行不会改变最终结果就会优先把它下沉。CN则退居总装车间的角色把各DN汇上来的中间结果做最后的归并和修正。这里有个容易被忽略的点数据本地性不只是省了网络流量还省了序列化和反序列化的开销。数据从DN存储引擎里读出来在本地内存里直接参与计算不需要转成网络传输格式再还原一次。这种隐性收益在线上海量查询场景下累积起来非常可观。2.2 下推的完整链条从优化器决策到每个分片上的子SQL下推不是执行器在执行阶段临时决定的而是在优化器生成物理执行计划时就定好了。我大致梳理一下这个链条第一步CN做语法解析和逻辑计划生成这个阶段还没有算子的位置概念。第二步优化器拿到表的元数据、分片信息、统计信息开始评估每个算子放到哪一层执行代价最小它能利用的信息包括分区键哪个分片存了什么范围的数据、DN的计算能力、网络带宽估测、表的行数基数等。第三步优化器把可以在DN上执行的算子下沉生成一个包含下推决策的物理计划。第四步执行器把下推部分转换成一段段针对单分片的子SQL派发给对应的DN执行。这个过程很关键的一个副产品是逻辑视图LogicalView。你在EXPLAIN结果里看到LogicalView节点就说明这部分扫描和计算被下推了——它代表每个分片上要执行的那段子SQL。所以判断一条SQL有没有下推不靠猜看执行计划里的节点类型和Extra字段就行这一点我后面专门展开。2.3 什么情况下下推收益最大用数量级思维预估优化空间我在实际排查慢SQL时很少第一时间去抠某条SQL的每个细节而是先用数量级思维判断最多能省多少。看三个问题的答案基本就能锁定优化方向第一WHERE条件里有没有能裁剪分片的字段如果条件包含分区键的等值或范围优化器可以直接跳过大量无关分片这比下推更前置是更根本的省。第二聚合或排序这类算子能不能让DN先做一遍再在CN汇总比如按状态分组统计订单数32个分片各自先GROUP BY返回的可能只有几十个分组这比拉几千万行原始数据回来再算要快几个数量级。第三两个表做Join时能不能保证每个分片内独立完成连接不需要跨节点搬数据这个稍后专门讲。有一个很实用的经验如果一个慢查询的执行计划里CN层出现的算子输入行数是所有分片扫描行数之和也就是几千万上亿的规模那么这个SQL基本没有吃到下推红利。目标应该是让CN层的算子输入行数缩小到分组数级别或者TopN条数级别。用这个标准去衡量一眼就能看出SQL还有没有救。3. 可下推与不可下推的算子边界一张清单看清存储层能力3.1 存储节点能独立完成的算子扫描、过滤、投影、聚合、排序、Limit并不是所有算子都能无脑下沉到DN但好消息是常用算子里的多数都可以。我先把能下推的列出来配上各自的应用场景。TableScan也就是LogicalView是必下推的。数据在DN上扫描天然发生在DN这个没有选择余地。Filter下推的收益非常直观一张表1亿行WHERE statusPAID只留下1%Filter在DN执行完CN拿到的就是过滤后的100万行传输量直接砍到1%。如果不下推1亿行原始数据先到CN再在CN里过滤前面那笔账就白算了。Project投影也值得说。很多开发习惯SELECT *如果宽表有100列业务实际只用3列投影下推可以让DN只返回3列未取用的列不占带宽不占内存。别小看这个优化列数从100降到3意味着同样一个查询网络传输量缩小三十几倍。聚合和排序属于部分下推的典型场景。比如GROUP BY user_idDN每个分片先把局部聚合做完返回每个user_id的局部汇总值CN再做一次最终合并。Sort和LIMIT也是同理比如取每个分片价格最高的前10条DN先各自排序取TOP10CN做归并排序后取最终TOP10。至于Limit这里有个语义细节要注意如果是LIMIT 10DN各取前10再在CN归并取前10结果是对的。但如果SQL里出现LIMIT 10 OFFSET 100000这种深度分页下推的意义就大打折扣了因为每分片截断会丢失全局偏移信息这类查询大概率还是要拉大量数据到CN重新处理。3.2 不能下推的算子语义正确性比性能更重要有些算子不能下推不是存储层不会做而是做了之后结果可能不对。数据库优化器有一条铁律任何优化不能改变SQL的结果语义。破坏语义的解释一下就清楚。第一类是跨越分片依赖的DISTINCT。如果去重字段的相同值可能分布在不同分片DN本地去重只能打掉本分片内部的重复跨分片的重复还得CN层二次去重。所以这类DISTINCT在CN层至少还留一个节点属于部分下推。第二类是窗口函数例如ROW_NUMBER() OVER(PARTITION BY user_id ORDER BY create_time)。窗口函数通常需要在整个分区的范围内做计算如果分区数据的全量分散在多个DN上DN本地做的窗口计算就不完整只能把数据拉到CN统一处理。第三类是非确定性函数这是最容易让人踩坑的点。NOW()、RAND()、UUID()这类函数每次调用的结果都可能不同。如果下推到DN执行32个分片各执行一次产生的32个值互相不一致比如WHERE create_time NOW()各分片用的是不同时间点最终结果就乱套了。优化器对这类函数非常保守。第四类是用户自定义函数UDF。CN上注册的UDF不一定存在于DN实例中如果DN没有这个函数定义想推也推不下去。我见过有人在WHERE条件里用自定义函数做字符串清洗结果这个自定义函数无法下推全表数据被拉到CN过滤变成了CN层的逐行调用那个查询慢得惨不忍睹。3.3 部分下推不是失败理解聚合的二次加工模型很多时候看到EXPLAIN里CN层还挂着聚合节点就以为下推失败了其实这是误解。聚合下推的完整形态本来就是局部聚合最终合并的两段式结构。打个比方总部分公司统计全国销售总额不是让每家门店把每一笔小票都寄到总部而是每家门店先算出自己的销售额报一个数字上来总部把各门店数字相加。这就是MapReduce里的Combiner思想也是分布式SQL引擎处理分布式聚合的标准做法。所以在执行计划里你经常能看到DN侧的Partial Aggregation和CN侧的Final Aggregation。只要DN侧把几亿行压缩成了几十个分组CN侧的最终聚合就是瞬时完成的事。判断一个查询有没有吃到下推红利关键是看CN层算子的输入规模而不是看CN层还有没有算子。CN层可以有算子但不能有海量数据输入。4. Join算子的下推策略广播、重分布与同分片就地连接4.1 三种Join策略的原理与选择逻辑Join是所有算子里最复杂的一个也是下推收益最容易体现、最容易翻车的地方。两个表的数据分散在不同的分片上要做到每一对匹配的行都能碰面只有三条路可以选。第一条路两表按相同维度分片且Join键恰好是分片键。这样匹配的行从一开始就住在同一个DN上每个DN完全可以在本地完成连接这就是同分片Join也叫Collocated Join。它是代价最低的Join方式因为根本不需要跨节点搬数据。第二条路一张大表Join一张小表。把小表完整广播到所有DN每个DN拿小表的副本和本地的大表分片做连接这叫Broadcast Join。代价是大表分片数乘以小表大小适合小表体量很小的场景。第三条路Join键和分片键不匹配或者无法复用现有分布。这时只能让两张表都按照Join键重新分发到对应节点再在节点上做本地连接这就是Shuffle Join也叫重分布Join。它需要把至少一张表的全部相关数据在网络里搬一遍是最贵的方案。优化器的选择逻辑不复杂就是分别估算三种策略的网络传输量、内存消耗和执行时间然后选一个总代价最小的。但正是因为估算依赖统计信息统计信息不准确时它也会犯错这点后面展开说。4.2 同分片Join是最便宜的连接方式但前提是设计阶段就想明白先说说最理想的Collocated Join。假设orders表按order_id分片order_items表也按order_id分片查询条件是orders.order_id order_items.order_id。由于两张表的分片键相同同一个order_id对应的订单数据行和订单明细行必然落在同一个DN的分片上。这样每个DN就可以只读本地数据完成Join全程不需要任何网络数据搬移。这种查询跑起来非常快快的原因不是DN的MySQL有多彪悍而是它压根没有产生分布式系统最怕的Shuffle流量。所以我经常在业务建模阶段就跟同事强调那些高频Join的表分片键一定要设计成同一种业务键。订单和订单明细按order_id分片用户和用户扩展信息按user_id分片会员和会员等级按member_id分片——这比事后优化SQL管用一万倍。注意一个容易误判的细节同分片Join要求两表的分片方式和分区键完全一致。如果orders按order_id分片order_items却按item_id分片即使Join条件是order_id优化器也不能用同分片Join必须先做一次数据重分布这就瞬间贵了起来。4.3 广播Join的适用条件与隐藏成本当两表不能同分片时Broadcast Join是第二优选。我到目前为止遇到最多的是大表和小表字典表的Join。比如1亿行订单表Join一个只有几千行的shop店铺维度表务必要走广播让每个DN拿一份shop表副本本地Join效果极好。但小表的定义很微妙。代价公式是小表大小乘以分片数等于需要传输的总流量。一张1MB的表、32个分片传播量32MB无所谓可如果所谓的小表有500MB32个分片就是16GB这个代价一下子就上去了。优化器通常有阈值判断但阈值参数不一定贴合你的实际网络带宽所以看到执行计划里出现Broadcast Join时千万别只看有没有下推要顺手估算一下广播的数据量级。还有一种情况小表其实可以反复广播给同一DN。如果每个DN上需要这份数据又恰好没有广播就是刚性需求。为了避免高频广播有些团队会把字典表做成复制表也就是在每个DN上都冗余一份完整数据。这是用存储换时间的经典做法在PolarDB-X的分布式表设计里也非常实用。4.4 Join下推失效时的性能落差一个让我印象深刻的案例讲一个真实又典型的排查经历。线上有个运营后台的查询要查每个用户最近3笔订单及用户姓名SQL大概是users表Join orders表Join条件是user_id。users表按user_id分片没问题但orders表是按order_id分片的。结果就是orders表里所有涉及到的行必须按user_id重新Shuffle一遍1亿行级别的数据在网络上搬家查询跑了两百多秒。排查时的核心证据是执行计划里出现了明显的重分布节点CN层挂了一个很大的Shuffle Hash Join。后来我们把思路从优化这条SQL转到了优化表设计为这个高频场景单独维护一张按user_id分片的订单冗余明细表业务写入时同步写入。改完之后这个查询从同分片Join直接变成Collocated Join执行时间从两百秒降到了几百毫秒。这个案例最大的价值是提醒我Join下推失效很多时候不是优化器笨而是表设计时压根没给优化器留出下推的机会。与其在SQL里加HINT硬凑不如在一开始就让数据分布跟业务访问模式对齐。5. 用EXPLAIN验证下推是否生效别信感觉看执行计划5.1 EXPLAIN输出里哪些字段是下推信号PolarDB-X兼容MySQL生态所以用EXPLAIN看执行计划是排查问题的第一手段。拿到执行计划之后我先干三件事第一看有没有LogicalView节点它代表下推给DN的逻辑视图第二看LogicalView的Extra字段里写了哪些Push down信息第三数一数CN层执行节点的输入量级。Extra字段里常见的下推信号包括Push down filter表示过滤条件下推到了存储层Push down aggregation表示聚合下推了Push down order和Push down limit表示排序和Limit下推了。如果某个算子该下推却没有对应关键词那它就是在CN层执行的此时要警惕输入数据量是否过大。这里要强调一个容易混淆的点LogicalView本身只代表扫描被下推不代表所有计算都被下推。真正决定性能的是LogicalView返回的数据已经加工到多细的程度。它返回的是几亿行原始数据还是几十个GROUP BY分组这才是下推深度的问题。5.2 两个典型执行计划对照一次下推成功一次只有扫描下推为了说清楚怎么判断我用两个简化的执行计划示例做对照。第一个查询SELECT order_id, amount FROM orders WHERE order_id 12345。执行计划里看到LogicalViewExtra里写着Push down filter。这说明WHERE条件被下推到DNDN扫描并过滤后只把命中行返回CN。这个查询的单分片裁剪也生效了几乎可以忽略不计地快。第二个查询SELECT user_id, COUNT() FROM orders WHERE status PAID GROUP BY user_id。执行计划里CN层有一个HashAggregate节点输入来自LogicalView。注意CN有聚合不代表下推失败。关键是DN侧的SQL实际是SELECT user_id, COUNT() FROM orders WHERE statusPAID GROUP BY user_id也就是DN先把几千万行压缩成了几百个分组CN只需要把各分片的几百个分组汇总。输入从千万行级别降到百行级别这就是一次很健康的部分下推。真正的失败案例长什么样如果逻辑视图的Extra里只有扫描而CN层出现了一个Filter算子输入行数是全表行数那说明过滤条件没有下推所有原始数据都拉到CN了。这类查询在数据量上来后必然变慢。判断工具除了EXPLAIN有条件的话还可以用EXPLAIN ANALYZE看实际执行统计包括每个算子实际处理的行数和耗时。我当时排查线上问题就是靠它定位到某一个算子吞掉了几乎所有时间再倒推数据为什么那么大。5.3 反推销技巧如何从DN侧看到真实下推的SQL执行计划是优化器的说法如果你想看真实发生的事还有一招到DN上去观察实际执行的SQL。第一种做法在DN侧开启通用日志或慢日志过滤出来自CN的连接执行的语句。你会发现下推过来的SQL往往带着分片条件或注释标识看到它带着GROUP BY、ORDER BY这些计算就说明下推实打实发生了。第二种做法用SHOW PROCESSLIST观察活跃会话。慢查询执行时到DN上看有没有正在跑、跑很久的SQL如果DN侧干干净净没有动静那说明大部分计算都发生在CN侧数据搬运被隐藏在计划里了。这在排障时特别直观也比反复琢磨执行计划图更贴近事实。我在一次排查慢查询时就是靠DN日志发现下推SQL里有一个函数根本没有出现在应用的原始SQL里——那是优化器为了保持语义自动补上的而那个函数恰好不能下推导致整个查询卡在CN。这种问题只看CN侧的计划看不到必须到DN反推。6. 下推失效的五个经典场景与排查经验6.1 分区键被函数包裹裁剪失效第一个坑我认为是最常见的在WHERE条件里对分区键套了函数。比如分区键是create_time查询写成WHERE DATE(create_time) 2024-01-01。优化器面对DATE函数时无法求逆没法把这个条件推导成分区裁剪的范围于是只能扫所有分片。正确的写法应该是WHERE create_time 2024-01-01 00:00:00 AND create_time 2024-01-02 00:00:00。两种写法业务效果相同但第二种能让优化器直接把范围映射到对应的分片过滤体量完全不同。排查方法很简单EXPLAIN看partitions列如果显示的是全部分片说明裁剪失效了。这是个只要写SQL稍不注意就会掉进去的坑而且数据量越大越致命。6.2 非确定性函数把过滤条件锁在CN层第二个坑也是我实际遇到过并且调了一整天的非确定性函数导致的无法下推。某位同事写了个WHERE RAND() 0.1想从订单表抽10%的样本。这个条件如果下推每个DN各自执行一次RAND()各分片抽出来的样本可能完全不对等结果的随机语义会被破坏。优化器很聪明地选择不下推但代价是全部行先搬到CN再逐行抽样一次抽样查询把数据库拖得很惨。类似还有NOW()。同一个事务或同一条SQL里NOW()的值要求保持一致下推到多个DN就会出现各分片取到不同时间点的问题。所以这类条件要么改写法比如先在外面算好时间变量再用常量传入要么接受它无法下推的代价。判断思路很简单凡是每次调用结果可能不同的函数都要怀疑它能否下推。6.3 Join条件与分片键对不上数据被迫全量重分布第三个坑也是前面4.4的正式总结Join键和分片键不一致导致的全量Shuffle。当两个表的高频Join条件不是表的分片键时优化器只能按Join键重新分布数据。这时所有的行都要在网络里流动一遍性能断崖式下跌。排查方法继续靠执行计划如果看到Shuffle、Exchange或者大行数的Join节点基本可以判定发生了数据重分布。治理方案有三个层级第一层业务建模时就把高频Join键设计成分片键第二层增加冗余表或冗余字段避免跨表Join第三层实在不行就用广播Join兜底前提是小表足够小。每次都从SQL和HINT层面想办法是治标不治本。6.4 排序规则不一致导致的执行计划伪下推第四个坑比较隐蔽中了的人多半一头雾水字符集排序规则不一致导致下推SQL和CN层执行结果存在一致性风险优化器在某些版本下会变得更保守。举例来说某列使用utf8mb4_general_ci排序规则但DN实例的参数配置和CN默认行为不一致。优化器在做下推决策时要做一次语义等价判断这个排序条件下推到DN执行和留在CN执行结果一样吗一旦它发现表达式里涉及自定义函数或者规则不确定就会放弃下推宁可把数据拉回来执行保证语义正确。遇到这类情况排查思路是把下推SQL单独拿到DN上跑一遍和CN上的结果对比不一致就检查排序规则、时区、字段字符集这些环境参数。时区不一致也会影响NOW()等时间函数的取值进而影响范围过滤的分区裁剪结果。这类坑通常查起来慢但一旦定位到环境参数问题修复也很快。6.5 统计信息缺失让优化器选错Join策略最后一个坑来自优化器本身统计信息不准或缺失导致代价估算失真。最典型的症状是大表和小表做Join时优化器放弃广播而选择了全量Shuffle。原因可能是小表的统计信息显示它很大或者大表实际行数和统计信息差了好几个量级。解决办法是主动收集统计信息。我在运维实践中会把ANALYZE TABLE加进定时任务数据大批量导入之后也手动跑一遍然后重新EXPLAIN对比计划是否变化。如果确认统计信息没问题依然选错策略才轮到HINT强制指定方案。但HINT是最后一根稻草不该是首选方案因为它等于绕过优化器一旦数据分布变化硬编码的策略可能会更糟。我的排查顺序通常是这样先EXPLAIN看裁剪和下推深度再EXPLAIN ANALYZE定位瓶颈算子然后ANALYZE TABLE刷新统计信息重看计划最后才考虑HINT和表设计改动。按照这个顺序走绝大多数下推失效的问题都能找到原因。最后分享一个这些年养成的习惯所有核心查询在发布前我都要求团队把EXPLAIN结果存一份到文档里重点检查LogicalView的Extra字段、partitions裁剪列、CN层算子输入量级这三项。一旦之后线上性能劣化翻出当时的执行计划一对比问题基本能锁定在下推失效还是数据量变化上。这比临时抓日志、猜原因要高效得多。算子下推这件事说到底是优化器在做一道在哪里计算最划算的决策题而你作为使用者至少要学会看懂它的答案。