升级指南

DataFusion 51.0.0

师成师成· 更新于 2026-09-28· 阅读 20 分钟· 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 columns

AggregateUDFImpl::is_ordered_set_aggregate 已重命名为 AggregateUDFImpl::supports_within_group_clause

此方法已重命名,以便更准确地反映其对聚合 UDF 实现的实际影响。相应的 AggregateUDF::is_ordered_set_aggregate 也已重命名为 AggregateUDF::supports_within_group_clause。该方法的功能没有任何变化,仍然仅用于表示是否允许该聚合函数使用 WITHIN GROUP SQL 语法。

评论

登录后参与评论

正在加载评论…