DataFusion 46.0.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 承载。
更多信息请参见:
- 设计工单
- 变更 PR PR devlive-community/knowforge#14224
- 升级示例 delta-rs 中的 PR
实战指南: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_elementSignature::array_and_element_and_optional_indexSignature::array_and_indexSignature::array
评论
登录后参与评论
KnowForge