ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache DataFusion 46.0.0 升级指南:四大破坏性变更的迁移实战

Apache DataFusion 46.0.0 升级指南:四大破坏性变更的迁移实战 大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载本篇指南聚焦 Apache DataFusion 46.0.0 引入的四个核心破坏性变更ScalarUDF 调用 API 统一为invoke_with_args()、内置 DataSource 体系重构为DataSourceExecFileSource、datafusion-cli字符串转义行为修正以及数组函数签名声明方式的全面改造。读完本文你将掌握上述每项变更的报错识别、迁移改写范式与底层源码原理可直接用于升级现有基于 DataFusion 的查询引擎或自定义 UDF 库。一、ScalarUDF 调用 API全面迁移到invoke_with_argsDataFusion 46.0.0 对标量用户自定义函数ScalarUDF的调用方式进行了统一所有函数实现改为实现ScalarUDFImpl::invoke_with_args()同时弃用了invoke()、invoke_batch()与invoke_no_args()三个旧接口。1.1 报错识别如果你在升级后运行查询时看到如下错误说明仍有旧 API 在起作用例如concat这类内置函数通过兼容层被调用时This feature is not implemented: Function concat does not implement invoke but called从源码看该错误源于执行路径中只能调用invoke_with_args这一统一入口见 datafusion/expr/src/udf.rs旧接口不再被运行时解析。1.2 迁移示例从invoke_batch到invoke_with_args假设原有自定义函数SparkConcat通过invoke_batch分别派发到不同的内部实现impl ScalarUDFImpl for SparkConcat { fn invoke_batch(self, args: [ColumnarValue], number_rows: usize) - ResultColumnarValue { 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) } } }迁移后签名改为接收统一的ScalarFunctionArgs结构体且内部代理调用也同步改为invoke_with_argsimpl ScalarUDFImpl for SparkConcat { fn invoke_with_args(self, args: ScalarFunctionArgs) - ResultColumnarValue { 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) } } }1.3 源码层面的关键信息唯一必选接口在ScalarUDFImpltrait 中invoke_with_args是唯一没有默认实现的调用方法datafusion/expr/src/udf.rs#L702签名固定为fn invoke_with_args(self, args: ScalarFunctionArgs) - ResultColumnarValue。ScalarFunctionArgs结构datafusion/expr/src/udf.rs#L414-L427统一携带了args: VecColumnarValue已求值的参数、arg_fields各参数对应 Field、number_rows当前 RecordBatch 行数、return_field规划期承诺的返回字段以及config_options执行期配置。这意味着函数实现不再需要分别从多个参数中拼凑信息。返回类型校验ScalarUDF::invoke_with_args在 debug 构建下会自动断言实际返回值类型与规划期承诺类型一致datafusion/expr/src/udf.rs#L267-L287因此自定义函数必须保证return_type与实际计算出的列类型严格相符。性能提示trait 文档明确指出当参数中存在常量值即ColumnarValue::Scalar时应优先处理以获得最佳性能若实现简单但较慢的路径可使用ColumnarValue::values_to_arrays将参数统一转为数组。二、DataSource 体系重构ParquetExec等内置 Executor 弃用DataFusion 46.0.0 对内置数据源的组织方式做了重大调整不再为每种文件格式维护独立的ExecutionPlan如ParquetExec、AvroExec、CsvExec、JsonExec而是统一收敛到DataSourceExec格式相关信息下沉到新的DataSource与FileSourcetrait 中。2.1 新的分层架构源码中的文档注释给出了清晰的层级关系datafusion/datasource/src/source.rs#L80-L127DataSourceExec唯一的执行计划入口负责读取数据。DataSourcetraitdatafusion/datasource/src/source.rs#L128抽象数据来源提供open()、repartitioned()、output_partitioning()、partition_statistics()、try_pushdown_filters()等能力有两个内建实现FileScanConfig文件列表与MemorySourceConfig内存中的 RecordBatch。FileSourcetrait承载文件格式相关的具体行为其具体实现包括ArrowSource、ParquetSource、JsonSource等。2.2 Cookbook如何识别执行计划中的扫描节点旧代码通过downcast_ref::ParquetExec()直接检查并提取信息if let Some(parquet_exec) plan.as_any().downcast_ref::ParquetExec() { // Do something with ParquetExec here }新代码需要先下转为DataSourceExec再经由data_source()下转到FileScanConfig最后下转格式对应的ParquetSourceif 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 } } }信息分布原则文件组FileGroups、投影、统计等信息位于FileScanConfig剪枝谓词PruningPredicates与 Parquet 专属选项位于ParquetSource。2.3 CookbookParquetExecBuilder构建方式迁移旧代码用ParquetExecBuilder::new(...)直接构建执行计划再通过with_predicate追加谓词下推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); } };新代码先构造ParquetSource承载格式选项与谓词下推再将其作为参数交给FileScanConfig最后调用build()生成执行计划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(),FileScanConfig::build()正是新体系中生成DataSourceExec的标准入口datafusion/datasource/src/file_scan_config/mod.rs#L529它会根据内部携带的DataSource此处即ParquetSource构造出可直接加入执行树的物理计划节点。这一重构带来的收益是文件格式差异被隔离在FileSource实现内部DataSourceExec对谓词下推、统计信息、分区裁剪等通用逻辑可以跨格式复用第三方格式如 delta-rs 的 Delta Lake 表接入成本也显著降低。三、datafusion-cli不再自动反转义字符串此前的datafusion-cli会错误地对字符串字面量执行反转义46.0.0 起移除了该行为SQL 语义与标准一致。3.1 单引号转义SQL 标准中字符串内的单引号用两个连续单引号表示 select its escaped; ---------------------- | Utf8(its escaped) | ---------------------- | its escaped | ---------------------- 1 row(s) fetched.3.2 特殊字符使用 E 字符串如需在字面量中包含换行符等特殊字符如\n不再依赖 CLI 自动处理而应使用E前缀字符串 select foo\nbar; ------------------ | Utf8(foo\nbar) | ------------------ | foo\nbar | ------------------ 1 row(s) fetched. Elapsed 0.005 seconds.注意对比输出结果不加E时\n被当作普通字符原样保留这正是不再自动反转义的直观体现。迁移要点很简单——检查所有依赖旧 CLI 反转义行为的脚本将\n、\t等转义序列改为E...形式将改为。四、数组标量函数签名从枚举选择到伪类型向量DataFusion 46.0.0 重构了数组类标量函数的签名声明方式。过去函数需要从预定义的ArrayFunctionSignature枚举变体中选择一个组合签名现在签名由一组伪类型pseudo-typeVecArrayFunctionArgument描述每个伪类型对应一个实参。4.1 伪类型语义ArrayFunctionArgument枚举的变体含义见 datafusion/expr-common/src/signature.rs#L568-L580Array参数类型为 List/LargeList/FixedSizeList所有 Array 参数必须可强制转换为同一类型。Element参数可被强制转换为 Array 参数的内层元素类型。IndexInt64类型的索引参数。新格式的载体是ArrayFunctionSignature::Array结构体变体包含两个字段arguments: VecArrayFunctionArgument与array_coercion: OptionListCoerciondatafusion/expr-common/src/signature.rs#L530-L537后者用于声明数组参数的强制转换策略如ListCoercion::FixedSizedListToList。4.2 旧变体到新格式的完整对照ArrayAndElement数组 元素use datafusion::common::utils::ListCoercion; use datafusion_expr_common::signature::{ArrayFunctionArgument, ArrayFunctionSignature, TypeSignature}; TypeSignature::ArraySignature(ArrayFunctionSignature::Array { arguments: vec![ArrayFunctionArgument::Array, ArrayFunctionArgument::Element], array_coercion: Some(ListCoercion::FixedSizedListToList), });ElementAndArray元素 数组TypeSignature::ArraySignature(ArrayFunctionSignature::Array { arguments: vec![ArrayFunctionArgument::Element, ArrayFunctionArgument::Array], array_coercion: Some(ListCoercion::FixedSizedListToList), });ArrayAndIndex数组 索引TypeSignature::ArraySignature(ArrayFunctionSignature::Array { arguments: vec![ArrayFunctionArgument::Array, ArrayFunctionArgument::Index], array_coercion: None, });ArrayAndElementAndOptionalIndex数组 元素 可选索引旧枚举中的可选索引在新格式下无法用一个变体表达需拆分为两个签名的TypeSignature::OneOf组合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, }), ]);Array仅数组TypeSignature::ArraySignature(ArrayFunctionSignature::Array { arguments: vec![ArrayFunctionArgument::Array], array_coercion: None, });4.3 更便捷的替代签名辅助函数如果不想手写ArrayFunctionSignature::Array { ... }结构可以直接使用Signature上预置的构造函数它们会替你完成TypeSignature的组装Signature::array_and_elementSignature::array_and_element_and_optional_indexSignature::array_and_indexSignature::array建议在新增数组函数时优先使用这些辅助函数减少手写结构的出错概率仅在需要自定义array_coercion时手写结构体形式。五、升级清单速查变更点旧写法新写法报错/现象特征UDF 调用invoke()/invoke_batch()/invoke_no_args()invoke_with_args(ScalarFunctionArgs)Function concat does not implement invoke but called扫描计划检查downcast_ref::ParquetExec()downcast_ref::DataSourceExec()→FileScanConfig→ParquetSource编译期类型不匹配扫描计划构建ParquetExecBuilder::new(FileScanConfig)ParquetSource::new(opts)FileScanConfig::new(..., Arc::new(source)).build()编译期类型不匹配CLI 字符串依赖自动反转义转义单引号、E...表达特殊字符字符串输出与预期不一致数组函数签名枚举组合变体VecArrayFunctionArgumentListCoercion编译期签名结构不匹配六、结语DataFusion 46.0.0 的这四项变更体现了两个长期方向接口收敛UDF 调用与文件扫描入口都收敛为单一统一 API与语义规范化CLI 字符串处理回归 SQL 标准。对于库作者而言invoke_with_args迁移大多只需机械改写而DataSource/FileSource重构则需要重新组织扫描计划的构造与检查逻辑对于 CLI 脚本用户重点排查字符串转义即可。完成上述迁移后你的代码将与 DataFusion 46 及后续版本保持长期兼容。若需要更细粒度的实现参考可继续阅读 ScalarFunctionArgs 定义、DataSource trait 与 ArrayFunctionArgument 定义。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐猫抓Cat-Catch浏览器视频资源嗅探的终极解决方案猫抓Cat Catch浏览器视频资源嗅探的终极解决方案 在浏览网页时你是否经常遇到想要保存的视频、音频或图片资源却苦于网站不提供下载功能猫抓Cat Ca大数据数据分析后端react-jsonschema-form v3.x 升级指南四大破坏性变更解析与迁移实践react jsonschema form v3.x 升级指南四大破坏性变更解析与迁移实践 本指南以当前仓库中保留的官方迁移文档v3.x upgrade g前端UI组件DataFusion 55.0.0 升级指南破坏性变更全景解析与迁移路径DataFusion 55.0.0 升级指南破坏性变更全景解析与迁移路径 DataFusion 55 是一次涉及面很广的破坏性发布从 SQL 层的谓词求值顺大数据数据分析后端上一篇FOSSASIA Labyrinth 项目常见问题解决方案下一篇【亲测免费】 Live2D Widget.js 常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表