- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
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 的核心改动可以归纳为四个方向:
- SQL 能力增强:新增
QUALIFY子句、窗口函数DISTINCT/FILTER支持、Spark 兼容函数大量补齐; - 存储与缓存:内置 Parquet 读取器元数据缓存、限制文件元数据缓存内存占用、动态 Parquet 加密解密属性;
- 执行框架重构:Nested Loop Join 执行器重写(约 5 倍提速、内存占用降至约 1%)、多级归并排序、SortMergeJoin 模块拆分;
- 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)。之后逻辑计划分两阶段重写:
- 聚合阶段:
qualify_expr经rebase_expr重写为引用聚合投影中的列(select.rs); - 窗口阶段:通过
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); - 内存与泄漏:FFI
RecordBatchStream内存泄漏(#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
相关推荐
Tinypool性能优化指南:如何压榨Node.js多线程的最大潜力
Tinypool性能优化指南:如何压榨Node.js多线程的最大潜力 Tinypool是一个轻量级的Node.js Worker Thread Pool实现,仅
如何在UE5中实现3D高斯泼溅实时渲染:XScene-UEPlugin完整指南
如何在UE5中实现3D高斯泼溅实时渲染:XScene UEPlugin完整指南 XScene UEPlugin是专为Unreal Engine 5设计的高斯泼溅
图形学3D渲染计算机视觉深度学习Apache DataFusion查询缓存键生成:包含参数策略
Apache DataFusion查询缓存键生成:包含参数策略 在数据处理场景中,重复执行相同或相似查询会导致资源浪费和性能下降。Apache DataFusi
大数据数据分析后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考