ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache DataFusion 50.0.0 升级指南:从 Hive 分区自动推断到 UDF 特质重构的完整迁移手册

Apache DataFusion 50.0.0 升级指南:从 Hive 分区自动推断到 UDF 特质重构的完整迁移手册 大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载本指南基于 Apache DataFusion 官方库使用指南docs/source/library-user-guide/upgrading/50.0.0.md整理全面覆盖 50.0.0 版本中影响库使用者的全部破坏性变更Breaking Changes包含ListingTable的 Hive 分区自动推断新行为与恢复旧行为的配置方法、MSRV 提升至 1.86.0、三个 UDF 特质对PartialEq/Eq/Hash的强制要求、AsyncScalarUDFImpl签名变更、ProjectionExpr从类型别名到结构体的重构以及ExecutionPlan::reset_state、FileOpenFuture错误类型、FFI UDAF 等底层 API 调整。读完本文你将能够逐项对照完成自有代码库向 DataFusion 50.0.0 的平滑迁移。升级前须知本文涉及的核心变更总览DataFusion 50.0.0 是一次面向库使用者的 API 调整版本变更横跨表工厂行为、Rust 工具链、用户自定义函数UDF/UDAF/UDWF、物理执行计划、数据源流式读取与 FFI 接口等多个层面。下表按影响面归纳本指南将逐一展开的变更项变更类别核心内容影响人群ListingTable 行为CREATE EXTERNAL TABLE自动推断 Hive 分区列使用 Hive 风格目录组织数据的用户Rust 工具链MSRV 提升至 1.86.0所有以 DataFusion 为依赖的项目UDF 特质ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl要求实现PartialEq、Eq、Hash所有自定义函数实现者异步 UDFAsyncScalarUDFImpl::invoke_async_with_args返回ColumnarValue并移除ConfigOptions参数注册远程/异步函数如调用 LLM的用户投影表达式ProjectionExpr由元组别名改为具名字段结构体直接使用ProjectionExec的物理计划开发者会话配置options()等 API 返回ArcConfigOptions访问会话配置的库使用者物理计划新增ExecutionPlan::reset_state自定义ExecutionPlan实现者表达式特质新增PhysicalExpr::is_volatile_node自定义PhysicalExpr实现者数据源FileOpenFuture错误类型改为DataFusionErrorFileOpener自定义实现者schema 重写schema_rewriter模块迁入新 crate依赖物理表达式适配器的用户FFIUDAF 从return_type迁移到return_fieldFFI 库开发者与测试作者ListingTable自动检测 Hive 分区表新行为分区列自动进入表 Schema在 50.0.0 之前通过ListingTableFactory或CREATE EXTERNAL TABLE创建ListingTable时采用 Hive 分区布局例如/table_root/column1value1/column2value2/data.parquet的数据集其分区列column1、column2不会反映在表 schema 中也不会出现在查询结果里。从 50.0.0 起ListingTableFactory会自动推断 Hive 分区将column1、column2作为分区列纳入表的 schema 与数据。仓库源码证实了这一行为在 listing_table_factory.rs 中当CREATE EXTERNAL TABLE未显式提供列定义cmd.schema.fields().is_empty()且未显式指定PARTITIONED BY列时工厂会读取会话配置并调用options.infer_partitions(session_state, first_path)从目录结构推断分区列推断得到的列以Dictionary(UInt16, Utf8)类型写入表 schemalet infer_parts session_state .config_options() .execution .listing_table_factory_infer_partitions; let part_cols if cmd.table_partition_cols.is_empty() infer_parts { options .infer_partitions(session_state, first_path) .await? .into_iter() } else { cmd.table_partition_cols.clone().into_iter() };恢复旧行为的配置开关如果你依赖旧行为分区列不出现在 schema 中可以将配置项datafusion.execution.listing_table_factory_infer_partitions设置为false来恢复。该配置项定义在 datafusion/common/src/config.rs属于datafusion.execution命名空间/// Should a ListingTable created through the ListingTableFactory infer table /// partitions from Hive compliant directories. Defaults to true (partition columns are /// inferred and will be represented in the table schema). pub listing_table_factory_infer_partitions: bool, default true配置方式与 DataFusion 其他执行期配置一致SQL 会话级SET datafusion.execution.listing_table_factory_infer_partitions false;SessionConfig 编程方式构造SessionConfig时通过config.set(...)或直接修改config_options().execution.listing_table_factory_infer_partitions字段。仓库的 sqllogictest 测试对两种行为都做了回归覆盖listing_table_partitions.slt 分别以false和true设置该开关并验证分区列的推断结果information_schema.slt 也记录了该配置的默认值与说明文本。详细讨论见上游 issue #17049。MSRV 更新至 1.86.0DataFusion 50.0.0 将最低支持的 Rust 版本MSRV提升到1.86.0。这意味着所有以 DataFusion 为依赖的项目其工具链版本不得低于 1.86.0否则无法编译。需要说明的是当前仓库主分支的rust-toolchain.toml已指定channel 1.98.1这是后续版本持续演进的结果对于锁定在 50.0.0 的用户请以该版本发布时的 MSRV1.86.0为准并在 CI 中相应调整rust-version或工具链配置。升级理由与讨论详见上游 PR #17230。ScalarUDFImpl、AggregateUDFImpl、WindowUDFImpl特质要求PartialEq、Eq与Hash变更动机消灭手写equals/hash_value的错误隐患此前ScalarUDFImpl::equals、AggregateUDFImpl::equals、WindowUDFImpl::equals以及配套的hash_value方法需要实现者手工编写相等性与哈希逻辑极易出错例如遗漏新加入的字段。50.0.0 移除了这三个特质的equals与hash_value方法改为强制要求实现 Rust 标准的PartialEq、Eq和Hash特质从而让函数相等性比较走上编译器保障的轨道。仓库当前代码印证了这一演进方向在 datafusion/expr/src/udf.rs、datafusion/expr/src/udaf.rs 与 datafusion/expr/src/udwf.rs 中三个特质均约束为Debug DynEq DynHash Send Sync Anydatafusion/expr/src/udf_eq.rs 中的DynEq/DynHash桥接特质通过dyn_eq与dyn_hash将标准特质转发到具体实现使ScalarUDF等封装类型得以在 trait object 层面执行相等性比较与哈希。迁移方法正则替换 vs 手工实现绝大多数标量函数是无状态的且结构体中只有signature字段。这类函数可以借助正则表达式批量迁移搜索模式\#\derive\(Debug\)\?struct \w \{\n *signature\: Signature\,\n *\})替换为#[derive(Debug, PartialEq, Eq, Hash)]$1替换完成后务必人工审查所有改动确认只影响了函数结构体没有误伤其他类型。对于含额外字段或需要自定义相等语义例如忽略运行时统计字段的函数则需手工实现PartialEq/Eq/Hash保证比较逻辑与字段变更保持同步。AsyncScalarUDFImpl::invoke_async_with_args的两次签名调整异步标量函数特质AsyncScalarUDFImpl用于注册远程函数如调用 LLM 的AskLLM示例在 50.0.0 中经历了两次相关调整。先看当前仓库中该特质的最终形态datafusion/expr/src/async_udf.rs#[async_trait] pub trait AsyncScalarUDFImpl: ScalarUDFImpl { ... async fn invoke_async_with_args( self, args: ScalarFunctionArgs, ) - ResultColumnarValue; }调整一返回类型从ArrayRef变为ColumnarValue为使返回值能够走单值scalar value优化并与其他 UDF API 保持一致invoke_async_with_args的返回类型由ArrayRef改为ColumnarValue。迁移只需在旧返回值外层包一层转换# /* comment to avoid running impl AsyncScalarUDFImpl for AskLLM { async fn invoke_async_with_args( self, args: ScalarFunctionArgs, _option: ConfigOptions, ) - ResultColumnarValue { .. return ColumnarValue::from(array_ref); // new codeColumnarValue::from 完成 ArrayRef - ColumnarValue 转换 } } # */调整二移除_option: ConfigOptions参数ConfigOptions已可通过ScalarFunctionArgs参数直接获取因此invoke_async_with_args中原有的_option: ConfigOptions参数被移除接口得到简化。迁移示例# /* comment to avoid running impl AsyncScalarUDFImpl for AskLLM { async fn invoke_async_with_args( self, args: ScalarFunctionArgs, ) - ResultColumnarValue { let options args.config_options; // 从参数中读取配置替代原 _option 参数 .. } ... } # */ScalarFunctionArgs结构体定义于 datafusion/expr/src/udf.rs现在携带args、arg_fields、number_rows、return_field以及config_options: ArcConfigOptions五个字段执行期配置由此通过参数透传无需再从外部单独注入。相关实现细节见上游 issue #16896。ProjectionExpr从类型别名重构为结构体变更前后对比ProjectionExpr投影表达式 别名从元组类型别名改为带具名字段的结构体以提升代码可读性与可维护性变更前类型别名pub type ProjectionExpr (Arcdyn PhysicalExpr, String);变更后结构体#[derive(Debug, Clone)] pub struct ProjectionExpr { pub expr: Arcdyn PhysicalExpr, pub alias: String, }迁移要点构造将元组构造(expr, alias)替换为ProjectionExpr::new(expr, alias)或结构体字面量ProjectionExpr { expr, alias }字段访问将.0、.1替换为.expr、.alias模式匹配将(expr, alias)模式更新为ProjectionExpr { expr, alias }。仓库中大量代码已按新形态使用例如 datafusion/core/src/physical_planner.rs 中的.map(|(expr, alias)| ProjectionExpr { expr, alias })以及 datafusion/core/tests/physical_optimizer/projection_pushdown.rs 中的ProjectionExpr::new(Arc::new(Column::new(b, 2)), b)等测试用法可作迁移参考。此变更主要影响ProjectionExec的使用者实现见上游 PR #17398。SessionState/SessionConfig/OptimizerConfig返回ArcConfigOptions为让ConfigOptions获得更广泛的访问路径并减少不必要的克隆options()等访问器由返回ConfigOptions改为返回ArcConfigOptions。借助Arc同一份ConfigOptions可在线程间共享只有真正需要修改时才克隆整个结构。多数情况下 Rust 编译器会自动解引用Arc因此大部分用户无感但在显式标注类型的场景需要补上.as_ref()# /* comment to avoid running let optimizer_config: ConfigOptions state.options(); // 旧写法 let optimizer_config: ConfigOptions state.options().as_ref(); // 新写法 # */ScalarFunctionArgs::config_options类型为ArcConfigOptions见 datafusion/expr/src/udf.rs即此模式的应用之一。详见上游 PR #16970。Schema Rewriter 模块迁移至新 crateschema_rewriter模块及其符号从datafusion_physical_expr迁出进入新 cratedatafusion_physical_expr_adapter。涉及以下符号DefaultPhysicalExprAdapterDefaultPhysicalExprAdapterFactoryPhysicalExprAdapterPhysicalExprAdapterFactory仓库当前结构印证了迁移结果新 crate 的 datafusion/physical-expr-adapter/src/lib.rs 直接导出这些符号而 schema_rewriter.rs 则承载了 schema 重写的工作职责包括为缺失列填充默认值、按物理 schema 投影重写表达式等。迁移时只需更新 import 路径use datafusion_physical_expr_adapter::{ DefaultPhysicalExprAdapter, DefaultPhysicalExprAdapterFactory, PhysicalExprAdapter, PhysicalExprAdapterFactory };Arrow 与 Parquet 升级到 56.0.0DataFusion 50.0.0 将底层 Apache Arrow 实现升级到56.0.0Parquet 同样为 56.0.0。如果你在代码中直接依赖arrow或parquetcrate 并引用了相关 API请同步升级依赖版本并以 Arrow 56.0.0 的发布说明为准核对变更。新增ExecutionPlan::reset_state方法变更背景DataFusion 49.0.0 曾存在一个 bug动态过滤器dynamic filter目前仅在出现类似ORDER BY ... LIMIT ...的查询时生成在递归查询recursive query如递归 CTE中产生错误结果。50.0.0 通过为ExecutionPlan特质新增reset_state方法修复此问题——凡是需要在执行计划树中维护内部状态、或持有对其他节点引用的ExecutionPlan都应实现该方法以在重新执行时重置状态。默认实现与实现要求reset_state的默认实现datafusion/physical-plan/src/execution_plan.rs只是用现有 children 调用replace_children重建一个同构实例并不重置内部状态因此任何有状态组件例如DynamicFilterPhysicalExpr的执行计划都必须覆写它。需要注意该方法不应递归重置 children因为框架期望在遍历执行计划树时对每个子节点逐个调用。SortExec的示例实现可参考上游 PR #17028本仓库的 sort.rs 中亦有对应实现。嵌套循环连接NLJ重写排序保持的取舍NestedLoopJoin算子被从零重写以提升性能与内存效率。官方文档给出的微基准结论是相比旧实现新实现在极端情况下可获得最高5 倍加速内存占用仅为原来的1%。但需要明确这一变更带来的行为差异新实现无法像旧版本那样保持输入的有序性input sort order。这是性能与内存效率优先于排序保持的根本性设计取舍而非 bug。如果你的查询依赖 NLJ 输出保持输入顺序例如在下游继续依赖该顺序的算子链请评估重写后的计划形态或显式引入排序算子。详见上游 PR #16996。LazyBatchGenerator新增as_any()方法为支持 protobuf 序列化LazyBatchGenerator特质新增了as_any()方法自定义实现需要补充# /* comment to avoid running impl LazyBatchGenerator for MyBatchGenerator { fn as_any(self) - dyn Any { self } ... } # */仓库中GenerateSeries等表函数实现datafusion/functions-table/src/generate_series.rs已按此模式提供as_any返回self可作为参考模板。详见上游 PR #17200。DataSource::try_swapping_with_projection重构DataSource::try_swapping_with_projection被重构以简化方法签名并尽量减少ExecutionPlan与DataSource抽象层之间的耦合泄漏。对于任何自定义DataSource重新实现该方法相对直接保持投影交换的目标语义同时遵循精简后的签名。详细说明见上游 PR #17395。FileOpenFuture错误类型从ArrowError改为DataFusionError变更前后对比FileOpenFuture类型别名的错误类型由ArrowError统一为DataFusionError。这影响FileOpener特质及所有涉及文件流式读取的实现。变更前pub type FileOpenFuture BoxFuturestatic, ResultBoxStreamstatic, ResultRecordBatch, ArrowError;变更后pub type FileOpenFuture BoxFuturestatic, ResultBoxStreamstatic, ResultRecordBatch;当前仓库中的定义datafusion/datasource/src/file_stream/mod.rs即为去掉内层错误类型参数的简化形态——内层Result与FileStream其余错误路径统一走DataFusionError。同时FileStreamState枚举的Open变体也做了相应调整。迁移动作如果你的代码有自定义FileOpener实现或直接处理FileOpenFuture需要将错误处理从ArrowError迁移到DataFusionError例如通过arrow_datafusion_err!宏或DataFusionError::from转换。详见上游 PR #17397。FFI 用户自定义聚合函数签名变更从return_type到return_fieldFFIForeign Function Interfacedatafusion/ffi/src/udaf/mod.rs中 UDAF 的 C ABI 结构体已改为调用底层聚合函数的return_field方法而非return_type以支持聚合函数返回字段的元数据处理。仓库中 datafusion/ffi/src/udaf/accumulator_args.rs 的AccumulatorArgs也携带了return_field: FieldRef与 FFI 侧新的return_fieldC 函数指针physical_expr/mod.rs配套。对大多数用户而言该变更透明但如果你编写过直接调用return_type的单元测试需要改为调用return_field。FFI 跨版本使用的注意事项这是 FFI API 的一次破坏性变更。当前最佳实践是确保所有交互的库使用相同的底层 Rust 版本以规避 ABI 不一致。上游 issue #17374 正在讨论稳定该接口使这些库未来可以跨不同 DataFusion 版本互操作。具体实现见上游 PR #17407。新增PhysicalExpr::is_volatile_node为正确标记易变volatile表达式——即每次求值可能返回不同结果的表达式典型如随机数函数——PhysicalExpr特质新增了is_volatile_node方法impl PhysicalExpr for MyRandomExpr { fn is_volatile_node(self) - bool { true } }该方法提供了默认值false以最小化破坏面但官方强烈建议所有PhysicalExpr实现者显式选择行为即使只是返回false。is_volatile_node在优化器中已有实际用途在公共子表达式消除CSE规则中datafusion/optimizer/src/common_subexpr_eliminate.rsis_valid检查通过!node.is_volatile_node()将易变表达式排除在可公共子表达式提取的范围之外避免对random()这类函数错误地做重复计算消除datafusion/optimizer/src/utils.rs 也将其用于表达式树遍历中的易变性判断。此外 datafusion/ffi/src/physical_expr/mod.rs 的 FFI 包装层同样暴露了该函数指针。测试覆盖可参见 datafusion/physical-expr/src/physical_expr.rs默认返回false与 higher_order_function.rs。更多讨论与示例实现见上游 PR #17351。迁移清单50.0.0 升级自查表为便于对照执行汇总一份升级自查清单工具链确认 Rust ≥ 1.86.0rust-toolchain.toml中当前主分支为 1.98.1锁定 50.0.0 时以 1.86.0 为准。分区表行为如不想要自动推断 Hive 分区设置datafusion.execution.listing_table_factory_infer_partitions false默认trueconfig.rs。UDF 特质为ScalarUDFImpl/AggregateUDFImpl/WindowUDFImpl实现PartialEq、Eq、Hash删除手写equals/hash_value无状态函数可用正则批量加 derive。异步 UDFinvoke_async_with_args返回ColumnarValue用ColumnarValue::from包裹ArrayRef并移除_option参数改从args.config_options读取。投影ProjectionExpr改为结构体用ProjectionExpr::new(expr, alias)构造、.expr/.alias访问、ProjectionExpr { expr, alias }匹配。配置访问标注类型的options()调用补.as_ref()返回ArcConfigOptions。schema 重写import 改为datafusion_physical_expr_adapter::{...}。Arrow/Parquet升级依赖到 56.0.0。自定义 ExecutionPlan按需实现reset_state有状态或持引用时必须。NLJ确认不依赖旧实现的输入排序保持必要时显式排序。LazyBatchGenerator补充as_any()返回self。DataSource按新签名重新实现try_swapping_with_projection。FileOpenFuture错误处理改用DataFusionError。FFI UDAF测试从return_type改为return_field保证交互库 Rust 版本一致。PhysicalExpr显式实现is_volatile_node易变表达式返回true。按此清单逐项排查即可完成向 DataFusion 50.0.0 的平滑升级。对于每项变更的更深入讨论与背景细节可继续查阅上文引用的对应仓库源码路径与上游 issue/PR。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐Social Analyzer安全加固终极指南10个关键防护策略防止未授权访问与滥用Social Analyzer安全加固终极指南10个关键防护策略防止未授权访问与滥用 Social Analyzer是一款强大的开源工具可通过API、CLI大数据数据分析后端Apache DataFusion 53.0.0 升级指南核心 API 破坏性变更与完整迁移手册Apache DataFusion 53.0.0 升级指南核心 API 破坏性变更与完整迁移手册 导读 本文以 Apache DataFusion 53.0.大数据数据分析后端探索no-neck-pain.nvim架构核心组件与实现原理探索no neck pain.nvim架构核心组件与实现原理 架构概览 no neck pain.nvim是一款旨在将当前聚焦的缓冲区居中显示在屏幕中间的Ne大数据数据仓库OLAP批处理后端上一篇jwt库进阶教程自定义Claims与高级验证策略下一篇Plus Jakarta Sans 字体终极使用指南如何为你的项目选择完美的开源字体创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表