升级指南

DataFusion 52.0.0

师成师成· 更新于 2026-09-28· 阅读 40 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

升级指南

DataFusion 52.0.0

DFSchema API 的变更

为了使规划更加高效,DFSchema 上的若干方法已改为返回底层的 [&FieldRef] 而非 [&Field] 的引用。这样规划器便可以通过 Arc::clone 以更低的成本复制这些引用,而不必克隆整个 Field 结构体。

你可能需要修改代码,使用 Arc::clone 而不是直接对 Field 调用 .as_ref().clone()。例如:

- let field = df_schema.field("my_column").as_ref().clone();
+ let field = Arc::clone(df_schema.field("my_column"));

ListingTableProvider 现在会缓存 LIST 命令

在之前的版本中,ListingTableProvider 每次需要为查询列出文件时,都会向底层对象存储发出 LIST 命令。为了提升性能,ListingTableProvider 现在会在 ListingTableProvider 实例的整个生命周期内缓存 LIST 命令的结果,直到缓存条目过期为止。

请注意,默认情况下缓存没有过期时间,因此如果底层对象存储中的文件被添加或删除,ListingTableProvider 将无法看到这些更改,直到该 ListingTableProvider 实例被销毁并重新创建。

你可以通过配置选项来设置缓存的最大大小和缓存条目的过期时间:

  • datafusion.runtime.list_files_cache_limit - 以字节为单位限制缓存的大小
  • datafusion.runtime.list_files_cache_ttl - 以分钟和/或秒为单位限制条目的 TTL(生存时间)

有关详细的配置信息,请参阅 DataFusion 运行时配置用户指南。

将该限制设置为 0 即可禁用缓存:

SET datafusion.runtime.list_files_cache_limit TO "0K";

请注意,内部 API 已改为使用 trait ListFilesCache,而不再使用类型别名。

newlines_in_values 从 FileScanConfig 移至 CsvOptions

CSV 专用的 newlines_in_values 配置选项已从 FileScanConfig 移至 CsvOptions,因为它只适用于 CSV 文件解析。

受影响的用户:

  • 通过 FileScanConfigBuilder::with_newlines_in_values() 设置 newlines_in_values 的用户

迁移指南:

在 CsvOptions 中设置 newlines_in_values,而不是在 FileScanConfigBuilder 上设置:

迁移前:

let source = Arc::new(CsvSource::new(file_schema.clone()));
let config = FileScanConfigBuilder::new(object_store_url, source)
    .with_newlines_in_values(true)
    .build();

之后:

let options = CsvOptions {
    newlines_in_values: Some(true),
    ..Default::default()
};
let source = Arc::new(CsvSource::new(file_schema.clone())
    .with_csv_options(options));
let config = FileScanConfigBuilder::new(object_store_url, source)
    .build();

移除 pyarrow 特性

pyarrow 特性标志已被移除。该功能自 44.0.0 版本起已迁移至 datafusion-python 仓库。

重构 FileSource 构造函数与 FileScanConfigBuilder 以提前传入 schema

向文件源和扫描配置传递 schema 的方式已进行大幅重构。文件源现在要求在构造时就提供 schema(包含分区列),并且 FileScanConfigBuilder 不再接受单独的 schema 参数。

受影响的用户:

  • 直接创建 FileScanConfig 或文件源(ParquetSource、CsvSource、JsonSource、AvroSource)的用户
  • 实现自定义 FileFormat 的用户

