升级指南

DataFusion 50.0.0

师成师成· 更新于 2026-09-28· 阅读 20 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

升级指南

DataFusion 50.0.0

ListingTable 自动检测 Hive 分区表

DataFusion 50.0.0 在使用 ListingTableFactory 和 CREATE EXTERNAL TABLE 时会自动推断 Hive 分区。此前,创建 ListingTable 时,使用 Hive 分区方式的数据集(例如 /table_root/column1=value1/column2=value2/data.parquet)其 Hive 列不会体现在表的 schema 或数据中。可以通过将 datafusion.execution.listing_table_factory_infer_partitions 配置项设置为 false 来恢复之前的行为。更多详情请参见问题 devlive-community/knowforge#17049。

MSRV 更新至 1.86.0

最低支持的 Rust 版本(MSRV)已更新为 1.86.0。详情请参见 devlive-community/knowforge#17230。

ScalarUDFImpl、AggregateUDFImpl 和 WindowUDFImpl trait 现在要求实现 PartialEq、Eq 和 Hash trait

为了解决 ScalarUDFImpl::equals、AggregateUDFImpl::equals 和 WindowUDFImpl::equals 方法容易出错的问题,并使函数相等性能够被正确、简便地实现,equals 和 hash_value 方法已从 ScalarUDFImpl、AggregateUDFImpl 和 WindowUDFImpl trait 中移除。取而代之的是,任何实现了 ScalarUDFImpl、AggregateUDFImpl 或 WindowUDFImpl 的类型都必须实现 PartialEq、Eq 和 Hash trait。更多详情请参见问题 devlive-community/knowforge#16677。

大多数标量函数是无状态的,并且带有 signature 字段。可以使用正则表达式进行迁移:

  • 搜索 \[derive\(Debug\)\](\n *(pub )?struct \w+ \{\n *signature\: Signature\,\n *\}),
  • 替换为 #[derive(Debug, PartialEq, Eq, Hash)]$1,
  • 检查所有改动,确保只有函数结构体发生了变更。

AsyncScalarUDFImpl::invoke_async_with_args 返回 ColumnarValue

为了支持单值优化,并与其他用户自定义函数 API 保持一致,AsyncScalarUDFImpl::invoke_async_with_args 方法现在返回 ColumnarValue 而不是 ArrayRef。

升级时,请修改你的实现的返回类型。

impl AsyncScalarUDFImpl for AskLLM {
    async fn invoke_async_with_args(
        &self,
        args: ScalarFunctionArgs,
        _option: &ConfigOptions,
    ) -> Result<ColumnarValue> {
        ..
      return array_ref; // old code
    }
}

要返回 ColumnarValue

impl AsyncScalarUDFImpl for AskLLM {
    async fn invoke_async_with_args(
        &self,
        args: ScalarFunctionArgs,
        _option: &ConfigOptions,
    ) -> Result<ColumnarValue> {
        ..
      return ColumnarValue::from(array_ref); // new code
    }
}

详情参见 devlive-community/knowforge#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 }

这主要影响 ProjectionExec 的使用。

此次变更在 devlive-community/knowforge#17398 中完成。

SessionState、SessionConfig 和 OptimizerConfig 返回 &Arc<ConfigOptions> 而非 &ConfigOptions

为了提供对 ConfigOptions 更广泛的访问并减少所需的克隆,部分 API 被改为返回 &Arc<ConfigOptions> 而不是 &ConfigOptions。这样可以在多个线程之间共享同一个 ConfigOptions,而无需克隆整个 ConfigOptions 结构,除非需要修改它。

大多数用户不会受到此次变更的影响,因为 Rust 编译器通常会在需要时自动解引用 Arc。不过在某些情况下,你可能需要修改代码以显式调用 as_ref(),例如从

let optimizer_config: &ConfigOptions = state.options();

到

let optimizer_config: &ConfigOptions = state.options().as_ref();

参见 PR devlive-community/knowforge#16970

AsyncScalarUDFImpl::invoke_async_with_args 的 API 变更

AsyncScalarUDFImpl trait 的 invoke_async_with_args 方法已更新,移除了 _option: &ConfigOptions 参数以简化 API,因为现在可以通过 ScalarFunctionArgs 参数访问 ConfigOptions。

你可以按如下方式修改代码

impl AsyncScalarUDFImpl for AskLLM {
    async fn invoke_async_with_args(
        &self,
        args: ScalarFunctionArgs,
        _option: &ConfigOptions,
    ) -> Result<ArrayRef> {
        ..
    }
    ...
}

改为:

impl AsyncScalarUDFImpl for AskLLM {
    async fn invoke_async_with_args(
        &self,
        args: ScalarFunctionArgs,
    ) -> Result<ArrayRef> {
        let options = &args.config_options;
        ..
    }
    ...
}

