news 2026/9/25 1:55:38

Apache DataFusion 50.0.0 升级指南:从 Hive 分区自动推断到 UDF 特质重构的完整迁移手册

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache DataFusion 50.0.0 升级指南:从 Hive 分区自动推断到 UDF 特质重构的完整迁移手册
  • 大数据
  • 数据分析
  • 后端

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

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

本指南基于 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所有自定义函数实现者
异步 UDFAsyncScalarUDFImpl::invoke_async_with_args返回ColumnarValue并移除ConfigOptions参数注册远程/异步函数(如调用 LLM)的用户
投影表达式ProjectionExpr由元组别名改为具名字段结构体直接使用ProjectionExec的物理计划开发者
会话配置options()等 API 返回&Arc<ConfigOptions>访问会话配置的库使用者
物理计划新增ExecutionPlan::reset_state自定义ExecutionPlan实现者
表达式特质新增PhysicalExpr::is_volatile_node自定义PhysicalExpr实现者
数据源FileOpenFuture错误类型改为DataFusionErrorFileOpener自定义实现者
schema 重写schema_rewriter模块迁入新 crate依赖物理表达式适配器的用户
FFIUDAF 从return_type迁移到return_fieldFFI 库开发者与测试作者

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。涉及以下符号:

  • DefaultPhysicalExprAdapter
  • DefaultPhysicalExprAdapterFactory
  • PhysicalExprAdapter
  • PhysicalExprAdapterFactory

仓库当前结构印证了迁移结果:新 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 升级自查表

为便于对照执行,汇总一份升级自查清单:

  1. 工具链:确认 Rust ≥ 1.86.0(rust-toolchain.toml中当前主分支为 1.98.1,锁定 50.0.0 时以 1.86.0 为准)。
  2. 分区表行为:如不想要自动推断 Hive 分区,设置datafusion.execution.listing_table_factory_infer_partitions = false(默认true,config.rs)。
  3. UDF 特质:为ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl实现PartialEq、Eq、Hash,删除手写equals/hash_value;无状态函数可用正则批量加 derive。
  4. 异步 UDF:invoke_async_with_args返回ColumnarValue(用ColumnarValue::from包裹ArrayRef),并移除_option参数,改从args.config_options读取。
  5. 投影:ProjectionExpr改为结构体,用ProjectionExpr::new(expr, alias)构造、.expr/.alias访问、ProjectionExpr { expr, alias }匹配。
  6. 配置访问:标注类型的options()调用补.as_ref()(返回&Arc<ConfigOptions>)。
  7. schema 重写:import 改为datafusion_physical_expr_adapter::{...}。
  8. Arrow/Parquet:升级依赖到 56.0.0。
  9. 自定义 ExecutionPlan:按需实现reset_state(有状态或持引用时必须)。
  10. NLJ:确认不依赖旧实现的输入排序保持;必要时显式排序。
  11. LazyBatchGenerator:补充as_any()返回self。
  12. DataSource:按新签名重新实现try_swapping_with_projection。
  13. FileOpenFuture:错误处理改用DataFusionError。
  14. FFI UDAF:测试从return_type改为return_field;保证交互库 Rust 版本一致。
  15. PhysicalExpr:显式实现is_volatile_node(易变表达式返回true)。

按此清单逐项排查,即可完成向 DataFusion 50.0.0 的平滑升级。对于每项变更的更深入讨论与背景细节,可继续查阅上文引用的对应仓库源码路径与上游 issue/PR。

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

【免费下载链接】datafusion

Apache DataFusion SQL Query Engine

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

相关推荐

上一篇:jwt库进阶教程:自定义Claims与高级验证策略
下一篇:Plus Jakarta Sans 字体终极使用指南:如何为你的项目选择完美的开源字体

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 1:54:20

突破100万token:长上下文大模型技术完全解析与TaoToken配置实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/25 1:51:35

ZXing批量生成DM二维码:工业追溯场景的实现与避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/25 1:51:24

ROCm多流调度实战:hipMemcpyAsync异步陷阱与拷贝计算重叠

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/25 1:50:31

网心云OES Plus刷Armbian后系统迁移至SATA硬盘扩容实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华