DataFusion 53.0.0
升级指南
DataFusion 53.0.0
将 arrow/parquet 升级到 58.0.0、object_store 升级到 0.13.0
DataFusion 53.0.0 使用 arrow 和 parquet 58.0.0 以及 object_store 0.13.0。如果你直接依赖这些 crate,可能需要相应更新你的 Cargo.toml。
有关这些版本中的破坏性变更详情,请参阅 Arrow 58.0.0 发布说明和 object_store 0.13.0 升级指南。
ExecutionPlan::statistics 已移除
已废弃的 ExecutionPlan::statistics() 方法已被移除。如果你实现了自定义的 ExecutionPlan,请从实现中删除该方法,改为实现 partition_statistics()。
修改前:
impl ExecutionPlan for MyExec {
// ...
fn statistics(&self) -> Result<Statistics> {
Ok(Statistics::new_unknown(&self.schema()))
}
}之后:
impl ExecutionPlan for MyExec {
// ...
fn partition_statistics(&self, _partition: Option<usize>) -> Result<Statistics> {
Ok(Statistics::new_unknown(&self.schema()))
}
}如果没有分区特定的统计数据,则对 None 和任意分区索引都返回相同的值。
ExecutionPlan::properties 现在返回 &Arc<PlanProperties>
现在,ExecutionPlan::properties() 返回的是 &Arc<PlanProperties>,而不是引用。这样可以低成本地克隆属性,并在多个 ExecutionPlan 之间复用。它还使得优化 [ExecutionPlan::with_new_children] 成为可能——当子计划没有变化时复用属性,从而显著减少复杂查询的规划时间。
ExecutionPlan::with_new_children
要进行迁移,在所有 ExecutionPlan 实现中,你很可能需要将存储的 PlanProperties 包装在 Arc 中:
- cache: PlanProperties,
+ cache: Arc<PlanProperties>,
...
- fn properties(&self) -> &PlanProperties {
+ fn properties(&self) -> &Arc<PlanProperties> {
&self.cache
}为了提升自定义 ExecutionPlan 实现中 with_new_children 的性能,你可以使用新的宏:check_if_same_properties。要使其生效,你需要实现函数 with_new_children_and_same_properties,其语义与 with_new_children 相同,但以子计划的属性未发生变化为前提进行操作。
以下是一个为 ProjectionExec 支持该优化的示例:
impl ProjectionExec {
+ fn with_new_children_and_same_properties(
+ &self,
+ mut children: Vec<Arc<dyn ExecutionPlan>>,
+ ) -> Self {
+ Self {
+ input: children.swap_remove(0),
+ metrics: ExecutionPlanMetricsSet::new(),
+ ..Self::clone(self)
+ }
+ }
}
impl ExecutionPlan for ProjectionExec {
fn with_new_children(
self: Arc<Self>,
mut children: Vec<Arc<dyn ExecutionPlan>>,
) -> Result<Arc<dyn ExecutionPlan>> {
+ check_if_same_properties!(self, children);
ProjectionExec::try_new(
self.projector.projection().into_iter().cloned(),
children.swap_remove(0),
)
.map(|p| Arc::new(p) as _)
}
}PlannerContext 外部查询 schema API 现已改用栈
PlannerContext 不再存储单个 outer_query_schema,而是维护一个外部关系 schema 的栈,使嵌套子查询能够访问非相邻的外部关系。
修改前:
let old_outer_query_schema =
planner_context.set_outer_query_schema(Some(input_schema.clone().into()));
let sub_plan = self.query_to_plan(subquery, planner_context)?;
planner_context.set_outer_query_schema(old_outer_query_schema);之后:
planner_context.append_outer_query_schema(input_schema.clone().into());
let sub_plan = self.query_to_plan(subquery, planner_context)?;
planner_context.pop_outer_query_schema();HashJoinExec::try_new 新增 null_aware 参数
HashJoinExec::try_new 现在额外接受一个 null_aware: bool 参数。该标志用于空值感知反连接(null-aware anti join),例如为 NOT IN 子查询生成的执行计划。
大多数调用方应传入 false,即保持原有行为。仅在执行空值感知的 JoinType::LeftAnti 连接时才应传入 true。
FileSinkConfig 新增 file_output_mode
FileSinkConfig 现在包含一个 file_output_mode: FileOutputMode 字段,用于控制单文件输出还是目录输出。任何通过结构体字面量构造 FileSinkConfig 的代码都必须初始化该字段。
FileOutputMode 枚举有三个变体:
Automatic(默认):根据 URL 推断输出模式(基于扩展名/末尾/的启发式判断)SingleFile:将内容写入输出路径所指定的单个文件Directory:写入目录,文件名自动生成
改动前:
FileSinkConfig {
// ...
file_extension: "parquet".into(),
}之后:
use datafusion_datasource::file_sink_config::FileOutputMode;
FileSinkConfig {
// ...
file_extension: "parquet".into(),
file_output_mode: FileOutputMode::Automatic,
}SimplifyInfo trait 已移除,SimplifyContext 现在使用构建器式 API
SimplifyInfo trait 已被移除,取而代之的是具体的 SimplifyContext 结构体。这简化了表达式简化的 API,并消除了对 trait 对象的需求。
受影响的用户:
- 实现了自定义
SimplifyInfo的用户 - 为自定义标量函数实现了
ScalarUDFImpl::simplify()的用户 - 直接使用
SimplifyContext或ExprSimplifier的用户
破坏性变更:
SimplifyInfotrait 已被完全移除SimplifyContext不再接受&ExecutionProps- 它现在使用带有直接字段的构建器式 APIScalarUDFImpl::simplify()现在接受&SimplifyContext而不是&dyn SimplifyInfo- 时间相关函数的简化(例如
now())现在是可选的 - 如果query_execution_start_time为None,这些函数将不会被简化
迁移指南:
如果你实现了自定义的 SimplifyInfo:
迁移前:
impl SimplifyInfo for MySimplifyInfo {
fn is_boolean_type(&self, expr: &Expr) -> Result<bool> { ... }
fn nullable(&self, expr: &Expr) -> Result<bool> { ... }
fn execution_props(&self) -> &ExecutionProps { ... }
fn get_data_type(&self, expr: &Expr) -> Result<DataType> { ... }
}之后:
直接通过构建器风格的 API 使用 SimplifyContext:
let context = SimplifyContext::default()
.with_schema(schema)
.with_config_options(config_options)
.with_query_execution_start_time(Some(Utc::now())); // or use .with_current_time()如果实现了 ScalarUDFImpl::simplify():
此前:
fn simplify(
&self,
args: Vec<Expr>,
info: &dyn SimplifyInfo,
) -> Result<ExprSimplifyResult> {
let now_ts = info.execution_props().query_execution_start_time;
// ...
}之后:
fn simplify(
&self,
args: Vec<Expr>,
info: &SimplifyContext,
) -> Result<ExprSimplifyResult> {
// query_execution_start_time is now Option<DateTime<Utc>>
// Return Original if time is not set (simplification skipped)
let Some(now_ts) = info.query_execution_start_time() else {
return Ok(ExprSimplifyResult::Original(args));
};
// ...
}如果你通过 ExecutionProps 创建 SimplifyContext:
修改前:
let props = ExecutionProps::new();
let context = SimplifyContext::new(&props).with_schema(schema);之后:
let context = SimplifyContext::default()
.with_schema(schema)
.with_config_options(config_options)
.with_current_time(); // Sets query_execution_start_time to Utc::now()更多详情请参阅 SimplifyContext 文档。
结构体现在要求字段名存在重叠
DataFusion 的结构体(struct)类型转换机制此前允许在字段名不同但字段数量相同的结构体之间进行转换。这种「按位置回退」的行为可能悄悄地错配字段,从而导致数据损坏。
破坏性变更:
从 DataFusion 53.0.0 开始,结构体转换现在要求源结构体与目标结构体之间至少有一个重叠的字段名。若字段名没有重叠,转换将在计划阶段被拒绝,并给出明确的错误信息。
受影响的用户:
- 在字段名完全没有重叠的结构体之间进行转换的应用程序
- 依赖按位置匹配结构体字段的查询(例如,仅凭位置将
struct(x, y)转换为struct(a, b)) - 以编程方式构造或转换结构体列的代码
迁移指南:
如果你遇到如下错误:
Cannot cast struct with 2 fields to 2 fields because there is no field name overlap你必须显式地重命名或映射字段,以确保至少有一个字段名能够匹配。以下是常见模式:
示例 1:源字段与目标字段的名称一致(基于名称的转换)
成功场景(字段名对齐):
-- source_col has schema: STRUCT<x INT, y INT>
-- Casting to the same field names succeeds (no-op or type validation only)
SELECT CAST(source_col AS STRUCT<x INT, y INT>) FROM table1;示例 2:源字段名与目标字段名不一致(迁移场景)
当前失败的情况(字段名没有任何重叠):
-- source_col has schema: STRUCT<a INT, b INT>
-- This FAILS because there is no field name overlap:
-- ❌ SELECT CAST(source_col AS STRUCT<x INT, y INT>) FROM table1;
-- Error: Cannot cast struct with 2 fields to 2 fields because there is no field name overlap迁移选项(必须使名称保持一致):
选项 A:使用 struct 构造器进行显式字段映射
-- source_col has schema: STRUCT<a INT, b INT>
-- Use STRUCT_CONSTRUCT with explicit field names
SELECT STRUCT_CONSTRUCT(
'x', source_col.a,
'y', source_col.b
) AS renamed_struct FROM table1;选项 B:在 CAST 目标中重命名以匹配源名称
-- source_col has schema: STRUCT<a INT, b INT>
-- Cast to target with matching field names
SELECT CAST(source_col AS STRUCT<a INT, b INT>) FROM table1;示例 3:在 Rust API 中使用结构体构造函数
如果需要以编程方式映射字段,可以显式地构建目标结构体:
// Build the target struct with explicit field names
let target_struct_type = DataType::Struct(vec![
FieldRef::new("x", DataType::Int32),
FieldRef::new("y", DataType::Utf8),
]);
// Use struct constructors rather than casting for field mapping
// This makes the field mapping explicit and unambiguous
// Use struct builders or row constructors that preserve your mapping logic本次变更的原因:
- 安全性: 字段名称现在是结构体兼容性的首要契约
- 明确性: 避免因位置假设而导致的静默数据错位
- 一致性: 与 DuckDB 的行为保持一致,并与其他基于名称匹配的 SQL 引擎保持统一
- 可调试性: 错误现在会在计划阶段报出,而不再表现为静默的数据损坏
更多详情请参阅 Issue devlive-community/knowforge#19841 和 PR devlive-community/knowforge#19955。
FilterExec 的构建方法已弃用
FilterExec 上的以下方法已弃用,建议改用 FilterExecBuilder:
with_projection()with_batch_size()
受影响的用户:
- 创建
FilterExec实例并使用这些方法进行配置的用户
迁移指南:
请使用 FilterExecBuilder,而不是在 FilterExec 上链式调用这些方法:
修改前:
let filter = FilterExec::try_new(predicate, input)?
.with_projection(Some(vec![0, 2]))?
.with_batch_size(8192)?;之后:
let filter = FilterExecBuilder::new(predicate, input)
.with_projection(Some(vec![0, 2]))
.with_batch_size(8192)
.build()?;构建器模式效率更高,因为它在 build() 期间只计算一次属性,而不是在每次方法调用时重新计算。
注意:with_default_selectology() 并未被弃用,因为它只是更新字段值,不需要构建器模式的额外开销。
新增 Protobuf 转换 trait
datafusion-proto crate 中新增了一个 trait:PhysicalProtoConverterExtension。它用于控制物理计划及其表达式与 protobuf 对应形式之间的相互转换过程。转换相关的方法现在需要一个额外的参数。
与该 crate 交互的主要 API 并未修改,因此大多数用户无需做出任何更改。如果你确实需要这个 trait,可以使用 DefaultPhysicalProtoConverter 实现。
例如,要转换一个排序表达式的 protobuf 节点,你可以进行如下更新:
更新前:
let sort_expr = parse_physical_sort_expr(
sort_proto,
ctx,
input_schema,
codec,
);之后:
let converter = DefaultPhysicalProtoConverter {};
let sort_expr = parse_physical_sort_expr(
sort_proto,
ctx,
input_schema,
codec,
&converter
);与将物理排序表达式转换为 protobuf 节点类似:
转换前:
let sort_proto = serialize_physical_sort_expr(
sort_expr,
codec,
);之后:
let converter = DefaultPhysicalProtoConverter {};
let sort_proto = serialize_physical_sort_expr(
sort_expr,
codec,
&converter,
);generate_series 与 range 表函数行为变更
generate_series 和 range 表函数在间隔(interval)无效时,现在会返回空集合,而不是抛出错误。这一行为与 PostgreSQL 等系统保持一致。
之前:
> select * from generate_series(0, -1);
Error during planning: Start is bigger than end, but increment is positive: Cannot generate infinite series
> select * from range(0, -1);
Error during planning: Start is bigger than end, but increment is positive: Cannot generate infinite series现在:
> select * from generate_series(0, -1);
+-------+
| value |
+-------+
+-------+
0 row(s) fetched.
> select * from range(0, -1);
+-------+
| value |
+-------+
+-------+
0 row(s) fetched.array_remove、array_remove_n、array_remove_all 在元素参数为 NULL 时现在返回 NULL
此前,调用 array_remove(array, NULL) 会尝试在数组中匹配并移除 NULL 元素,返回去除 NULL 后的数组。现在,当传入 NULL 作为待移除的元素时,该函数会返回 NULL,这与标准 SQL 的 NULL 传播语义保持一致。
同样的变更适用于 array_remove_n(别名:list_remove_n)和 array_remove_all(别名:list_remove_all)。
受影响的用户:
- 以 NULL 元素参数(字面量或列派生值)调用
array_remove、array_remove_n或array_remove_all的查询。
行为变更:
| 表达式 | 旧结果 | 新结果 |
|---|---|---|
array_remove(make_array(1, NULL, 2), NULL) | [1, 2] | NULL |
array_remove(make_array(1, NULL, 2, NULL), NULL) | [1, 2, NULL] | NULL |
array_remove_n(make_array(1, 2, 2, 1, 1), NULL, 2) | [1, 2, 2, 1, 1] | NULL |
array_remove_all(make_array(1, 2, 2, 1, 1), NULL) | [1, 2, 2, 1, 1] | NULL |
迁移指南:
如果你的查询依赖旧行为来去除数组中的 NULL,请改用带非 NULL 哨兵值的 array_remove_all,或改用 array_filter:
-- Before (removed NULLs from array):
SELECT array_remove(make_array(1, NULL, 2), NULL);
-- Old result: [1, 2]
-- After (returns NULL due to NULL propagation):
SELECT array_remove(make_array(1, NULL, 2), NULL);
-- New result: NULL详情请参阅 devlive-community/knowforge#21011 和 PR devlive-community/knowforge#21013。
评论
登录后参与评论
KnowForge