DataFusion 51.0.0
升级指南
DataFusion 51.0.0
arrow / parquet 已升级至 57.0.0
升级到 arrow 57.0.0 和 parquet 57.0.0
此版本的 DataFusion 将底层的 Apache Arrow 实现升级到 57.0.0 版本,同时升级了若干依赖 crate,包括 prost、tonic、pyo3 和 substrait。详情请参阅发行说明。
MSRV 已更新至 1.88.0
最低支持的 Rust 版本(MSRV)已更新为 1.88.0。
FunctionRegistry 暴露了两个额外方法
FunctionRegistry 暴露了两个额外方法 udafs 和 udwfs,分别用于返回已注册的用户自定义聚合函数和窗口函数的名称集合。升级时请实现这两个方法,使其返回已注册函数名的集合:
impl FunctionRegistry for FunctionRegistryImpl {
fn udfs(&self) -> HashSet<String> {
self.scalar_functions.keys().cloned().collect()
}
+ fn udafs(&self) -> HashSet<String> {
+ self.aggregate_functions.keys().cloned().collect()
+ }
+
+ fn udwfs(&self) -> HashSet<String> {
+ self.window_functions.keys().cloned().collect()
+ }
}datafusion-proto 在物理计划序列化方法中使用 TaskContext 而非 SessionContext
datafusion-proto 中处理物理计划序列化/反序列化的公开 API 方法发生了变更。
physical_plan_from_bytes、parse_physical_expr 等方法现在期望传入 TaskContext,而不是 SessionContext。
- let plan2 = physical_plan_from_bytes(&bytes, &ctx)?;
+ let plan2 = physical_plan_from_bytes(&bytes, &ctx.task_ctx())?;由于 TaskContext 中已包含 RuntimeEnv,诸如 try_into_physical_plan 之类的方法将不再带有显式的 RuntimeEnv 参数。
let result_exec_plan: Arc<dyn ExecutionPlan> = proto
- .try_into_physical_plan(&ctx, runtime.deref(), &composed_codec)
+. .try_into_physical_plan(&ctx.task_ctx(), &composed_codec)PhysicalExtensionCodec::try_decode() 现在期望接收 TaskContext,而不是 FunctionRegistry:
pub trait PhysicalExtensionCodec {
fn try_decode(
&self,
buf: &[u8],
inputs: &[Arc<dyn ExecutionPlan>],
- registry: &dyn FunctionRegistry,
+ ctx: &TaskContext,
) -> Result<Arc<dyn ExecutionPlan>>;有关更多详情,请参见 issue devlive-community/knowforge#17601。
SessionState 的 sql_to_statement 方法改用 Dialect 而非 str
datafusion::execution::session_state::SessionState 中定义的 sql_to_statement 方法的 dialect 参数已从 &str 更改为 &Dialect。Dialect 是定义在 datafusion-common crate 的 config 模块中的一个枚举,它为 SQL 方言的选择提供了类型安全和更好的校验。
将 ListingTable 重组到 datafusion-catalog-listing crate 中
长期以来一直有请求希望将 ListingTable 等功能从 datafusion crate 中移出,以加快构建速度。结构体 ListingOptions、ListingTable 和 ListingTableConfig 现在位于 datafusion-catalog-listing crate 中。它们在 datafusion crate 中被重新导出,因此对现有用户的影响应该是很小的。
有关更多详情,请参见 issue devlive-community/knowforge#14462 和 issue devlive-community/knowforge#17713。
将 ArrowSource 重组到 datafusion-datasource-arrow crate 中
为支持 issue devlive-community/knowforge#17713,ArrowSource 的代码已从 datafusion 核心 crate 中移出,放入其独立的 datafusion-datasource-arrow crate 中。这遵循了 AVRO、CSV、JSON 和 Parquet 数据源的模式。用户可能需要更新其路径以适应这些变更。
有关更多详情,请参见 issue devlive-community/knowforge#17713。
FileScanConfig::projection 重命名为 FileScanConfig::projection_exprs
FileScanConfig 中的 projection 字段已重命名为 projection_exprs,其类型也从 Option<Vec<usize>> 变为 Option<ProjectionExprs>。这一变更支持任意物理表达式而不仅仅是列索引,从而实现了更强大的投影下推能力。
对直接字段访问的影响:
如果你直接访问 projection 字段:
let config: FileScanConfig = ...;
let projection = config.projection;你应该更新到:
let config: FileScanConfig = ...;
let projection_exprs = config.projection_exprs;对构建器的影响:
FileScanConfigBuilder::with_projection() 方法已被弃用,请改用 with_projection_indices():
let config = FileScanConfigBuilder::new(url, file_source)
- .with_projection(Some(vec![0, 2, 3]))
+ .with_projection_indices(Some(vec![0, 2, 3]))
.build();注:with_projection() 仍然可用,但已被弃用,并将在未来版本中移除。
ProjectionExprs 是什么?
ProjectionExprs 是一种新类型,用于表示投影(projection)所对应的物理表达式列表。虽然它可以由列索引构建(with_projection_indices 内部就是这样做的),但它同样支持任意的物理表达式,从而能够实现扫描期间求值表达式等高级功能。
如果需要,你可以通过 ProjectionExprs 的方法来获取其中的列索引:
let projection_exprs: ProjectionExprs = ...;
// Get the column indices if the projection only contains simple column references
let indices = projection_exprs.column_indices();DESCRIBE query 支持
此前 DESCRIBE query 是 EXPLAIN query 的别名,用于输出查询的执行计划。在本次发布中,DESCRIBE query 现在输出查询计算得到的模式(schema),这与 DESCRIBE table_name 的行为保持一致。
datafusion.execution.time_zone 默认配置变更
datafusion.execution.time_zone 的默认值此前是字符串值 +00:00(GMT/祖鲁时间)。现已改为 Option<String>,默认值为 None。如果你想将时区改回之前的值,可以执行以下 SQL:
SET
TIMEZONE = '+00:00';此次改动旨在更好地支持在 now、current_date、current_time 和 to_timestamp 等标量 UDF 函数中使用默认时区。
引入 TableSchema 及 FileSource::with_schema() 方法的变更
datafusion-datasource crate 中引入了一个新的 TableSchema 结构体,以便更好地管理包含分区列的表结构。该结构体用于区分以下几种概念:
- 文件结构(File schema):磁盘上实际数据文件的结构
- 分区列(Partition columns):由目录结构派生的列(例如 Hive 风格分区)
- 表结构(Table schema):文件结构与分区列合并后的完整结构
作为此次改动的一部分,FileSource::with_schema() 方法的签名从接收 SchemaRef 变更为接收 TableSchema。
受影响的人群:
- 已实现自定义
FileSource的用户需要更新其代码 - 仅使用内置文件源(Parquet、CSV、JSON、AVRO、Arrow)的用户不受影响
自定义 FileSource 实现的迁移指南:
use datafusion_datasource::file::FileSource;
-use arrow::datatypes::SchemaRef;
+use datafusion_datasource::TableSchema;
impl FileSource for MyCustomSource {
- fn with_schema(&self, schema: SchemaRef) -> Arc<dyn FileSource> {
+ fn with_schema(&self, schema: TableSchema) -> Arc<dyn FileSource> {
Arc::new(Self {
- schema: Some(schema),
+ // Use schema.file_schema() to get the file schema without partition columns
+ schema: Some(Arc::clone(schema.file_schema())),
..self.clone()
})
}
}对于需要访问分区列的实现:
fn with_schema(&self, schema: TableSchema) -> Arc<dyn FileSource> {
Arc::new(Self {
file_schema: Arc::clone(schema.file_schema()),
partition_cols: schema.table_partition_cols().clone(),
table_schema: Arc::clone(schema.table_schema()),
..self.clone()
})
}注意:大多数 FileSource 实现只需要存储文件 schema(不含分区列),如第一个示例所示。存储全部三个 schema 组成部分的第二种模式通常只在高级场景中才需要,即需要针对不同操作访问不同的 schema 表示形式时(例如,ParquetSource 使用文件 schema 来构建裁剪谓词,但需要表 schema 来实现过滤下推逻辑)。
直接使用 TableSchema:
如果你正在构造 FileScanConfig,或处理表 schema 与分区列,现在可以直接使用 TableSchema:
use datafusion_datasource::TableSchema;
use arrow::datatypes::{Schema, Field, DataType};
use std::sync::Arc;
// Create a TableSchema with partition columns
let file_schema = Arc::new(Schema::new(vec![
Field::new("user_id", DataType::Int64, false),
Field::new("amount", DataType::Float64, false),
]));
let partition_cols = vec![
Arc::new(Field::new("date", DataType::Utf8, false)),
Arc::new(Field::new("region", DataType::Utf8, false)),
];
let table_schema = TableSchema::new(file_schema, partition_cols);
// Access different schema representations
let file_schema_ref = table_schema.file_schema(); // Schema without partition columns
let full_schema = table_schema.table_schema(); // Complete schema with partition columns
let partition_cols_ref = table_schema.table_partition_cols(); // Just the partition columnsAggregateUDFImpl::is_ordered_set_aggregate 已重命名为 AggregateUDFImpl::supports_within_group_clause
此方法已重命名,以便更准确地反映其对聚合 UDF 实现的实际影响。相应的 AggregateUDF::is_ordered_set_aggregate 也已重命名为 AggregateUDF::supports_within_group_clause。该方法的功能没有任何变化,仍然仅用于表示是否允许该聚合函数使用 WITHIN GROUP SQL 语法。
评论
登录后参与评论
KnowForge