Apache DataFusion 5.0.0 版本深度解析:Join 语义、窗口函数、Parquet 剪枝与执行器重构

原创2026-09-24 23:32:161,445 阅读
文章标签:大数据数据分析后端

Apache DataFusion 5.0.0 版本深度解析:Join 语义、窗口函数、Parquet 剪枝与执行器重构

Apache DataFusion 5.0.0 发布于 2021-08-10,是该查询引擎在快速演进期的第一个大版本:它一口气引入了完整的窗口函数体系(row_number / rank / dense_rank / lead / lag / first_value 等)、SELECT DISTINCT 与反连接/半连接支持,重写了 Join 条件解析与 Parquet 行组剪枝逻辑,并完成了 MergeExec → CoalescePartitionsExec 的重命名。读完本文,你将掌握 5.0.0 引入的每个破坏性变更与关键新能力,理解其底层实现原理(结合当前仓库源码佐证),并能正确迁移到新 API。

变更日志的完整原始条目位于仓库的 dev/changelog/5.0.0.md,本文以该文档为主线,按「破坏性变更 → 新能力 → Bug 修复 → 性能优化」四个维度展开,并在每个主题下补充当前仓库中的源码级证据。

一、破坏性变更(Breaking Changes):迁移前必须知道的事

5.0.0 是一次「语义修正 + 大规模重构」并存的版本。从 dev/changelog/5.0.0.md 可以看到,官方列出了 22 条 breaking changes,可以归纳为五类:

