DataFusion 56.0.0
升级指南
DataFusion 56.0.0
注意: DataFusion 56.0.0 尚未发布。本节提供的信息涉及已合并到主分支、并将在该版本中发布的内容与变更。
将 arrow/parquet 升级到 60.0.0,将 object_store 升级到 0.14.2
DataFusion 56.0.0 使用 arrow 和 parquet 60.0.0,以及 object_store 0.14.2。如果你直接依赖这些 crate,可能需要更新 Cargo.toml。
有关这些版本中的破坏性变更详情,请参阅 Arrow 60.0.0 发布说明和 object_store 0.14.2 升级指南。
字段和 schema 元数据现在使用 arrow 的 Metadata 类型
Arrow 60 引入了专门的 Metadata 类型,用于替代 HashMap<String, String> 来表示字段和 schema 的元数据。继该变更之后,DataFusion 的多个 API 也改用了 Metadata,例如 DFSchema::metadata 和 ExprSchema::metadata 返回 &Metadata,CastExpr::target_metadata 和 TryCastExpr::target_metadata 返回 Option<&Metadata>。
迁移指南:
// Before
let meta: &HashMap<String, String> = schema.metadata();
let field = Field::new("a", DataType::Int64, true)
.with_metadata([("k".to_string(), "v".to_string())].into());
// After
let meta: &Metadata = schema.metadata();
let field = Field::new("a", DataType::Int64, true)
.with_metadata(Metadata::new().with("k", "v"));Metadata 支持像 HashMap 一样的 .get()、.iter()、.is_empty() 和 .extend(),并通过 From 与 HashMap<String, String>、BTreeMap<String, String> 相互转换,因此大多数调用点只需更改类型即可。当仍然需要 HashMap 时,请使用 FieldMetadata::to_hashmap。
缺失的 Parquet null_count 被视为未知
DataFusion 现在会将省略的 Parquet null_count 统计信息保留为未知。此前,它会假设空值数为零,这可能导致 IS NULL 过滤条件丢弃本应匹配的 NULL 行,或使 ORDER BY ... NULLS FIRST LIMIT 查询产生错误结果,也可能基于元数据计算出偏高的 COUNT(column)。
对受影响文件执行的查询可能会读取或排序更多数据,因为以下优化都需要已知的空值计数:
- 为
IS NULL过滤条件裁剪行组。 - 仅使用文件元数据计算
COUNT(column)。 - 使用动态的
NULLS FIRSTTopK 过滤器裁剪行组。 - 在已排序且不重叠的文件上消除可空列的排序。此类查询现在可能会保留完整的
SortExec,即使数据中并不包含 NULL。
53.1.0 之前的 parquet-rs 版本会省略为零的空值计数。这包括使用 DataFusion 42.1.0 之前版本所依赖的旧版 parquet-rs 生成的文件。Arrow devlive-community/knowforge#6490 修改了写入器,使其记录已知的零计数。带有显式计数的文件保持原有行为不变。
要在 DataFusion CLI 中查找具有边界统计信息但没有空值计数的列块,请运行以下查询:
SELECT row_group_id, path_in_schema
FROM parquet_metadata('data.parquet')
WHERE (stats_min IS NOT NULL OR stats_max IS NOT NULL)
AND stats_null_count IS NULL;使用当前 DataFusion 版本重写受影响的文件,例如在启用 Parquet 统计信息的情况下使用 COPY ... TO,可以补录缺失的计数并恢复依赖这些计数的优化。重写已排序的文件时,请保留所需的数据顺序以及顺序元数据。
ForeignSession::create_physical_plan 不再受支持
ForeignSession::create_physical_plan 不再转发给拥有该会话的库,而是返回 NotImplemented 错误,因为转发可能会重新进入已安装的外部计划器(planner),而旧回调返回的执行计划句柄无法为向下转型(downcasting)恢复本地 Rust 类型标识。原始的 FFI_SessionRef 回调槽位仍然保留,以兼容 DataFusion 55 的使用者,但调用该回调也会返回相同的错误。
拥有会话的库应当改为在安装外部计划器之前,将其原始计划器导出为 datafusion_ffi::query_planner::FFI_QueryPlanner。外部计划器可以持有并调用该句柄,以接收一个以本地类型标识重建的序列化物理计划。完整的委托模式请参阅 datafusion_ffi::query_planner 模块文档。ForeignSession::query_planner、optimize 和 physical_optimizers 仍会跨 FFI 边界转发到拥有该会话的会话对象。
GroupColumn 现在要求实现 values_preserving
公共 GroupColumn trait 的自定义实现必须实现 values_preserving。该方法在不改变所存储的值及其分组索引的情况下返回选中的行。它保留所要求的顺序,并支持重复的索引。
迁移指南:
实现 values_preserving,在读取之前调用 selection.validate_num_groups(self.len())?,并按照 selection.iter() 的顺序返回行,且不改变构建器(builder)的状态。
datafusion-proto:常见的选项与约束转换可能失败
CsvOptions、JsonOptions、ParquetCdcOptions、Constraint 和 Constraints 的 Protobuf 转换现在会拒绝无法用 usize 表示的整数值。它们原本不可能失败的 From 实现已替换为 TryFrom。ParquetOptions 和 TableParquetOptions 现有的 TryFrom 转换现在也会校验所有以 usize 支撑的字段。没有 constraint_mode 的约束会返回错误,而不再发生 panic。
迁移指南:
// Before
let csv = CsvOptions::from(&proto_csv);
let cdc = ParquetCdcOptions::from(proto_cdc);
let constraints: Constraints = proto_constraints.into();
// After
let csv = CsvOptions::try_from(&proto_csv)?;
let cdc = ParquetCdcOptions::try_from(proto_cdc)?;
let constraints = Constraints::try_from(proto_constraints)?;详情参见 issue devlive-community/knowforge#24170。
ExecutionOptions 新增 enable_nlj_coordinated_fallback 字段
ExecutionOptions 新增了一个公开字段 enable_nlj_coordinated_fallback: bool(默认为 true),用于控制内存受限的 NestedLoopJoinExec 回退实现是否在各探测分区之间共享每个分块的构建侧状态。正是这一机制使得当右侧存在多个分区时,LEFT、LEFT SEMI、LEFT ANTI、LEFT MARK 和 FULL 连接可以溢写(spill),而不是以 ResourcesExhausted 失败。
受影响的用户:
- 使用穷举式结构体字面量构造
ExecutionOptions(或ConfigOptions)的用户。读取配置或执行SET不受影响。 - 将每个输出分区作为独立任务运行的分布式引擎。该协调机制假设所有探测分区都运行在同一进程中;如果每个任务都有自己的协调器,共享的探测线程计数器永远不会归零,回退流程将陷入停滞。此类引擎应将该标志设置为
false,这样受影响的连接类型将保持原有的快速失败(fail-fast)行为。
迁移指南:
设置新增字段,或从 Default 填充:
// Before
ExecutionOptions {
batch_size: 8192,
// ... every other field ...
}
// After: set it explicitly
ExecutionOptions {
batch_size: 8192,
enable_nlj_coordinated_fallback: true,
// ... every other field ...
}
// After: or let the remaining fields come from Default, which also keeps
// future field additions from breaking the literal
ExecutionOptions {
batch_size: 8192,
..Default::default()
}分布式引擎的退出机制:
let mut config = SessionConfig::new();
config.options_mut().execution.enable_nlj_coordinated_fallback = false;datafusion.optimizer.use_statistics_registry 已废弃并被忽略
配置项 datafusion.optimizer.use_statistics_registry 已废弃,现在会被忽略(设置该配置会触发一条废弃警告)。在单次统计信息遍历过程中,始终会查询可插拔的 StatisticsRegistry:在会话上注册的提供者会直接生效。如果没有注册任何提供者,注册表便等同于空操作,因此默认行为与之前的默认值(use_statistics_registry = false)保持一致。
受影响的用户:
- 任何设置过
datafusion.optimizer.use_statistics_registry的用户。该设置仍被接受(并会发出废弃警告),但已不再产生任何效果,并将在未来的版本中移除。
迁移指南:
- 删除所有
use_statistics_registry设置。 - 若要使用统计信息提供者,请通过
SessionStateBuilder::with_statistics_registry(...)在会话上注册它们。若要保留该标志此前启用的内置 NDV 感知提供者,请注册StatisticsRegistry::default_with_builtin_providers()。 - 若要选择不使用,不注册任何提供者即可(这是默认行为)。
生成的 protobuf JsonWriterOptions 新增了 compression_level 字段
生成的公共 protobuf 结构体 JsonWriterOptions 现在包含一个可选的 compression_level 字段,使 JSON sink 计划在序列化过程中能够保留显式配置的压缩级别。
受影响的用户:
- 使用穷举结构体字面量构造生成的
JsonWriterOptions值的用户。
迁移指南:
显式设置 compression_level,或通过 Default 填充它:
// Before
JsonWriterOptions { compression }
// After
JsonWriterOptions {
compression,
compression_level: None,
}
// or
JsonWriterOptions {
compression,
..Default::default()
}Protobuf 线格式保持向后兼容。
详情参见 PR devlive-community/knowforge#24945。
生成的 protobuf Parquet 结构体新增了状态字段
生成的公开 protobuf 结构体 ParquetScanExecNode 现在携带一个可选的 metadata_size_hint,ParquetSink 现在携带一个可选的 sorting_columns。这些字段会在物理计划序列化过程中保留对应的 Parquet source 和 sink 设置。
受影响的用户:
- 使用穷举式结构体字面量构造上述任一生成结构体的用户。
迁移指南:
在 ParquetScanExecNode 字面量中添加 metadata_size_hint: None,在 ParquetSink 字面量中添加 sorting_columns: None,以保持之前的行为;或者对未指定的字段使用 ..Default::default()。
Protobuf 线格式保持向后兼容。
详情参见 PR devlive-community/knowforge#25057。
生成的 protobuf CsvWriterOptions 发生变更
生成的 CsvWriterOptions 现在包含 compression_level、timestamp_tz_format 和 terminator。其原有的四个格式字段从 String 变为 Option<String>,以区分未设置与显式设置为空字符串的情况。这会影响穷举式结构体字面量以及对格式字段的直接访问。
将现有的格式值包裹在 Some 中,或使用 None,并初始化新增字段(显式初始化或通过 Default 初始化):
// Before (unchanged fields omitted)
CsvWriterOptions {
date_format: "%Y-%m-%d".to_string(),
datetime_format: String::new(),
timestamp_format: String::new(),
time_format: String::new(),
// ...
}
// After
CsvWriterOptions {
date_format: Some("%Y-%m-%d".to_string()),
datetime_format: None,
timestamp_format: None,
time_format: None,
compression_level: None,
timestamp_tz_format: None,
terminator: Vec::new(),
// ...
}Protobuf 线格式保持向后兼容。
详情请参阅 PR devlive-community/knowforge#25058。
DefaultStatisticsProvider 已弃用
datafusion_physical_plan::operator_statistics::DefaultStatisticsProvider 已被弃用。它是多余的:当 provider 链委派或为空时,统计遍历(StatisticsContext)会回退到各算子的 statistics_from_inputs,因此不再需要一个末尾的「默认」provider。它也已不再属于 StatisticsRegistry::default_with_builtin_providers() 的一部分。
迁移指南:
- 从任何自定义 provider 链中移除
DefaultStatisticsProvider,不再注册任何末尾 provider(遍历会自行回退)。
StatisticsRegistry::compute 和 compute_base 已弃用
请改用遍历(walk):
StatisticsContext::new_with_registry(registry)
.compute_extended(plan, &StatisticsArgs::new())?; // or .compute(...) for core Statisticsfloor 与 ceil UDF 的 API 变更
floor 和 ceil UDF 的输出类型已从与输入完全相同的类型改为位宽相同的重新缩放后的类型。例如,对于输入 Decimal32(7,2),floor 现在返回 Decimal32(6,0),其中新精度为 p - s + 1,与 Spark 的行为一致。详见 [devlive-community/knowforge#24703]。
受影响的用户:
- 将这些 UDF 的查询结果存储到固定 schema 中的用户
- 依赖
arrow_typeof获取这些 UDF 返回类型的用户
迁移指南:
修改期望的类型,或用 CAST 包裹该表达式。建议不要依赖十进制数的精确精度和标度。
map_extract / element_at 对不存在的键返回空列表
map_extract(及其别名 element_at)此前在键不存在于 map 中时,会返回包含 NULL 的单元素列表。现在它返回空列表,与文档描述的行为以及 DuckDB 保持一致。同时还有两个相关场景的行为也发生了变化,同样与 DuckDB 保持一致:
NULL的 map 输入现在产生NULL,而不是[NULL]。NULL的查询键现在产生[],而不是[NULL]。
值为 NULL 的键仍会返回 [NULL],因此现在可以区分「键不存在」和「值为 NULL」这两种情况。
迁移指南:
-- Before
SELECT map_extract(MAP {'a': 1}, 'missing'); -- [NULL]
-- After
SELECT map_extract(MAP {'a': 1}, 'missing'); -- []那些假定结果恰好只有一个元素的表达式——例如通过检查其长度、对其执行展开(unnest)操作,或将其与 [NULL] 进行比较——现在应将空列表视为键不存在的情况。
详情请参阅 issue devlive-community/knowforge#24981 和 issue devlive-community/knowforge#24983。
DataFrame::from_columns 接受 IntoIterator
DataFrame::from_columns 现在接受任意的 IntoIterator<Item = (&str, ArrayRef)>,而不再只接受 Vec<(&str, ArrayRef)>。
// Existing Vec usage continues to work
let df = DataFrame::from_columns(vec![
("id", id),
("name", name),
])?;
// Arrays can now be used directly
let df = DataFrame::from_columns([
("id", id),
("name", name),
])?;大多数使用 Vec 的现有调用点无需修改。依赖 DataFrame::from_columns 精确非泛型函数签名的代码,可能需要更新以适配新的泛型 API。
聚合排序状态字段名称现按命名空间组织
作为顶层聚合状态字段暴露的排序表达式,现在使用由聚合名称及其排序位置派生出的名称。
例如:
first_value(value)[ordering_0]现在改为使用限定列名,而不是不带限定的排序字段名,例如:
timestamp@0这保证了聚合状态字段名的唯一性,并使得 DFSchema::check_names() 得以被强制执行。
受影响的对象:
直接检查聚合状态字段名的用户或集成,包括自定义 UDAF 和 FFI 集成。
物理过滤下推按位置解析列
datafusion_physical_plan::filter_pushdown::ChildFilterDescription::from_child 和 FilterDescription::from_children 现在按位置解析过滤列,而不再按名称查找。这可以防止子模式包含重复列名时产生错误结果。from_child 要求每个被引用位置上的子字段与过滤列同名。
ChildFilterDescription::from_child_with_allowed_indices 已被弃用,但仍保留其原先基于名称映射到第一个匹配子字段的行为。请迁移到 from_child_with_column_mapping,因为当子模式包含重复字段名时,按名称解析存在歧义。
迁移指南:
当父模式与子模式的列位置和名称一致时,请使用 from_child(多个子节点则使用 from_children)。当某个节点会投影、重排或重命名列时,请使用 from_child_with_column_mapping,并显式提供从父节点输出索引到子节点输入索引的映射。未出现在映射中的列无法下推。
例如,若父节点输出 [a, b],子节点输出 [b, a],则针对 a@0 的过滤必须映射到子列 a@1:
use std::collections::{HashMap, HashSet};
use datafusion_physical_plan::filter_pushdown::ChildFilterDescription;
// Before: allow parent column 0 and resolve "a" by name in the child.
let description = ChildFilterDescription::from_child_with_allowed_indices(
&parent_filters,
HashSet::from([0]),
&child,
)?;
// After: explicitly map parent column 0 to child column 1.
let description = ChildFilterDescription::from_child_with_column_mapping(
&parent_filters,
HashMap::from([(0, 1)]),
&child,
)?;默认 datafusion.sql_parser.recursion_limit 由 50 提高到 51
DataFusion 现在使用 sqlparser 0.63.0,该版本为更多解析函数(例如数据类型解析和 INTERVAL 解析)增加了递归保护。因此,同一条 SQL 语句会比使用之前的 sqlparser 版本消耗稍多的解析器递归预算。具体来说,每个表达式现在需要多出一层的递归深度。
为避免拒绝此前能够成功解析的查询,datafusion.sql_parser.recursion_limit 配置项以及 DFParserBuilder::with_recursion_limit 的默认值已从 50 提高到 51。
如果你显式地将 datafusion.sql_parser.recursion_limit 设置为接近查询深度的值,升级后可能会遇到 RecursionLimitExceeded 错误,此时应相应地调高该限制。
详情请参阅 PR devlive-community/knowforge#25278。
ensure_distribution 已弃用,改用 ensure_distribution_with_stats
datafusion_physical_optimizer::enforce_distribution::ensure_distribution 现在从调用方获取其 StatisticsContext,因此一个上下文(及其记忆化缓存)可以由整个自底向上的遍历过程共享。旧的双参数形式被保留为一个已弃用的包装器,它在每次调用时都会分配一个新的上下文,从而导致每个共享子树的统计信息会针对其每个祖先重新计算一次。
// Before
let plan = DistributionContext::new_default(plan)
.transform_up(|ctx| ensure_distribution(ctx, config))
.data()?;
// After: one context for the whole walk. `StatsCache` is keyed by raw plan-node
// pointers, so reset it whenever a node's plan pointer actually changed.
let stats_ctx = StatisticsContext::new();
let plan = DistributionContext::new_default(plan)
.transform_up(|ctx| {
let before = Arc::clone(&ctx.plan);
let result = ensure_distribution_with_stats(ctx, config, &stats_ctx)?;
if !Arc::ptr_eq(&before, &result.data.plan) {
stats_ctx.reset_cache();
}
Ok(result)
})
.data()?;只需要先前行为的调用方,可以继续使用已弃用的形式,也可以在每次调用时传入新构造的 StatisticsContext。
MergeIntoOp 需要 SQL 可见的目标限定符
MergeIntoOp 现在将 SQL 可见的目标限定符与目标提供者标识分开存储。这样可以避免带别名的目标被误认为是使用了目标表名的源关系。
该结构体现在标记为非穷举(non-exhaustive)。请将 55.0 版本中的结构体字面量替换为 MergeIntoOp::new:
// DataFusion 55.0
let op = MergeIntoOp { on, clauses };
// DataFusion 56.0
let op = MergeIntoOp::new(target_qualifier, on, clauses);target_qualifier 必须是 MERGE 表达式中引用的关系名称:存在别名时为该目标别名,否则为目标表引用。DmlStatement::table_name 仍表示提供者/目录标识。
MergeIntoOpNode protobuf 字面量必须提供 target_qualifier
生成的 MergeIntoOpNode 类型新增了一个可选字段。使用结构体字面量构造该生成类型的代码必须提供该字段:
// DataFusion 55.0
let node = protobuf::MergeIntoOpNode { on, clauses };
// DataFusion 56.0
let node = protobuf::MergeIntoOpNode {
on,
clauses,
target_qualifier: Some(protobuf::TableReference::from(target_qualifier)),
};线协议兼容性是单向的:
- 56.0 的读取方可以接受 55.0 的负载。当
target_qualifier缺失时,它会回退到DmlNode.table_name,与 55.0 的表示方式一致。 - 55.0 的读取方不能消费保留了别名的 56.0 MERGE 负载。它会忽略这个未知字段,但无法保留表达式所需的限定符,从而可能导致解析失败或错误的重新绑定。
评论
登录后参与评论
KnowForge