☰
Apache DataFusion 50.0.0 版本解析:QUALIFY 子句、Parquet 元数据缓存与查询引擎关键改进
2026/9/25 3:06:53 网站建设 项目流程
  • 大数据
  • 数据分析
  • 后端

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

项目地址:https://gitcode.com/gh_mirrors/datafu/datafusion
点击查看免费下载

Apache DataFusion 50.0.0 是围绕 SQL 查询引擎核心能力的一次大版本更新,累计合并 315 个 PR(来自 79 位贡献者),涵盖 SQL 语法、物理执行、存储 I/O、UDF 框架与性能优化等多个层面。本文以仓库内 dev/changelog/50.0.0.md 为主线,逐类拆解 Breaking Changes、性能相关改动、新特性与缺陷修复,并深入对应源码与测试,帮助你评估升级影响、理解底层实现,并为 50.x 迁移做准备。

版本概况与升级提示

50.0.0 的核心改动可以归纳为四个方向:

  1. SQL 能力增强:新增QUALIFY子句、窗口函数DISTINCT/FILTER支持、Spark 兼容函数大量补齐;
  2. 存储与缓存:内置 Parquet 读取器元数据缓存、限制文件元数据缓存内存占用、动态 Parquet 加密解密属性;
  3. 执行框架重构:Nested Loop Join 执行器重写(约 5 倍提速、内存占用降至约 1%)、多级归并排序、SortMergeJoin 模块拆分;
  4. UDF/表达式一致性:ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl全面基于Eq/Hash派生相等性与哈希,废弃了一批旧 API。

