DataFusion 50.0.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 中。受影响的符号如下:
DefaultPhysicalExprAdapterDefaultPhysicalExprAdapterFactoryPhysicalExprAdapterPhysicalExprAdapterFactory
如需升级,请将导入语句修改为:
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 中查看更多讨论和示例实现。
评论
登录后参与评论
KnowForge