Schema Rewriter 模块已迁移至新 Crate

schema_rewriter 模块及其相关符号已从 datafusion_physical_expr 迁移到新的 crate datafusion_physical_expr_adapter 中。受影响的符号如下:

  • DefaultPhysicalExprAdapter
  • DefaultPhysicalExprAdapterFactory
  • PhysicalExprAdapter
  • PhysicalExprAdapterFactory

如需升级,请将导入语句修改为:

use datafusion_physical_expr_adapter::{
    DefaultPhysicalExprAdapter, DefaultPhysicalExprAdapterFactory,
    PhysicalExprAdapter, PhysicalExprAdapterFactory
};

升级至 arrow 56.0.0 与 parquet 56.0.0

此版本的 DataFusion 将底层的 Apache Arrow 实现升级至 56.0.0 版本。更多详情请参阅发行说明。

新增 ExecutionPlan::reset_state

为了修复 DataFusion 49.0.0 中的一个缺陷——动态过滤器(目前仅在存在类似 ORDER BY ... LIMIT ... 的查询时生成)在递归查询中会产生错误结果——ExecutionPlan trait 新增了 reset_state 方法。

任何需要维护内部状态或对执行计划树中其他节点的引用的 ExecutionPlan,都应实现该方法以重置这些状态。更多详情以及 SortExec 的示例实现请参阅 devlive-community/knowforge#17028。

Nested Loop Join 无法保留输入排序顺序

Nested Loop Join 算子已从头重写,以提升性能和内存效率。根据微基准测试:与旧实现相比,此改动在极端情况下可带来最高 5 倍的加速,且内存占用仅为原来的 1%。

然而,新实现无法像旧版本那样保留输入的排序顺序。这是一个根本性的设计权衡,即以排序顺序保留为代价,优先保证性能和内存效率。

详情请参阅 devlive-community/knowforge#16996。

为 LazyBatchGenerator 添加 as_any() 方法

为便于进行 protobuf 序列化,LazyBatchGenerator trait 中新增了 as_any() 方法。这意味着你需要在自己的 LazyBatchGenerator 实现中添加 as_any():

impl LazyBatchGenerator for MyBatchGenerator {
    fn as_any(&self) -> &dyn Any {
        self
    }

    ...
}

详情请参阅 devlive-community/knowforge#17200。

重构 DataSource::try_swapping_with_projection

我们重构了 DataSource::try_swapping_with_projection,以简化该方法,并尽可能减少 ExecutionPlan 与 DataSource 抽象层之间的细节泄漏。对任何自定义 DataSource 进行重新实现都相对简单,详情请参阅 devlive-community/knowforge#17395。

FileOpenFuture 现在使用 DataFusionError 而非 ArrowError

FileOpenFuture 类型别名已更新,其错误类型由 ArrowError 改为 DataFusionError。此变更会影响 FileOpener trait 以及所有涉及文件流式操作的实现。

变更前:

pub type FileOpenFuture = BoxFuture<'static, Result<BoxStream<'static, Result<RecordBatch, ArrowError>>>>;

之后:

pub type FileOpenFuture = BoxFuture<'static, Result<BoxStream<'static, Result<RecordBatch>>>>;

如果你自定义实现了 FileOpener 或直接使用 FileOpenFuture,就需要将错误处理从 ArrowError 改为使用 DataFusionError。FileStreamState 枚举的 Open 变体也已相应更新。更多详情请参阅 devlive-community/knowforge#17397。

FFI 用户自定义聚合函数签名变更

用户自定义聚合函数的外部函数接口(FFI)签名已更新,现在会在底层聚合函数上调用 return_field 而非 return_type。这样做的目的是支持对这些聚合函数的元数据处理。此变更对大多数用户来说是透明的。如果你编写过直接调用 return_type 的单元测试,可能需要将其改为调用 return_field。

此次更新是对 FFI API 的破坏性变更。目前使用 FFI crate 的最佳实践是确保所有相互交互的库都使用相同的底层 Rust 版本。已创建 Issue devlive-community/knowforge#17374,用于讨论该接口的稳定化,以便这些库能够在不同的 DataFusion 版本之间通用。

详情请参阅 devlive-community/knowforge#17407。

新增 PhysicalExpr::is_volatile_node

我们为 PhysicalExpr 新增了一个方法,用于将某个 PhysicalExpr 标记为易变的:

impl PhysicalExpr for MyRandomExpr {
  fn is_volatile_node(&self) -> bool {
    true
  }
}

我们以默认值 false 发布了此功能,以尽量减少破坏性变更,但我们强烈建议 PhysicalExpr 的实现者选择加入这一行为,即使它返回 false。

你可以在 devlive-community/knowforge#17351 中查看更多讨论和示例实现。

评论

登录后参与评论

正在加载评论…