关键变更:

  1. FileSource 构造函数现在需要 TableSchema:所有内置文件源现在都在构造函数中接收 schema:

    - let source = ParquetSource::default();
    + let source = ParquetSource::new(table_schema);
  2. FileScanConfigBuilder 不再将 schema 作为参数接收:schema 现在通过 FileSource 传递:

    - FileScanConfigBuilder::new(url, schema, source)
    + FileScanConfigBuilder::new(url, source)
  3. 分区列现在是 TableSchema 的一部分:with_table_partition_cols() 方法已从 FileScanConfigBuilder 中移除。分区列现在作为 TableSchema 的一部分传递给 FileSource 构造函数:

    + let table_schema = TableSchema::new(
    +     file_schema,
    +     vec![Arc::new(Field::new("date", DataType::Utf8, false))],
    + );
    + let source = ParquetSource::new(table_schema);
      let config = FileScanConfigBuilder::new(url, source)
    -     .with_table_partition_cols(vec![Field::new("date", DataType::Utf8, false)])
          .with_file(partitioned_file)
          .build();
  4. FileFormat::file_source() 现在接收 TableSchema 参数:自定义的 FileFormat 实现必须相应更新:

    impl FileFormat for MyFileFormat {
    -   fn file_source(&self) -> Arc<dyn FileSource> {
    +   fn file_source(&self, table_schema: TableSchema) -> Arc<dyn FileSource> {
    -       Arc::new(MyFileSource::default())
    +       Arc::new(MyFileSource::new(table_schema))
        }
    }

迁移示例:

针对 Parquet 文件:

- let source = Arc::new(ParquetSource::default());
- let config = FileScanConfigBuilder::new(url, schema, source)
+ let table_schema = TableSchema::new(schema, vec![]);
+ let source = Arc::new(ParquetSource::new(table_schema));
+ let config = FileScanConfigBuilder::new(url, source)
      .with_file(partitioned_file)
      .build();

对于包含分区列的 CSV 文件:

- let source = Arc::new(CsvSource::new(true, b',', b'"'));
- let config = FileScanConfigBuilder::new(url, file_schema, source)
-     .with_table_partition_cols(vec![Field::new("year", DataType::Int32, false)])
+ let options = CsvOptions {
+     has_header: Some(true),
+     delimiter: b',',
+     quote: b'"',
+     ..Default::default()
+ };
+ let table_schema = TableSchema::new(
+     file_schema,
+     vec![Arc::new(Field::new("year", DataType::Int32, false))],
+ );
+ let source = Arc::new(CsvSource::new(table_schema).with_csv_options(options));
+ let config = FileScanConfigBuilder::new(url, source)
      .build();

Parquet 过滤下推中的自适应过滤器表示

从 Arrow 57.1.0 开始,DataFusion 在为 Parquet 文件求值下推的过滤器时采用了一种新的自适应过滤器策略。对于某些类型的查询,过滤结果用位掩码(bitmask)表示比用选择集(selection)表示更加高效,这种新策略可以提升这类查询的性能。详见 arrow-rs devlive-community/knowforge#5523。

此更改仅适用于启用了过滤下推功能的内置 Parquet 数据源(该功能目前尚未成为默认行为)。

你可以通过将 datafusion.execution.parquet.force_filter_selections 配置项设置为 true 来禁用这一新行为。

> set datafusion.execution.parquet.force_filter_selections = true;

统计信息处理从 FileSource 移至 FileScanConfig

现在统计信息由 FileScanConfig 直接管理,不再委托给 FileSource 实现。这简化了 FileSource trait,并使所有文件格式的统计信息处理更加一致。

受影响的用户:

  • 已实现自定义 FileSource 的用户

破坏性变更:

FileSource trait 中移除了两个方法:

  • with_statistics(&self, statistics: Statistics) -> Arc<dyn FileSource>
  • statistics(&self) -> Result<Statistics>

迁移指南:

如果你有自定义的 FileSource 实现,需要:

  1. 移除 with_statistics 方法的实现
  2. 移除 statistics 方法的实现
  3. 移除所有用于存储统计信息的内部状态

修改前:

#[derive(Clone)]
struct MyCustomSource {
    table_schema: TableSchema,
    projected_statistics: Option<Statistics>,
    // other fields...
}

impl FileSource for MyCustomSource {
    fn with_statistics(&self, statistics: Statistics) -> Arc<dyn FileSource> {
        Arc::new(Self {
            table_schema: self.table_schema.clone(),
            projected_statistics: Some(statistics),
            // other fields...
        })
    }

    fn statistics(&self) -> Result<Statistics> {
        Ok(self.projected_statistics.clone().unwrap_or_else(||
            Statistics::new_unknown(self.table_schema.file_schema())
        ))
    }

    // other methods...
}

之后:

#[derive(Clone)]
struct MyCustomSource {
    table_schema: TableSchema,
    // projected_statistics field removed
    // other fields...
}

impl FileSource for MyCustomSource {
    // with_statistics method removed
    // statistics method removed

    // other methods...
}

访问统计信息:

现在通过 FileScanConfig 而非 FileSource 来访问统计信息:

- let stats = config.file_source.statistics()?;
+ let stats = config.statistics();

注意,FileScanConfig::statistics() 会在存在过滤条件时自动将统计信息标记为不精确,从而在过滤条件下推时保证正确性。

分区列处理已移出 PhysicalExprAdapter

分区列替换现在是一个独立的预处理步骤,在通过 PhysicalExprAdapter 进行表达式重写之前执行。这一变更带来了更好的职责划分,使适配器更专注于处理 schema 差异,而非分区值替换。

受影响的用户:

  • 自定义实现了 PhysicalExprAdapterFactory 并在其中处理分区列的用户
  • 直接使用 FilePruner API 的用户

破坏性变更:

  1. FilePruner::try_new() 的签名已变更:由于分区列处理现在单独进行,partition_fields 参数已被移除
  2. 分区列替换现在必须在表达式传入适配器之前,通过 replace_columns_with_literals() 完成

迁移指南:

如果你的代码中存在使用分区字段创建 FilePruner 的逻辑:

变更前:

use datafusion_pruning::FilePruner;

let pruner = FilePruner::try_new(
    predicate,
    file_schema,
    partition_fields,  // This parameter is removed
    file_stats,
)?;

之后:

use datafusion_pruning::FilePruner;

// Partition fields are no longer needed
let pruner = FilePruner::try_new(
    predicate,
    file_schema,
    file_stats,
)?;

如果你的自定义代码依赖 PhysicalExprAdapter 来处理分区列,现在则需要单独调用 replace_columns_with_literals():

修改前:

// Adapter handled partition column replacement internally
let adapted_expr = adapter.rewrite(expr)?;

之后:

use datafusion_physical_expr_adapter::replace_columns_with_literals;

// Replace partition columns first
let expr_with_literals = replace_columns_with_literals(expr, &partition_values)?;
// Then apply the adapter
let adapted_expr = adapter.rewrite(expr_with_literals)?;

build_row_filter 签名简化

datafusion-datasource-parquet 中的 build_row_filter 函数已简化,现在只接受单个 schema 参数,而非两个。现在期望在将过滤器传入该函数之前,它已经被适配到物理文件 schema(parquet 文件 schema 的 arrow 表示),例如可以使用 PhysicalExprAdapter 完成适配。

受影响的用户:

  • 直接调用 build_row_filter 的用户

破坏性变更:

函数签名由:

pub fn build_row_filter(
    expr: &Arc<dyn PhysicalExpr>,
    physical_file_schema: &SchemaRef,
    predicate_file_schema: &SchemaRef,  // removed
    metadata: &ParquetMetaData,
    reorder_predicates: bool,
    file_metrics: &ParquetFileMetrics,
) -> Result<Option<RowFilter>>

收件人:

pub fn build_row_filter(
    expr: &Arc<dyn PhysicalExpr>,
    file_schema: &SchemaRef,
    metadata: &ParquetMetaData,
    reorder_predicates: bool,
    file_metrics: &ParquetFileMetrics,
) -> Result<Option<RowFilter>>

迁移指南:

从你的调用中移除重复的 schema 参数:

- build_row_filter(&predicate, &file_schema, &file_schema, metadata, reorder, metrics)
+ build_row_filter(&predicate, &file_schema, metadata, reorder, metrics)

规划器现在要求显式选择启用 WITHIN GROUP 语法

SQL 规划器现在更严格地强制执行聚合 UDF 的约定:只有当聚合 UDAF 通过让 AggregateUDFImpl::supports_within_group_clause() 返回 true 来明确声明支持时,才会接受 WITHIN GROUP (ORDER BY ...) 语法。

此前,即使顺序敏感的聚合并未实现有序集语义,规划器也会将 WITHIN GROUP 子句转发给它们,这可能导致 SUM(x) WITHIN GROUP (ORDER BY x) 之类的查询被成功规划。这种行为过于宽松,现已改为与 PostgreSQL 及文档中所述语义保持一致。

迁移说明:如果你的 UDAF 有意实现有序集语义并希望接受 WITHIN GROUP SQL 语法,请更新实现,让 supports_within_group_clause() 返回 true,并在累加器实现中处理排序语义。如果你的 UDAF 只是顺序敏感(而非有序集聚合),则不要声明支持 supports_within_group_clause(),客户端应改用其他函数签名(例如,将显式排序作为函数参数传入)。

AggregateUDFImpl::supports_null_handling_clause 现在默认为 false

该方法指定聚合函数在 SQL 解析期间是否允许 IGNORE NULLS/RESPECT NULLS,这意味着它在计算过程中会遵循这些配置。

在之前的版本中,大多数 DataFusion 聚合函数并不使用这一语法,而该语法默认是被允许的,因此它们会静默忽略它。我们现在改为:只有少数确实遵循该子句的函数(例如 array_agg、first_value、last_value)才需要实现它。

用户自定义的聚合函数如果在使用了该语法时未通过重写该方法显式声明支持,同样会报错。

例如,此类查询的 SQL 解析现在将会失败:

SELECT median(c1) IGNORE NULLS FROM table

而是静默成功。

CacheAccessor trait 的 API 变更

remove API 不再要求可变实例

FFI crate 更新

datafusion-ffi crate 中的许多结构体已更新,以便更轻松地转换为它们所代表的底层 trait 类型。这简化了部分代码路径,同时也为库代码通过外部函数接口(FFI)往返调用的场景带来了额外的改进。

要更新你的代码,假设你有一个名为 ffi_provider 的 FFI_SchemaProvider,并且希望将其作为 SchemaProvider 使用。在旧的做法中,你可能会这样写:

let foreign_provider: ForeignSchemaProvider = ffi_provider.into();
    let foreign_provider = Arc::new(foreign_provider) as Arc<dyn SchemaProvider>;

这段代码现在应写成:

let foreign_provider: Arc<dyn SchemaProvider + Send> = ffi_provider.into();
    let foreign_provider = foreign_provider as Arc<dyn SchemaProvider>;

对于用户自定义函数,更新方式类似,但你可能需要修改创建 ScalarUDF 的调用方式。聚合函数和窗口函数遵循相同的模式。

之前你可能会这样写:

let foreign_udf: ForeignScalarUDF = ffi_udf.try_into()?;
    let foreign_udf: ScalarUDF = foreign_udf.into();

现在这应该改为:

let foreign_udf: Arc<dyn ScalarUDFImpl> = ffi_udf.into();
    let foreign_udf = ScalarUDF::new_from_shared_impl(foreign_udf);

在创建以下任一结构体时,现在要求用户提供 TaskContextProvider,并可选择性地提供 LogicalExtensionCodec:

  • FFI_CatalogListProvider
  • FFI_CatalogProvider
  • FFI_SchemaProvider
  • FFI_TableProvider
  • FFI_TableFunction

这些结构体中的每一个都提供了 new() 和 new_with_ffi_codec() 两个方法用于实例化。例如,以前你会这样写:

let table = Arc::new(MyTableProvider::new());
   let ffi_table = FFI_TableProvider::new(table, None);

现在你需要提供一个 TaskContextProvider。该 trait 最常见的实现是 SessionContext。

let ctx = Arc::new(SessionContext::default());
   let table = Arc::new(MyTableProvider::new());
   let ffi_table = FFI_TableProvider::new(table, None, ctx, None);

如果要执行大量此类操作,使用另一种函数来创建这些结构可能会更方便。FFI_LogicalExtensionCodec 同时也会保存 TaskContextProvider。

let codec = Arc::new(DefaultLogicalExtensionCodec {});
   let ctx = Arc::new(SessionContext::default());
   let ffi_codec = FFI_LogicalExtensionCodec::new(codec, None, ctx);
   let table = Arc::new(MyTableProvider::new());
   let ffi_table = FFI_TableProvider::new_with_ffi_codec(table, None, ffi_codec);

关于 TaskContextProvider 的更多用法,可以参阅该 crate 的 README。

此外,标量 UDF 的 FFI 结构不再包含 return_type 调用。由于 ForeignScalarUDF 结构体实现了 return_field_from_args,这段代码早已不再使用。

投影处理从 FileScanConfig 移至 FileSource

投影处理已从 FileScanConfig 移入 FileSource 实现中。这样可以实现针对特定格式的投影下推(例如,Parquet 可以下推结构体字段访问,Vortex 可以把计算表达式下推到尚未解码的数据中)。

受影响的用户:

  • 已实现自定义 FileSource 的用户
  • 直接使用 FileScanConfigBuilder::with_projection_indices 的用户

破坏性变更:

  1. FileSource::with_projection 被 try_pushdown_projection 取代:

    with_projection(&self, config: &FileScanConfig) -> Arc<dyn FileSource> 方法已被移除,取而代之的是 try_pushdown_projection(&self, projection: &ProjectionExprs) -> Result<Option<Arc<dyn FileSource>>>。

  2. FileScanConfig.projection_exprs 字段被移除:

    投影现在直接存储在 FileSource 中,而不再存储在 FileScanConfig 中。FileScanConfig 上用于访问投影信息的各种公开辅助方法也已被移除。

  3. FileScanConfigBuilder::with_projection_indices 现在返回 Result<Self>:

    如果投影下推失败,该方法现在可能会失败。

  4. FileSource::create_file_opener 现在返回 Result<Arc<dyn FileOpener>>:

    此前直接返回 Arc<dyn FileOpener>。任何可能在创建 FileOpener 时失败的 FileSource 实现,现在都应返回相应的错误。

  5. DataSource::try_swapping_with_projection 签名变更:

    参数类型从 &[ProjectionExpr] 变为 &ProjectionExprs。

迁移指南:

如果你有自定义的 FileSource 实现:

变更前:

impl FileSource for MyCustomSource {
    fn with_projection(&self, config: &FileScanConfig) -> Arc<dyn FileSource> {
        // Apply projection from config
        Arc::new(Self { /* ... */ })
    }

    fn create_file_opener(
        &self,
        object_store: Arc<dyn ObjectStore>,
        base_config: &FileScanConfig,
        partition: usize,
    ) -> Arc<dyn FileOpener> {
        Arc::new(MyOpener { /* ... */ })
    }
}

之后:

impl FileSource for MyCustomSource {
    fn try_pushdown_projection(
        &self,
        projection: &ProjectionExprs,
    ) -> Result<Option<Arc<dyn FileSource>>> {
        // Return None if projection cannot be pushed down
        // Return Some(new_source) with projection applied if it can
        Ok(Some(Arc::new(Self {
            projection: Some(projection.clone()),
            /* ... */
        })))
    }

    fn projection(&self) -> Option<&ProjectionExprs> {
        self.projection.as_ref()
    }

    fn create_file_opener(
        &self,
        object_store: Arc<dyn ObjectStore>,
        base_config: &FileScanConfig,
        partition: usize,
    ) -> Result<Arc<dyn FileOpener>> {
        Ok(Arc::new(MyOpener { /* ... */ }))
    }
}

建议你查看 devlive-community/knowforge#18627,该 PR 引入了这些变更,其中包含更多关于各种内置文件源如何处理这一问题的示例。

我们新增了 SplitProjection 和 ProjectionOpener 两个辅助工具,以便更轻松地在你的 FileSource 实现中处理投影。

对于只能处理简单列选择(而非计算表达式)的文件源,请使用 SplitProjection 和 ProjectionOpener 辅助工具,将投影拆分为可下推和不可下推两部分:

use datafusion_datasource::projection::{SplitProjection, ProjectionOpener};

// In try_pushdown_projection:
let split = SplitProjection::new(projection, self.table_schema())?;
// Use split.file_projection() for what to push down to the file format
// The ProjectionOpener wrapper will handle the rest

对于 FileScanConfigBuilder 用户:

let config = FileScanConfigBuilder::new(url, source)
-   .with_projection_indices(Some(vec![0, 2, 3]))
+   .with_projection_indices(Some(vec![0, 2, 3]))?
    .build();

SchemaAdapter 与 SchemaAdapterFactory 已被完全移除

继 DataFusion 49.0.0 中宣布弃用之后,SchemaAdapterFactory 已从 Parquet 扫描中被完全移除。这适用于以下两者:

以下符号已被弃用,并将在下一个版本中移除:

  • SchemaAdapter trait
  • SchemaAdapterFactory trait
  • SchemaMapper trait
  • SchemaMapping 结构体
  • DefaultSchemaAdapterFactory 结构体

这些类型过去用于在读取文件时调整记录批次(record batch)的模式(schema)。该功能现已被 PhysicalExprAdapterFactory 取代,后者在规划阶段重写表达式,而不是在运行时转换批次。如果你此前使用自定义的 SchemaAdapterFactory 来做模式适配(例如默认列值、类型强制转换),现在应改为实现 PhysicalExprAdapterFactory。有关如何实现自定义 PhysicalExprAdapterFactory,请参阅默认列值示例。

迁移指南:

如果你实现了自定义的 SchemaAdapterFactory,请迁移到 PhysicalExprAdapterFactory。完整的实现请参阅默认列值示例。

评论

登录后参与评论

正在加载评论…