- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
本指南基于 Apache DataFusion 官方库使用指南(
docs/source/library-user-guide/upgrading/50.0.0.md)整理,全面覆盖 50.0.0 版本中影响库使用者的全部破坏性变更(Breaking Changes):包含ListingTable的 Hive 分区自动推断新行为与恢复旧行为的配置方法、MSRV 提升至 1.86.0、三个 UDF 特质对PartialEq/Eq/Hash的强制要求、AsyncScalarUDFImpl签名变更、ProjectionExpr从类型别名到结构体的重构,以及ExecutionPlan::reset_state、FileOpenFuture错误类型、FFI UDAF 等底层 API 调整。读完本文,你将能够逐项对照完成自有代码库向 DataFusion 50.0.0 的平滑迁移。
升级前须知:本文涉及的核心变更总览
DataFusion 50.0.0 是一次面向库使用者的 API 调整版本,变更横跨表工厂行为、Rust 工具链、用户自定义函数(UDF/UDAF/UDWF)、物理执行计划、数据源流式读取与 FFI 接口等多个层面。下表按影响面归纳本指南将逐一展开的变更项:
| 变更类别 | 核心内容 | 影响人群 |
|---|---|---|
| ListingTable 行为 | CREATE EXTERNAL TABLE自动推断 Hive 分区列 | 使用 Hive 风格目录组织数据的用户 |
| Rust 工具链 | MSRV 提升至 1.86.0 | 所有以 DataFusion 为依赖的项目 |
| UDF 特质 | ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl要求实现PartialEq、Eq、Hash | 所有自定义函数实现者 |
| 异步 UDF | AsyncScalarUDFImpl::invoke_async_with_args返回ColumnarValue并移除ConfigOptions参数 | 注册远程/异步函数(如调用 LLM)的用户 |
| 投影表达式 | ProjectionExpr由元组别名改为具名字段结构体 | 直接使用ProjectionExec的物理计划开发者 |
| 会话配置 | options()等 API 返回&Arc<ConfigOptions> | 访问会话配置的库使用者 |
| 物理计划 | 新增ExecutionPlan::reset_state | 自定义ExecutionPlan实现者 |
| 表达式特质 | 新增PhysicalExpr::is_volatile_node | 自定义PhysicalExpr实现者 |
| 数据源 | FileOpenFuture错误类型改为DataFusionError | FileOpener自定义实现者 |
| schema 重写 | schema_rewriter模块迁入新 crate | 依赖物理表达式适配器的用户 |
| FFI | UDAF 从return_type迁移到return_field | FFI 库开发者与测试作者 |
ListingTable自动检测 Hive 分区表
新行为:分区列自动进入表 Schema
在 50.0.0 之前,通过ListingTableFactory或CREATE EXTERNAL TABLE创建ListingTable时,采用 Hive 分区布局(例如/table_root/column1=value1/column2=value2/data.parquet)的数据集,其分区列(column1、column2)不会反映在表 schema 中,也不会出现在查询结果里。从 50.0.0 起,ListingTableFactory会自动推断 Hive 分区,将column1、column2作为分区列纳入表的 schema 与数据。
仓库源码证实了这一行为:在 listing_table_factory.rs 中,当CREATE EXTERNAL TABLE未显式提供列定义(cmd.schema.fields().is_empty())且未显式指定PARTITIONED BY列时,工厂会读取会话配置并调用options.infer_partitions(session_state, first_path)从目录结构推断分区列,推断得到的列以Dictionary(UInt16, Utf8)类型写入表 schema:
let infer_parts = session_state .config_options() .execution .listing_table_factory_infer_partitions; let part_cols = if cmd.table_partition_cols.is_empty() && infer_parts { options .infer_partitions(session_state, first_path) .await? .into_iter() } else { cmd.table_partition_cols.clone().into_iter() };恢复旧行为的配置开关
如果你依赖旧行为(分区列不出现在 schema 中),可以将配置项datafusion.execution.listing_table_factory_infer_partitions设置为false来恢复。该配置项定义在 datafusion/common/src/config.rs,属于datafusion.execution命名空间:
/// Should a `ListingTable` created through the `ListingTableFactory` infer table /// partitions from Hive compliant directories. Defaults to true (partition columns are /// inferred and will be represented in the table schema). pub listing_table_factory_infer_partitions: bool, default = true配置方式与 DataFusion 其他执行期配置一致:
- SQL 会话级:
SET datafusion.execution.listing_table_factory_infer_partitions = false; - SessionConfig 编程方式:构造
SessionConfig时通过config.set(...)或直接修改config_options().execution.listing_table_factory_infer_partitions字段。
仓库的 sqllogictest 测试对两种行为都做了回归覆盖:listing_table_partitions.slt 分别以false和true设置该开关并验证分区列的推断结果;information_schema.slt 也记录了该配置的默认值与说明文本。详细讨论见上游 issue #17049。
MSRV 更新至 1.86.0
DataFusion 50.0.0 将最低支持的 Rust 版本(MSRV)提升到1.86.0。这意味着所有以 DataFusion 为依赖的项目,其工具链版本不得低于 1.86.0,否则无法编译。
需要说明的是:当前仓库主分支的rust-toolchain.toml已指定channel = "1.98.1",这是后续版本持续演进的结果;对于锁定在 50.0.0 的用户,请以该版本发布时的 MSRV(1.86.0)为准,并在 CI 中相应调整rust-version或工具链配置。升级理由与讨论详见上游 PR #17230。
ScalarUDFImpl、AggregateUDFImpl、WindowUDFImpl特质要求PartialEq、Eq与Hash
变更动机:消灭手写equals/hash_value的错误隐患
此前ScalarUDFImpl::equals、AggregateUDFImpl::equals、WindowUDFImpl::equals以及配套的hash_value方法需要实现者手工编写相等性与哈希逻辑,极易出错(例如遗漏新加入的字段)。50.0.0 移除了这三个特质的equals与hash_value方法,改为强制要求实现 Rust 标准的PartialEq、Eq和Hash特质,从而让函数相等性比较走上编译器保障的轨道。
仓库当前代码印证了这一演进方向:在 datafusion/expr/src/udf.rs、datafusion/expr/src/udaf.rs 与 datafusion/expr/src/udwf.rs 中,三个特质均约束为Debug + DynEq + DynHash + Send + Sync + Any;datafusion/expr/src/udf_eq.rs 中的DynEq/DynHash桥接特质通过dyn_eq与dyn_hash将标准特质转发到具体实现,使ScalarUDF等封装类型得以在 trait object 层面执行相等性比较与哈希。
迁移方法:正则替换 vs 手工实现
绝大多数标量函数是无状态的,且结构体中只有signature字段。这类函数可以借助正则表达式批量迁移:
- 搜索模式:
\#\derive\(Debug\)\?struct \w+ \{\n *signature\: Signature\,\n *\}) - 替换为:
#[derive(Debug, PartialEq, Eq, Hash)]$1
替换完成后务必人工审查所有改动,确认只影响了函数结构体,没有误伤其他类型。对于含额外字段或需要自定义相等语义(例如忽略运行时统计字段)的函数,则需手工实现PartialEq/Eq/Hash,保证比较逻辑与字段变更保持同步。
AsyncScalarUDFImpl::invoke_async_with_args的两次签名调整
异步标量函数特质(AsyncScalarUDFImpl,用于注册远程函数,如调用 LLM 的AskLLM示例)在 50.0.0 中经历了两次相关调整。先看当前仓库中该特质的最终形态(datafusion/expr/src/async_udf.rs):
#[async_trait] pub trait AsyncScalarUDFImpl: ScalarUDFImpl { ... async fn invoke_async_with_args( &self, args: ScalarFunctionArgs, ) -> Result<ColumnarValue>; }调整一:返回类型从ArrayRef变为ColumnarValue
为使返回值能够走单值(scalar value)优化并与其他 UDF API 保持一致,invoke_async_with_args的返回类型由ArrayRef改为ColumnarValue。迁移只需在旧返回值外层包一层转换:
# /* comment to avoid running impl AsyncScalarUDFImpl for AskLLM { async fn invoke_async_with_args( &self, args: ScalarFunctionArgs, _option: &ConfigOptions, ) -> Result<ColumnarValue> { .. return ColumnarValue::from(array_ref); // new code:ColumnarValue::from 完成 ArrayRef -> ColumnarValue 转换 } } # */调整二:移除_option: &ConfigOptions参数
ConfigOptions已可通过ScalarFunctionArgs参数直接获取,因此invoke_async_with_args中原有的_option: &ConfigOptions参数被移除,接口得到简化。迁移示例:
# /* comment to avoid running impl AsyncScalarUDFImpl for AskLLM { async fn invoke_async_with_args( &self, args: ScalarFunctionArgs, ) -> Result<ColumnarValue> { let options = &args.config_options; // 从参数中读取配置,替代原 _option 参数 .. } ... } # */ScalarFunctionArgs结构体(定义于 datafusion/expr/src/udf.rs)现在携带args、arg_fields、number_rows、return_field以及config_options: Arc<ConfigOptions>五个字段,执行期配置由此通过参数透传,无需再从外部单独注入。相关实现细节见上游 issue #16896。
ProjectionExpr从类型别名重构为结构体
变更前后对比
ProjectionExpr(投影表达式 + 别名)从元组类型别名改为带具名字段的结构体,以提升代码可读性与可维护性:
变更前(类型别名):
pub type ProjectionExpr = (Arc<dyn PhysicalExpr>, String);变更后(结构体):
#[derive(Debug, Clone)] pub struct ProjectionExpr { pub expr: Arc<dyn PhysicalExpr>, pub alias: String, }迁移要点
- 构造:将元组构造
(expr, alias)替换为ProjectionExpr::new(expr, alias)或结构体字面量ProjectionExpr { expr, alias }; - 字段访问:将
.0、.1替换为.expr、.alias; - 模式匹配:将
(expr, alias)模式更新为ProjectionExpr { expr, alias }。
仓库中大量代码已按新形态使用,例如 datafusion/core/src/physical_planner.rs 中的.map(|(expr, alias)| ProjectionExpr { expr, alias }),以及 datafusion/core/tests/physical_optimizer/projection_pushdown.rs 中的ProjectionExpr::new(Arc::new(Column::new("b", 2)), "b")等测试用法,可作迁移参考。此变更主要影响ProjectionExec的使用者,实现见上游 PR #17398。
SessionState/SessionConfig/OptimizerConfig返回&Arc<ConfigOptions>
为让ConfigOptions获得更广泛的访问路径并减少不必要的克隆,options()等访问器由返回&ConfigOptions改为返回&Arc<ConfigOptions>。借助Arc,同一份ConfigOptions可在线程间共享,只有真正需要修改时才克隆整个结构。
多数情况下 Rust 编译器会自动解引用Arc,因此大部分用户无感;但在显式标注类型的场景需要补上.as_ref():
# /* comment to avoid running let optimizer_config: &ConfigOptions = state.options(); // 旧写法 let optimizer_config: &ConfigOptions = state.options().as_ref(); // 新写法 # */ScalarFunctionArgs::config_options(类型为Arc<ConfigOptions>,见 datafusion/expr/src/udf.rs)即此模式的应用之一。详见上游 PR #16970。
Schema Rewriter 模块迁移至新 crate
schema_rewriter模块及其符号从datafusion_physical_expr迁出,进入新 cratedatafusion_physical_expr_adapter。涉及以下符号:
DefaultPhysicalExprAdapterDefaultPhysicalExprAdapterFactoryPhysicalExprAdapterPhysicalExprAdapterFactory
仓库当前结构印证了迁移结果:新 crate 的 datafusion/physical-expr-adapter/src/lib.rs 直接导出这些符号,而 schema_rewriter.rs 则承载了 schema 重写的工作职责(包括为缺失列填充默认值、按物理 schema 投影重写表达式等)。
迁移时只需更新 import 路径:
use datafusion_physical_expr_adapter::{ DefaultPhysicalExprAdapter, DefaultPhysicalExprAdapterFactory, PhysicalExprAdapter, PhysicalExprAdapterFactory };Arrow 与 Parquet 升级到 56.0.0
DataFusion 50.0.0 将底层 Apache Arrow 实现升级到56.0.0(Parquet 同样为 56.0.0)。如果你在代码中直接依赖arrow或parquetcrate 并引用了相关 API,请同步升级依赖版本,并以 Arrow 56.0.0 的发布说明为准核对变更。
新增ExecutionPlan::reset_state方法
变更背景
DataFusion 49.0.0 曾存在一个 bug:动态过滤器(dynamic filter,目前仅在出现类似ORDER BY ... LIMIT ...的查询时生成)在递归查询(recursive query,如递归 CTE)中产生错误结果。50.0.0 通过为ExecutionPlan特质新增reset_state方法修复此问题——凡是需要在执行计划树中维护内部状态、或持有对其他节点引用的ExecutionPlan,都应实现该方法以在重新执行时重置状态。
默认实现与实现要求
reset_state的默认实现(datafusion/physical-plan/src/execution_plan.rs)只是用现有 children 调用replace_children重建一个同构实例,并不重置内部状态;因此任何有状态组件(例如DynamicFilterPhysicalExpr)的执行计划都必须覆写它。需要注意:该方法不应递归重置 children,因为框架期望在遍历执行计划树时对每个子节点逐个调用。SortExec的示例实现可参考上游 PR #17028,本仓库的 sort.rs 中亦有对应实现。
嵌套循环连接(NLJ)重写:排序保持的取舍
NestedLoopJoin算子被从零重写以提升性能与内存效率。官方文档给出的微基准结论是:相比旧实现,新实现在极端情况下可获得最高5 倍加速,内存占用仅为原来的1%。
但需要明确这一变更带来的行为差异:新实现无法像旧版本那样保持输入的有序性(input sort order)。这是性能与内存效率优先于排序保持的根本性设计取舍,而非 bug。如果你的查询依赖 NLJ 输出保持输入顺序(例如在下游继续依赖该顺序的算子链),请评估重写后的计划形态或显式引入排序算子。详见上游 PR #16996。
LazyBatchGenerator新增as_any()方法
为支持 protobuf 序列化,LazyBatchGenerator特质新增了as_any()方法,自定义实现需要补充:
# /* comment to avoid running impl LazyBatchGenerator for MyBatchGenerator { fn as_any(&self) -> &dyn Any { self } ... } # */仓库中GenerateSeries等表函数实现(datafusion/functions-table/src/generate_series.rs)已按此模式提供as_any返回self,可作为参考模板。详见上游 PR #17200。
DataSource::try_swapping_with_projection重构
DataSource::try_swapping_with_projection被重构以简化方法签名,并尽量减少ExecutionPlan与DataSource抽象层之间的耦合泄漏。对于任何自定义DataSource,重新实现该方法相对直接:保持投影交换的目标语义,同时遵循精简后的签名。详细说明见上游 PR #17395。
FileOpenFuture错误类型从ArrowError改为DataFusionError
变更前后对比
FileOpenFuture类型别名的错误类型由ArrowError统一为DataFusionError。这影响FileOpener特质及所有涉及文件流式读取的实现。
变更前:
pub type FileOpenFuture = BoxFuture<'static, Result<BoxStream<'static, Result<RecordBatch, ArrowError>>>>;变更后:
pub type FileOpenFuture = BoxFuture<'static, Result<BoxStream<'static, Result<RecordBatch>>>>;当前仓库中的定义(datafusion/datasource/src/file_stream/mod.rs)即为去掉内层错误类型参数的简化形态——内层Result与FileStream其余错误路径统一走DataFusionError。同时,FileStreamState枚举的Open变体也做了相应调整。
迁移动作
如果你的代码有自定义FileOpener实现或直接处理FileOpenFuture,需要将错误处理从ArrowError迁移到DataFusionError(例如通过arrow_datafusion_err!宏或DataFusionError::from转换)。详见上游 PR #17397。
FFI 用户自定义聚合函数签名变更
从return_type到return_field
FFI(Foreign Function Interface,datafusion/ffi/src/udaf/mod.rs)中 UDAF 的 C ABI 结构体已改为调用底层聚合函数的return_field方法(而非return_type),以支持聚合函数返回字段的元数据处理。仓库中 datafusion/ffi/src/udaf/accumulator_args.rs 的AccumulatorArgs也携带了return_field: FieldRef,与 FFI 侧新的return_fieldC 函数指针(physical_expr/mod.rs)配套。
对大多数用户而言该变更透明;但如果你编写过直接调用return_type的单元测试,需要改为调用return_field。
FFI 跨版本使用的注意事项
这是 FFI API 的一次破坏性变更。当前最佳实践是:确保所有交互的库使用相同的底层 Rust 版本,以规避 ABI 不一致。上游 issue #17374 正在讨论稳定该接口,使这些库未来可以跨不同 DataFusion 版本互操作。具体实现见上游 PR #17407。
新增PhysicalExpr::is_volatile_node
为正确标记易变(volatile)表达式——即每次求值可能返回不同结果的表达式(典型如随机数函数)——PhysicalExpr特质新增了is_volatile_node方法:
impl PhysicalExpr for MyRandomExpr { fn is_volatile_node(&self) -> bool { true } }该方法提供了默认值false以最小化破坏面,但官方强烈建议所有PhysicalExpr实现者显式选择行为(即使只是返回false)。
is_volatile_node在优化器中已有实际用途:在公共子表达式消除(CSE)规则中(datafusion/optimizer/src/common_subexpr_eliminate.rs),is_valid检查通过!node.is_volatile_node()将易变表达式排除在可公共子表达式提取的范围之外,避免对random()这类函数错误地做重复计算消除;datafusion/optimizer/src/utils.rs 也将其用于表达式树遍历中的易变性判断。此外 datafusion/ffi/src/physical_expr/mod.rs 的 FFI 包装层同样暴露了该函数指针。测试覆盖可参见 datafusion/physical-expr/src/physical_expr.rs(默认返回false)与 higher_order_function.rs。更多讨论与示例实现见上游 PR #17351。
迁移清单:50.0.0 升级自查表
为便于对照执行,汇总一份升级自查清单:
- 工具链:确认 Rust ≥ 1.86.0(
rust-toolchain.toml中当前主分支为 1.98.1,锁定 50.0.0 时以 1.86.0 为准)。 - 分区表行为:如不想要自动推断 Hive 分区,设置
datafusion.execution.listing_table_factory_infer_partitions = false(默认true,config.rs)。 - UDF 特质:为
ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl实现PartialEq、Eq、Hash,删除手写equals/hash_value;无状态函数可用正则批量加 derive。 - 异步 UDF:
invoke_async_with_args返回ColumnarValue(用ColumnarValue::from包裹ArrayRef),并移除_option参数,改从args.config_options读取。 - 投影:
ProjectionExpr改为结构体,用ProjectionExpr::new(expr, alias)构造、.expr/.alias访问、ProjectionExpr { expr, alias }匹配。 - 配置访问:标注类型的
options()调用补.as_ref()(返回&Arc<ConfigOptions>)。 - schema 重写:import 改为
datafusion_physical_expr_adapter::{...}。 - Arrow/Parquet:升级依赖到 56.0.0。
- 自定义 ExecutionPlan:按需实现
reset_state(有状态或持引用时必须)。 - NLJ:确认不依赖旧实现的输入排序保持;必要时显式排序。
- LazyBatchGenerator:补充
as_any()返回self。 - DataSource:按新签名重新实现
try_swapping_with_projection。 - FileOpenFuture:错误处理改用
DataFusionError。 - FFI UDAF:测试从
return_type改为return_field;保证交互库 Rust 版本一致。 - PhysicalExpr:显式实现
is_volatile_node(易变表达式返回true)。
按此清单逐项排查,即可完成向 DataFusion 50.0.0 的平滑升级。对于每项变更的更深入讨论与背景细节,可继续查阅上文引用的对应仓库源码路径与上游 issue/PR。
- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
相关推荐
Social Analyzer安全加固终极指南:10个关键防护策略防止未授权访问与滥用
Social Analyzer安全加固终极指南:10个关键防护策略防止未授权访问与滥用 Social Analyzer是一款强大的开源工具,可通过API、CLI
大数据数据分析后端Apache DataFusion 53.0.0 升级指南:核心 API 破坏性变更与完整迁移手册
Apache DataFusion 53.0.0 升级指南:核心 API 破坏性变更与完整迁移手册 导读 本文以 Apache DataFusion 53.0.
大数据数据分析后端探索no-neck-pain.nvim架构:核心组件与实现原理
探索no neck pain.nvim架构:核心组件与实现原理 架构概览 no neck pain.nvim是一款旨在将当前聚焦的缓冲区居中显示在屏幕中间的Ne
大数据数据仓库OLAP批处理后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考