Apache DataFusion 5.0.0 版本深度解析:Join 语义、窗口函数、Parquet 剪枝与执行器重构
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 最重量级的特性,经历了一个清晰的渐进过程:
- 框架搭建(#334):先建立 window expression 的逻辑/物理规划与 EXPLAIN 结构,仅支持空
OVER (); order by与partition by子句(#463、#501、#520、#558):逐步加入排序分区语义;- 窗口帧(window frame)支持(#506、#530):加入
ROWS/RANGE帧构造、边界检查与错误处理; - 内置窗口函数全集(#403、#429、#631、#687):
first_value、last_value、nth_value、lead、lag(支持 offset 与 default 参数)、row_number、rank、dense_rank; - 性能与正确性(#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):剪枝逻辑整体改用
PruningStatisticstrait + 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 处理;使用
RawTableAPI 提升 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,请优先处理以下破坏性变更:
- 改名:
MergeExec→CoalescePartitionsExec;RowGroupPredicateBuilder→PruningPredicateBuilder。 - Join:改写多表 join 的连接条件书写顺序;非等值条件允许出现在 join 条件中(会转为 filter)。
- 函数语义:
log(x)现在是log10(x),自然对数请显式使用ln(x)。 - API:扩展物理规划的入口需要接收逻辑计划;
SQLMetric构造与命名方式变更;PredicateBuilder返回Vec<bool>。 - 依赖:arrow 系列依赖统一来自 crates.io 发布版(4.x → 5.0),不再依赖 GitHub SHA。
- 剪枝:如需禁用 parquet 剪枝,使用
ExecutionConfig中新增的开关(当前版本对应配置项为datafusion.execution.parquet.enable_pruning,默认开启)。
六、延伸阅读
- 完整条目列表:dev/changelog/5.0.0.md
- 当前仓库的物理执行器:coalesce_partitions.rs、sorts/sort_preserving_merge.rs、joins/hash_join/exec.rs
- 剪枝体系:datafusion/pruning/src/pruning_predicate.rs
- 窗口函数:datafusion/functions-window/src
- 时间函数:datetime/now.rs、datetime/date_trunc.rs
- 测试证据:dataframe 集成测试、merge fuzz 测试
说明:文中「当前仓库」的代码为演进后的最新实现,与 5.0.0 时期的历史代码在文件组织上已有差异,但核心设计(如
CoalescePartitionsExec的合并语义、PruningPredicate的容器剪枝模型、now()的 Stable 语义)保持一致,可作为理解 5.0.0 设计意图的参考依据。