升级指南

DataFusion 46.0.0

qianmoQqianmoQ· 更新于 2026-09-29· 阅读 17 分钟· 0 次阅读

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

升级指南

DataFusion 46.0.0

使用 invoke_with_args 代替 invoke() 和 invoke_batch()

DataFusion 正在转向统一的 ScalarUDF 调用 API,即 ScalarUDFImpl::invoke_with_args(),并弃用了 ScalarUDFImpl::invoke()、ScalarUDFImpl::invoke_batch() 以及 ScalarUDFImpl::invoke_no_args()。

如果你看到如下所示的错误,说明代码中仍在使用旧版 API:

This feature is not implemented: Function concat does not implement invoke but called

要修复此错误,请改用 ScalarUDFImpl::invoke_with_args(),如下所示。参见 PR 14876 中的示例。

假设已有如下代码:

impl ScalarUDFImpl for SparkConcat {
...
    fn invoke_batch(&self, args: &[ColumnarValue], number_rows: usize) -> Result<ColumnarValue> {
        if args
            .iter()
            .any(|arg| matches!(arg.data_type(), DataType::List(_)))
        {
            ArrayConcat::new().invoke_batch(args, number_rows)
        } else {
            ConcatFunc::new().invoke_batch(args, number_rows)
        }
    }
}

到

impl ScalarUDFImpl for SparkConcat {
    ...
    fn invoke_with_args(&self, args: ScalarFunctionArgs) -> Result<ColumnarValue> {
        if args
            .args
            .iter()
            .any(|arg| matches!(arg.data_type(), DataType::List(_)))
        {
            ArrayConcat::new().invoke_with_args(args)
        } else {
            ConcatFunc::new().invoke_with_args(args)
        }
    }
}

ParquetExec、AvroExec、CsvExec、JsonExec 已弃用

DataFusion 46 对内置 DataSource 的组织方式做了重大变更。不同文件格式不再各自拥有独立的 ExecutionPlan,而是统一使用 DataSourceExec,与格式相关的信息由新的 trait DataSource 和 FileSource 承载。

更多信息请参见:

实战指南:ParquetExecBuilder 的变更

类似下面这样查找 ParquetExec 的代码将不再有效:

    if let Some(parquet_exec) = plan.as_any().downcast_ref::<ParquetExec>() {
        // Do something with ParquetExec here
    }

相反,使用 DataSourceExec 后,同样的信息现在位于 FileScanConfig 和 ParquetSource 上。等效代码如下:

if let Some(datasource_exec) = plan.as_any().downcast_ref::<DataSourceExec>() {
  if let Some(scan_config) = datasource_exec.data_source().as_any().downcast_ref::<FileScanConfig>() {
    // FileGroups, and other information is on the FileScanConfig
    // parquet
    if let Some(parquet_source) = scan_config.file_source.as_any().downcast_ref::<ParquetSource>()
    {
      // Information on PruningPredicates and parquet options are here
    }
}

实操指南:ParquetExecBuilder 的变更

同理,使用 ParquetExecBuilder 构建 ParquetExec 的代码(如下所示)也必须进行修改:

let mut exec_plan_builder = ParquetExecBuilder::new(
    FileScanConfig::new(self.log_store.object_store_url(), file_schema)
        .with_projection(self.projection.cloned())
        .with_limit(self.limit)
        .with_table_partition_cols(table_partition_cols),
)
.with_schema_adapter_factory(Arc::new(DeltaSchemaAdapterFactory {}))
.with_table_parquet_options(parquet_options);

// Add filter
if let Some(predicate) = logical_filter {
    if config.enable_parquet_pushdown {
        exec_plan_builder = exec_plan_builder.with_predicate(predicate);
    }
};

新代码应使用 FileScanConfig 来构建相应的 DataSourceExec:

let mut file_source = ParquetSource::new(parquet_options)
    .with_schema_adapter_factory(Arc::new(DeltaSchemaAdapterFactory {}));

// Add filter
if let Some(predicate) = logical_filter {
    if config.enable_parquet_pushdown {
        file_source = file_source.with_predicate(predicate);
    }
};

let file_scan_config = FileScanConfig::new(
    self.log_store.object_store_url(),
    file_schema,
    Arc::new(file_source),
)
.with_statistics(stats)
.with_projection(self.projection.cloned())
.with_limit(self.limit)
.with_table_partition_cols(table_partition_cols);

// Build the actual scan like this
parquet_scan: file_scan_config.build(),

datafusion-cli 不再自动对字符串进行反转义

此前 datafusion-cli 会错误地对字符串字面量进行反转义(详见工单)。

要在 SQL 字面量中转义 ',请使用 '':

> select 'it''s escaped';
+----------------------+
| Utf8("it's escaped") |
+----------------------+
| it's escaped         |
+----------------------+
1 row(s) fetched.

要包含特殊字符(例如用 \n 表示换行),可以使用 E 前缀字符串字面量。例如

> select 'foo\nbar';
+------------------+
| Utf8("foo\nbar") |
+------------------+
| foo\nbar         |
+------------------+
1 row(s) fetched.
Elapsed 0.005 seconds.

数组标量函数签名的变更

DataFusion 46 改变了标量数组函数签名的声明方式。此前,函数需要从 ArrayFunctionSignature 枚举中预定义的签名列表里进行选择。现在,签名可以通过一个 Vec 形式的伪类型列表来定义,其中每个伪类型对应一个参数。这些伪类型是 ArrayFunctionArgument 枚举的变体,具体如下:

  • Array:类型为 List/LargeList/FixedSizeList 的参数。所有 Array 参数必须可强制转换为同一种类型。
  • Element:可强制转换为 Array 参数内部元素类型的参数。
  • Index:一个 Int64 参数。

旧的各个变体可以按如下方式转换为新格式:

TypeSignature::ArraySignature(ArrayFunctionSignature::ArrayAndElement):

TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
    arguments: vec![ArrayFunctionArgument::Array, ArrayFunctionArgument::Element],
    array_coercion: Some(ListCoercion::FixedSizedListToList),
});

TypeSignature::ArraySignature(ArrayFunctionSignature::ElementAndArray):

TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
    arguments: vec![ArrayFunctionArgument::Element, ArrayFunctionArgument::Array],
    array_coercion: Some(ListCoercion::FixedSizedListToList),
});

TypeSignature::ArraySignature(ArrayFunctionSignature::ArrayAndIndex):

TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
    arguments: vec![ArrayFunctionArgument::Array, ArrayFunctionArgument::Index],
    array_coercion: None,
});

TypeSignature::ArraySignature(ArrayFunctionSignature::ArrayAndElementAndOptionalIndex):

TypeSignature::OneOf(vec![
    TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
        arguments: vec![ArrayFunctionArgument::Array, ArrayFunctionArgument::Element],
        array_coercion: None,
    }),
    TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
        arguments: vec![
            ArrayFunctionArgument::Array,
            ArrayFunctionArgument::Element,
            ArrayFunctionArgument::Index,
        ],
        array_coercion: None,
    }),
]);

TypeSignature::ArraySignature(ArrayFunctionSignature::Array):

TypeSignature::ArraySignature(ArrayFunctionSignature::Array {
    arguments: vec![ArrayFunctionArgument::Array],
    array_coercion: None,
});

或者,你也可以改用以下函数之一,它们会为你构造好 TypeSignature:

  • Signature::array_and_element
  • Signature::array_and_element_and_optional_index
  • Signature::array_and_index
  • Signature::array

评论

登录后参与评论

正在加载评论…