1.1 Join 语义与列处理逻辑修正

  • JOIN 条件顺序相关(#778):join 条件不再被视为对称的,ON / USING 约束下的连接列处理逻辑被重写(#605),同时修复了「Qualified field resolution too strict」(#810)与「Better join order resolution logic」(#797)两个问题。对用户而言,多表 join 时列限定符(表名前缀)的解析变得更宽松也更准确。
  • 支持非等值 join 条件(#660):5.0.0 之前 join 条件中只允许等值谓词,现在非等值过滤条件(如 <、>)可以作为 join 条件的一部分,无法下推的 unsupported 条件会被转换为 filter 保留(#796)。从当前仓库的物理计划实现看,Hash Join 只处理等值连接键,其余条件落在 hash_join 执行器 之外的过滤节点上,与当时的拆分思路一脉相承。

1.2 执行器与规划器重命名、API 调整

  • MergeExec 更名为 CoalescePartitionsExec(#635):这一更名消除了「Merge」一词与 Sort Merge Join、SortPreservingMerge 的歧义。当前仓库中该执行器位于 coalesce_partitions.rs,其职责保持不变:并行执行多个分区输入,再合并成一个分区输出(不保证结果顺序,可通过 with_fetch 提前终止输出)。它在计划树中承担「多分区结果汇拢」的角色,例如 ORDER BY 全局排序之上必然接一个 CoalescePartitionsExec,这一结构在 dataframe 集成测试 的 EXPLAIN 输出中可以看到。
  • 扩展规划 API 纳入逻辑计划(#643):扩展节点的物理规划签名需要接收完整的逻辑计划,自定义数据源 / 自定义算子的接入方需要同步更新。
  • SQLMetric 改为原子操作(#25):指标实现从 Mutex 保护的共享状态改为原子计数(AtomicU64 等),降低了热路径上的锁竞争;name 字段被移除,Metrics 的构造方式随之变化。
  • PredicateBuilder 返回 Vec<bool>(#370):剪枝谓词构建的返回值从闭包 Fn 改为显式布尔向量,简化了行组剪枝的调用链。
  • RowGroupPredicateBuilder 重构为 PruningPredicateBuilder(#365):为 5.0.0 的剪枝重写(见 2.4 节)铺路。

1.3 表达式与类型系统调整

  • ScalarValue::List 改为 Box 存储(#788):List 变体的元素由直接内嵌改为装箱,官方说明「reduce size by half size」,大幅降低了 ScalarValue 在内存中的占用——在列式引擎中 ScalarValue 会被大量复制进聚合哈希表,这一改动对内存敏感场景收益明显。
  • log 语义修正为 log10,并新增 ln(#271):此前 log(x) 的实现实际上是自然对数(或与惯例不符),5.0.0 将其修正为以 10 为底的对数,并显式提供 ln(x)。如果你的 SQL 里依赖 log 的旧行为,迁移时务必检查。

1.4 依赖与数据源

  • 切换到 crates.io 上的 arrow-rs 4.x(#395):不再依赖 GitHub SHA 快照,构建系统从「跟随 arrow 主干」变为「锁定发布版」,后续升级到 arrow 5.0(#721)也更平滑。
  • 新增 NdJson(JSON Lines)数据源支持(#404):datafusion-cli 同步支持读取 ndjson(#427),为当时 JSON 数据接入提供了统一入口。

1.5 其他值得注意的变更

  • EXPLAIN VERBOSE 展示所有优化器 pass 的结果(#759):调试查询计划时可直接看到每个优化规则前后的变化,配合 #744(EXPLAIN 展示优化后的物理/逻辑计划)与 #751(limit push down 的 explain verbose 测试)一起使用效果更佳。
  • 支持带限定符的列名(#55):SELECT t.a FROM t 这类写法全面可用。
  • 支持 SELECT DISTINCT(#262):见 2.2 节。
  • CSV 可从 stdin 或内存读取(#54)。
  • now() 函数支持(#288):见 2.3 节。

二、重大新能力(Implemented Enhancements)

5.0.0 的增强清单规模庞大(从 dev/changelog/5.0.0.md 可以看到 90 余条),这里按功能域组织成六组。

2.1 窗口函数:从零到完整的体系化落地

这是 5.0.0 最重量级的特性,经历了一个清晰的渐进过程:

  1. 框架搭建(#334):先建立 window expression 的逻辑/物理规划与 EXPLAIN 结构,仅支持空 OVER ();
  2. order by 与 partition by 子句(#463、#501、#520、#558):逐步加入排序分区语义;
  3. 窗口帧(window frame)支持(#506、#530):加入 ROWS/RANGE 帧构造、边界检查与错误处理;
  4. 内置窗口函数全集(#403、#429、#631、#687):first_value、last_value、nth_value、lead、lag(支持 offset 与 default 参数)、row_number、rank、dense_rank;
  5. 性能与正确性(#569、#571、#595):窗口函数使用 repartition 并行化、ORDER BY 折叠进窗口表达式在逻辑层完成排序、find_ranges_in_range 优化。

当前仓库中这些函数沉淀在 functions-window crate,例如 row_number、rank、dense_rank 分别对应 row_number.rs、rank.rs 中的 RowNumber、Rank 等结构体,nth_value 对应 nth_value.rs 的 NthValue。当时的调用示例(现在仍然有效):

SELECT
  id,
  row_number() OVER (ORDER BY score DESC) AS rn,
  rank()       OVER (ORDER BY score DESC) AS rk,
  dense_rank() OVER (ORDER BY score DESC) AS drk,
  lead(score, 1, 0) OVER (PARTITION BY dept ORDER BY score) AS next_score
FROM students;

2.2 SELECT DISTINCT 与聚合能力补齐

  • SELECT DISTINCT(#262):查询级去重正式可用;随后 #250 的 SELECT DISTINCT 答案错误被 #268 的「wrong projection optimization」修复,两者共同保证了正确性。当前 SELECT DISTINCT 在优化器中被展开为 Aggregate(group by 所有投影列),再经过 distinct 专用优化规则处理。
  • COUNT DISTINCT 全面铺开:布尔(#230)、浮点(#252)、字典数组(#256)、时间戳(#319)相继支持,配合 #812 的 DictionaryArray 向量化哈希,聚合的输入类型覆盖接近完备。
  • Hash 分区聚合(#320):实现 hash-partitioned aggregation,解决 #107「高基数 GROUP BY 跑不完」的问题,是后续分布式聚合(Ballista #634)的基础。
  • GROUP BY 细节:支持按列位置分组(#519)、解析 group by 表达式别名(#485)、修复 GROUP BY NULL 的正确答案(#793),并移除 GroupByScalar 改用 ScalarValue 以支持 null 分组键(#786)。

2.3 时间与数学函数

  • now()(#288):与 PostgreSQL 兼容的当前时间函数。当前实现位于 datetime/now.rs,其 Signature 被声明为 Volatility::Stable(稳定函数)——这意味着 now() 在查询计划开始时取值,同一查询内无论在哪里执行都返回同一时间戳,这是可预测性设计的关键;返回值时区由会话时区(datafusion.execution.time_zone)决定,默认 None。
  • date_trunc()(#203):按粒度截断时间。当前实现位于 datetime/date_trunc.rs,支持 microsecond / millisecond / second / minute / hour / day / week / month / quarter / year 等粒度,并校验非法粒度(#L55-L78)。示例:
> SELECT date_trunc('month', '2024-05-15T10:30:00');
-- 2024-05-01T00:00:00
> SELECT date_trunc('hour', '2024-05-15T10:30:00');
-- 2024-05-15T10:00:00
  • 时间戳转换三件套(#567):to_timestamp_millis()、to_timestamp_micros()、to_timestamp_seconds(),并在规划阶段常量折叠优化(#387)。
  • 数学函数(#309、#716):一元数学表达式改用 arrow 的 unary kernel 实现并修复返回类型冲突;random 函数导出(#303、#389)。

2.4 Parquet 剪枝:基于 PruningStatistics 的彻底重写

5.0.0 对 Parquet 行组剪枝做了重构级改动:

  • 核心重写(#426):剪枝逻辑整体改用 PruningStatistics trait + arrow Array trait 表达,这是后续所有剪枝能力的地基。当前仓库中 PruningPredicate 的文档化示例清晰说明了其工作原理(见 pruning_predicate.rs):给定 x = 5 的谓词与容器 A {x_min=0, x_max=4}、B {x_min=2, x_max=10}、C {x_min=5, x_max=8},剪枝谓词可证明 A 中不可能有行满足条件,从而跳过整个容器;B、C 保留读取。
  • 布尔列剪枝(#490、#500):boolean 类型列首次进入剪枝范围。
  • 不等值谓词剪枝(#544、#561):!= 谓词可参与剪枝,并修复了 != 剪枝错误。
  • Date 剪枝修复(#690):修复 Date32 / Date64 的 parquet 行组剪枝。
  • 去限定符(#689、#619、#617):移除下推谓词上的限定符以修复剪枝,同时避免剪掉带非限定引用的必要列。
  • ExecutionConfig 开关(#749、#764):新增选项用于启用/禁用 parquet 剪枝(#723 需求),随后将剪枝规则限制在简单表达式上以避免误剪(#764),并配套剪枝禁用测试(#754)。
  • 端到端验证(#657):新增 parquet 剪枝端到端测试,并为 ParquetExec 增加 metrics,可直接观察跳过的行组数量。

2.5 Join 家族:反连接、半连接与跨连接

  • Semi Join / Anti Join(#470、#482):SEMI / ANTI 两种连接类型落地,补齐了 Join 语义拼图。
  • Cross Join(#11):纯笛卡尔积实现(此 PR 虽在合并列表,但 5.0.0 汇总中标注为 cross join implementation)。
  • Hash Join 内部优化(#24、#827):优化内部结构并修复 null 处理;使用 RawTable API 提升 hash join 性能;hash 函数抽离至 hash_utils.rs(#807),浮点数组按无符号整数哈希(#556)。
  • 可扩展分布式 Join(#634):Ballista 实现可扩展分布式 join,并移除硬编码的 PartitionMode(#637)。

2.6 执行器、DataFrame 与 CLI 增强

  • SortPreservingMergeExec(#362、#378、#379):新增保序归并算子,支持多分区的 SortExec;随后 #722、#691 持续优化其物化性能。它在当前物理计划库中位于 sorts/sort_preserving_merge.rs,并配套 fuzz 测试(merge_fuzz.rs)保证多输入流合并的正确性。
  • RepartitionExec 指标(#397、#398):为 repartition 增加 SQL metrics;#521、#576 修复其错误传播与输出挂断问题。
  • HashJoinExec / ParquetExec 指标(#657、#664):执行器级 metrics 体系成型,benchmark 可在查询结束后打印带指标的执行计划(#396、#662)。
  • DataFrame 流式 collect(#789):新增流式版本的 collect 方法,避免大数据集全量物化。
  • datafusion-cli(#284、#285、#289、#292、#295、#296、#323、#427):参数校验、文件参数、csv/tsv/json 打印格式、-- 注释、--quiet/-q 静默模式与计时开关、ndjson 支持相继加入;#514 修复执行时间显示,Docker 镜像体积从 2.16GB 降至 89.9MB(#266)。

三、关键 Bug 修复解读

5.0.0 修复了 40 余个 bug(见 dev/changelog/5.0.0.md),以下是最值得关注的几类:

3.1 Join 正确性

  • LEFT/FULL JOIN 多批数据错误(#238、#845):修复右侧 0 批或多批时左连接的实现错误,以及左侧存在多个不匹配行时 right/full join 的处理——这类问题在列式批量处理引擎中非常隐蔽,属于典型的批次边界 bug。
  • Join schema 重复字段(#311、#601):修复 join 时 schema 出现重复非限定字段名导致的 panic,以及 join datatypes 的 panic。
  • Join 条件顺序依赖(#778):正式文档化为 breaking change,要求 join 条件按固定顺序书写(见 1.1 节)。

3.2 剪枝与优化器正确性

  • 投影下推误删列(#617、#619):投影下推会删除「虽非限定但仍被使用」的列,5.0.0 通过 RFC 级修复(#619)停止剪除带非限定引用的列。
  • Filter 被错误移除(#225):修复 where 子句中无列名的过滤条件被优化 pass 删除的问题。
  • 重复过滤条件去重(#409、#436):c>5 AND c>5 折叠为 c>5,同时避免重复添加已有过滤条件。

3.3 窗口与聚合

  • 窗口函数别名与区分(#454、#592、#607、#622):窗口聚合支持别名;未命名窗口函数按 partition/order by 子句区分;窗口函数内支持别名。
  • 聚合参数错误提示(#505):给出正确的聚合参数错误信息,改善可诊断性。

3.4 数据读取与类型

  • 多字节列名(#357):SQL planner 支持多字节(非 ASCII)列名。
  • 空表 SELECT *(#613):修复空表上 select * 的问题。
  • COUNT DISTINCT 布尔/时间戳(#230、#314、#319):类型覆盖补全,修复 Unexpected DataType for list 报错。
  • CSV/Parquet 扫描表名(#629):Ballista 计划序列化时尊重 csv/parquet scan 的表名。

四、性能优化盘点

性能类 PR 集中在 dev/changelog/5.0.0.md 与 Merged PR 列表,归纳如下:

领域 PR 手段
in-list 求值 #813 字符串与原生类型的 inlist 加速
SortPreservingMerge #691、#722 物化与整体性能优化
min/max 聚合 #719 使用表统计信息直接回答 min/max 查询
count(*) #620 使用表统计信息直接回答 count(*)
窗口函数 #569、#571、#595 repartition 并行化、排序折叠进逻辑层、ranges 查找优化
to_timestamp #387 规划阶段常量折叠
批次构建 #339 create_batch_from_map 提速
数学表达式 #309 改用 unary kernel
哈希 #556、#807、#827 浮点哈希修正、hash 抽离、RawTable API

其中「用统计信息回答 min/max/count(*)」的思路在后续版本演化为更通用的「统计信息下推」框架,直到今天仍然能在 operator_statistics 与 statistics.rs 等模块中找到对应实现。

五、升级迁移速查

如果你正从 4.x 迁移到 5.0.0,请优先处理以下破坏性变更:

  1. 改名:MergeExec → CoalescePartitionsExec;RowGroupPredicateBuilder → PruningPredicateBuilder。
  2. Join:改写多表 join 的连接条件书写顺序;非等值条件允许出现在 join 条件中(会转为 filter)。
  3. 函数语义:log(x) 现在是 log10(x),自然对数请显式使用 ln(x)。
  4. API:扩展物理规划的入口需要接收逻辑计划;SQLMetric 构造与命名方式变更;PredicateBuilder 返回 Vec<bool>。
  5. 依赖:arrow 系列依赖统一来自 crates.io 发布版(4.x → 5.0),不再依赖 GitHub SHA。
  6. 剪枝:如需禁用 parquet 剪枝,使用 ExecutionConfig 中新增的开关(当前版本对应配置项为 datafusion.execution.parquet.enable_pruning,默认开启)。

六、延伸阅读

说明:文中「当前仓库」的代码为演进后的最新实现,与 5.0.0 时期的历史代码在文件组织上已有差异,但核心设计(如 CoalescePartitionsExec 的合并语义、PruningPredicate 的容器剪枝模型、now() 的 Stable 语义)保持一致,可作为理解 5.0.0 设计意图的参考依据。

登录后查看全文
datafusion