此外本版本将 MSRV(最低支持 Rust 版本)提升至1.86.0(PR #17230),Arrow/Parquet 依赖升级到56.0.0(PR #16690),sqlparser 升级到0.58(PR #16456)。升级前请先阅读官方升级指南并确认工具链版本满足要求。

破坏性变更(Breaking Changes):升级时必须注意

多 ORDER BY 的 array_agg 聚合

PR #16625 支持多个有序的array_agg聚合,即同一个查询中可以在不同array_agg上各自指定ORDER BY。相应地,ArrayAgg函数从参数化函数改为基于Eq/Hash派生相等性,datafusion/sqllogictest/test_files/string_agg.slt等测试也补充了array_agg带排序的回归用例(见 PR #17033)。

异步 UDF 与窗口 UDF 接口一致性

  • PR #16902 将AsyncScalarUDFImpl::invoke_async_with_args调整为与ScalarUDFImpl::invoke_with_args一致,异步 UDF 编写方式随之统一;
  • PR #17081 让WindowUDFImpl的相等性与哈希直接由Eq、Hash派生,自定义窗口 UDF 若字段不实现这两个 trait 将无法编译;
  • 后续 PR #17130、#17164 分别将相同方案推广到AggregateUDFImpl与ScalarUDFImpl,这是本版本对扩展点最重要的行为收敛。

文件与表达式 API 调整

  • PR #17398 将ProjectionExpr由元组改为结构体,依赖字段解构的外部代码需要适配;
  • PR #17397 中FileOpenFuture的错误类型从ArrowError改为DataFusionError;
  • PR #17407 聚合函数经 FFI 调用时改为使用return_field而非return_type;
  • PR #17199 从扩展的check_invariants中移除了冗余的plan参数。

parquet_encryption 改为非默认特性

PR #17137 将parquet_encryption改为非默认 feature,默认构建不再包含加密能力。需要使用 Parquet 加解密的用户必须在Cargo.toml中显式启用该 feature。

性能相关改进:查询引擎提速清单

本版本性能改动集中在表达式求值、聚合去重与字符串处理上:

改动PR要点
LiteralGuarantee 增强#16762对(a=1 AND b=1) OR (a=2 AND b=3)这类「同一列多值」复合谓词给出更强的常量保证,利于谓词推导
initcap优化#16878通过避免内存分配加速首字母大写函数
date_trunc提速#16859部分场景下约 7 倍提速
Hash Expr 性能#16977改进表达式哈希计算路径
NULL 比较/二元运算简化#17088简化涉及 NULL 的比较与二元运算
消除冗余聚合#17139在逻辑优化阶段剔除可证明冗余的聚合
StringView 无缓冲快速路径#17008移植 arrow-rs 的get_buffer_memory_size优化,为 GC string view 增加无缓冲快速路径

新特性详解(一):QUALIFY 子句落地

PR #16933 实现了标准 SQL 的QUALIFY子句:在窗口函数计算完成之后、ORDER BY/LIMIT之前进行过滤,解决「对窗口函数结果直接过滤」的经典痛点——此前只能写子查询包一层,现在可以在主查询直接过滤。

语法与语义

SELECT id, name, ROW_NUMBER() OVER (PARTITION BY dept ORDER BY salary DESC) AS rn FROM users QUALIFY rn = 1 ORDER BY dept, id;

QUALIFY的求值顺序位于HAVING(聚合后)之后、窗口计算之后,可以引用 SELECT 列表中出现的窗口函数别名(如上面的rn),也可以直接引用窗口函数表达式。

源码中的实现路径

QUALIFY 在 SQL 层由 datafusion/sql/src/select.rs 解析:qualify_expr先通过sql_expr_to_logical_expr转为逻辑表达式,再用resolve_aliases_to_exprs解析 SELECT 别名、normalize_col规范化列引用,随后与HAVING、聚合表达式一起通过find_aggregate_exprs收集聚合函数(select.rs)。之后逻辑计划分两阶段重写:

  1. 聚合阶段:qualify_expr经rebase_expr重写为引用聚合投影中的列(select.rs);
  2. 窗口阶段:通过find_window_exprs收集窗口函数,若 QUALIFY 中不存在窗口表达式则直接按普通过滤处理,否则将表达式重写到窗口计划之上,最终在窗口计划之上挂上filter(qualify_expr_post_window)(select.rs)。

对应地,如果窗口函数出现在WHERE或HAVING中,解析器会给出指向性错误提示:Move the condition that uses this window function to a QUALIFY clause, which is evaluated after window functions are computed(见 datafusion/sql/src/expr/mod.rs)。

实测验证:sqllogictest 覆盖

datafusion/sqllogictest/test_files/qualify.slt 用一张 8 行的users表系统验证了 QUALIFY 的各种形态:

  • 基础ROW_NUMBER+QUALIFY rn = 1:取每个部门薪资最高的人;
  • RANK/DENSE_RANK/NTILE/PERCENT_RANK/CUME_DIST组合;
  • 复合条件:QUALIFY rn <= 2 AND age_rank <= 5;
  • 引用LAG/LEAD结果:QUALIFY prev_salary IS NOT NULL AND salary > prev_salary;
  • 同时引用多个窗口函数别名(rn、age_rank、dept_age_rank三者组合过滤)。

新特性详解(二):窗口函数 DISTINCT 与 FILTER

  • PR #16925 支持窗口聚合的DISTINCT(首批支持sum(distinct ...)见 PR #16943,后续还补充了avg(distinct)对float64的支持,PR #17255);
  • PR #17378 支持窗口聚合函数上的FILTER子句,即COUNT(x) FILTER (WHERE ...) OVER (...)形态;
  • 配套的 protobuf 序列化在 PR #17235 中保证窗口表达式distinct与ignore_nulls标记在datafusion-proto往返过程中不丢失。

新特性详解(三):Parquet 元数据缓存与加密

内置 Parquet 读取器缓存元数据

PR #16971 让内置 Parquet 读取器缓存文件元数据,配合 PR #17031 对缓存内存上限的约束,避免海量小文件场景下重复读取 footer。实现上,ParquetFormat从RuntimeEnv的CacheManager获取文件元数据缓存并注入到读取配置中(见 datafusion/core/src/datasource/file_format/parquet.rs)。

缓存容量通过运行时配置datafusion.runtime.metadata_cache_limit控制,默认值为 50MB(DEFAULT_METADATA_CACHE_LIMIT: usize = 50 * 1024 * 1024,见 datafusion/execution/src/cache/cache_manager.rs),在 datafusion/core/src/execution/context/mod.rs 中解析进CacheManagerConfig::with_metadata_cache_limit。

相关配套改动还包括:使用缓存元数据计算ListingTable统计信息(PR #17022)、提供查看元数据缓存内容的能力(PR #17126),以及将 Parquet 元数据统一收敛为DFParquetMetadata结构(PR #17127)。旧配置项datafusion.execution.parquet.cache_metadata已被移除(PR #17062),统一走运行时缓存路径。

动态 Parquet 加密/解密

PR #16779 引入动态 Parquet 加密与解密属性,PR #17342 进一步将ParquetEncryptionFactory改为异步接口,配合动态属性可在运行时按需获取密钥与加密配置。PR #17426 重新启用了加密 Parquet 的 page index 支持。注意使用该能力需显式开启parquet_encryptionfeature(见上文 Breaking Changes)。

新特性详解(四):执行器与查询框架演进

  • Nested Loop Join 重写(PR #16996):执行器重写后官方标注约 5 倍提速、内存占用约降至 1%,升级指南(PR #17202)对 NLJ 行为变化作了说明,涉及swap_inputs()语义的注意事项见 PR #17373;
  • 多级归并排序(PR #15700):新增「始终适配内存」的多级 merge sort 路径,配合SpillManager与DiskManager使用;DiskManager默认临时目录上限为 100GB(DEFAULT_MAX_TEMP_DIRECTORY_SIZE,见 datafusion/execution/src/disk_manager.rs);
  • SortMergeJoin 增强:支持二进制类型的on子句(PR #17431)、protobuf 序列化(PR #17296)、模块拆分(PR #17304)、基于BufferedBatchState枚举的 spilling 重构(PR #17429);
  • 动态过滤器(bounds)下推(PR #16445):HashJoinExec 支持动态范围过滤下推,配套修复了分区查询下侧向信息传递(PR #17197)、TopK 动态过滤器仅在新过滤器更选择性时更新(PR #16433)以及 bounds accumulator 重置问题(PR #17371);
  • spawn/spawn_blocking 支持(PR #17239):RecordBatchReceiverStreamBuilder允许在外部提供的 runtime 上执行任务。

新特性详解(五):Spark 兼容函数大规模补齐

本版本在datafusion/spark(即 datafusion-spark 兼容层)中实现了一大批 Spark 内置函数,并在 datafusion/spark/src/function 下可找到对应实现:

  • 字符串类:luhn_check(#16848)、like/ilike(#16962)、parse_url(#16937)、hex(支持 utf8view,#16885);
  • 日期时间类:last_day(#16828)、next_day(#16780)、date_add/date_sub(#17024);
  • 数学类:rint(#16924)、mod/pmod(#16829)、bit_get/bit_count(#16942)、width_bucket(#17331);
  • 哈希类:crc32/sha1(#17032);
  • 位图与条件类:bitmap_count(#17179)、if(#16946)、array(#16936)。

通用 SQL 层同样有增强:date_part新增isodow(ISO 星期几,Monday=0)支持(PR #17112,实现在 datafusion/functions/src/datetime/date_part.rs);approx_percentile_cont(_with_weight)恢复旧语法并新增centroids配置(PR #16999、#17003)。

缺陷修复要点(Fixed Bugs)

本版本修复了多个与执行正确性相关的缺陷,升级后值得回归验证:

  • 排序输出一致性:sort 保证每批输出固定batch_size行(#17244);Partial Sort 批次间切片位置修复(#16881);
  • 统计信息:PlaceholderRowExec::partition_statistics(#16851)、合并统计时distinct_count置为 Absent(#17385)、EquivalenceProperties::constants返回全部常量(#17404);
  • 序列化/反序列化:FilterExec含 INLIST 谓词的反序列化(#17224)、物理计划序列化的三个缺陷(#16858)、ComposedPhysicalExtensionCodec编解码不一致(#16986)、AnalyzeExecprotobuf 往返(#17234)、内存数据源 protobuf 支持(#17290);
  • 类型与空值对齐:cast decimal→timestamp 标量与数组不一致(#16539)、array_has标量空值缓冲对齐(#17272)、map_keys可空性标志对齐(#17454)、case else 惰性求值(#17311);
  • 内存与泄漏:FFIRecordBatchStream内存泄漏(#17190)、slicedStringViewArray内存计量错误(#17315)、array_agg 内存过量计量(#16816);
  • 其他:NULL IN ()优化错误(#17092)、结构体 unnest 谓词跳过(#16790)、Windows 路径崩溃(#17231)、row group 元数据 inexact 标志(#16412)。

新运行时配置与文档更新

  • datafusion.runtime.temp_directory与datafusion.runtime.max_temp_directory_size:控制溢写临时目录与磁盘上限(PR #16934,解析逻辑见 datafusion/core/src/execution/context/mod.rs);配套测试位于 datafusion/core/tests/sql/runtime_config.rs,例如SET datafusion.runtime.max_temp_directory_size = '0K'时溢写会被拒绝;
  • datafusion.sql_parser.default_null_ordering:自定义默认 NULL 排序(PR #16963,配置定义见 datafusion/common/src/config.rs,默认"nulls_max");
  • 文档方面新增两份调优指南:小数据/短查询调优(#17040)与大于内存查询调优(#17069),并在配置选项页面补充了示例(#17039);
  • 其他基础设施:ExecutionPlan::reset_state(#17028)、物理优化器日志降为 debug 级别(#17383)、增加is_volatile检查禁止易变函数下推(#16861、#17351)。

升级建议小结

  • 自定义 UDF/UDWF/UDAAF 作者:优先检查实现了Eq/Hash的派生条件,并适配invoke_async_with_args与return_field语义变化;
  • 依赖 Parquet 加密的用户:显式开启parquet_encryptionfeature,并将加密工厂实现迁移到异步接口;
  • 使用投影表达式、文件打开回调等扩展点的用户:适配ProjectionExpr结构体化与FileOpenFuture错误类型变更;
  • 关注性能的查询:QUALIFY可简化窗口过滤查询,Parquet 元数据缓存与metadata_cache_limit适合小文件密集场景,max_temp_directory_size为溢写磁盘用量提供硬约束。

50.0.0 的完整改动清单、PR 明细与贡献者名单可继续查阅 dev/changelog/50.0.0.md,后续 50.x 补丁(50.1.0 至 50.3.0)的变更记录也保存在 dev/changelog 目录中。

  • 大数据
  • 数据分析
  • 后端

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

项目地址:https://gitcode.com/gh_mirrors/datafu/datafusion
点击查看免费下载
上一篇:如何扩展LIRE:自定义图像特征提取器的开发指南 🚀
下一篇:CANN/ge TensorDesc张量描述API

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询