- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
本文围绕 DataFusion 官方文档《Extending Operators》展开,讲解如何通过自定义 OptimizerRule 重写LogicalPlan、以及如何实现ExecutionPlan接入物理执行这两个扩展点来定制查询行为。文中以社区项目 DataFusion µWheel 的真实改造为例,并结合本仓库中的源码接口定义与端到端测试 user_defined_plan.rs,帮助读者掌握从“逻辑计划重写”到“自定义算子产出结果”的完整链路。
一、扩展点概览:在哪一层扩展算子
DataFusion 中一条查询的处理链路大致是:SQL 解析生成LogicalPlan,优化器(Optimizer)通过一系列规则将其变换为语义等价但更高效的计划,再由物理规划器生成ExecutionPlan并执行。文档(见 extending-operators.md)给出的核心思路是:
- 逻辑层扩展:实现
OptimizerRule,在优化阶段识别特定模式(如时间序列上的聚合),把子树替换为基于自定义索引/缓存的新节点,甚至直接产出“计划期已算好的结果”; - 物理层扩展:实现
ExecutionPlan(或扩展节点ExtensionPlanNode),自定义物理算子,让重写后的逻辑节点在执行期真正落地。
官方文档选择 µWheel 项目作为贯穿示例,因为它同时演示了“计划期聚合并替换为内存表扫描”这种较激进的逻辑重写方式,下面按“规则定义 → 注册方式 → 实例剖析 → 物理扩展”的顺序展开。
二、OptimizerRule 接口:源码级解读
自定义优化规则的核心接口是OptimizerRule,定义在 optimizer.rs:
pub trait OptimizerRule: Debug { /// 规则的可读名称 fn name(&self) -> &str; /// 规则的应用方式;返回 None(默认值)时, /// 规则需要自行处理递归 fn apply_order(&self) -> Option<ApplyOrder> { None } /// 尝试将 plan 重写为优化后的形式: /// 重写成功返回 Transformed::yes,未重写返回 Transformed::no fn rewrite( &self, _plan: LogicalPlan, _config: &dyn OptimizerConfig, ) -> Result<Transformed<LogicalPlan>, DataFusionError>; }从 trait 的文档注释(同文件 L102-L142)可以提炼出三条实现约定,它们对保证优化器正确性至关重要:
- 无变化必须原样返回。优化器会反复调用
rewrite直到达到不动点(fixed point),因此当输入计划没有可转换的模式时,必须返回Transformed::no并原样返回计划,否则会触发反复重写甚至无限循环; - 避免按函数名字面匹配。注释明确建议通过
ScalarUDFImpl/AggregateUDFImpl等 trait 方法判断函数语义,而不是比较func.name() == "sum"这类字符串,因为注册函数可能被覆写、且同名函数的语义可能不同(注释中给出了datafusion-spark中sum的例子); - OptimizerRule 只改变效率、不改变语义。trait 注释指明:如果需要改变
LogicalPlan的语义,应实现AnalyzerRule而非优化器规则。
另外注意当前仓库中rewrite的返回类型是Result<Transformed<LogicalPlan>, DataFusionError>(见 optimizer.rs),Transformed来自datafusion_common::tree_node;早期文档中的示例签名省略了错误类型,接入新版 DataFusion 时需以仓库中的实际签名为准。
三、注册自定义规则:add_optimizer_rule 与构建期配置
规则写好之后需要注入会话。仓库提供了两种粒度(均可从源码确认):
运行时动态追加/移除—— 在 SessionContext 上:
/// 把一条优化规则追加到现有规则列表末尾 pub fn add_optimizer_rule( &self, optimizer_rule: Arc<dyn OptimizerRule + Send + Sync>, ) { self.state.write().append_optimizer_rule(optimizer_rule); } /// 按名称移除规则,返回是否移除成功 pub fn remove_optimizer_rule(&self, name: &str) -> bool { self.state.write().remove_optimizer_rule(name) }典型用法是ctx.add_optimizer_rule(Arc::new(MyRule::new())),追加的自定义规则会在 DataFusion 内置规则之后执行。同一文件中还提供add_analyzer_rule用于注册分析器规则(见 mod.rs)。
构建期静态配置—— 通过 SessionStateBuilder:
with_optimizer_rules(...):在SessionState构建阶段替换/追加逻辑优化器规则列表;with_physical_optimizer_rules(...):对应物理优化器规则(PhysicalOptimizerRule),用于在ExecutionPlan层做等价重写。
从源码结构看,SessionStateBuilder中optimizer_rules与physical_optimizer_rules是两组相互独立的可选列表(见 session_state.rs),构建过程中会逐条装配进最终的规则链。这意味着“逻辑层重写”与“物理层重写”可以在同一个会话中共存,µWheel 这类项目通常只需要前者,而需要更贴近执行的改写则选后者。
四、实例剖析:DataFusion µWheel 的计划期聚合
µWheel 是一个与 DataFusion 社区合作集成的原生优化器,面向时间序列分析场景:它用自定义的“时间轮”索引(wheel index)保存按时间粒度预计算的聚合值,使SELECT sum(v) FROM t WHERE ts BETWEEN ...这类查询可以在不扫描明细数据的情况下直接命中索引结果。官方文档给出的关键实现有两段。
4.1 逻辑计划重写入口
µWheel 实现了OptimizerRule::rewrite:
fn rewrite( &self, plan: LogicalPlan, _config: &dyn OptimizerConfig, ) -> Result<Transformed<LogicalPlan>> { // 尝试把逻辑计划改写为 uwheel 计划: // 要么提供计划期聚合,要么基于 min/max 剪枝跳过执行 if let Some(rewritten) = self.try_rewrite(&plan) { Ok(Transformed::yes(rewritten)) } else { Ok(Transformed::no(plan)) } }这段代码精确体现了第二节总结的两条约定:命中时间模式(时间谓词 + 匹配的聚合函数 + 已建立的轮索引)时返回Transformed::yes并交出重写的计划;未命中时返回Transformed::no并原样透传,让查询继续走 DataFusion 的标准执行路径。
其工作方式为:rewrite识别出时间谓词与聚合模式后,查询对应索引取回预计算聚合值;没有可用索引匹配时,则退化为 min/max 剪枝(pruning)判断能否直接跳过执行;仍不行就保持原计划不变。
4.2 把聚合结果“物化”为 TableScan
命中索引后,µWheel 并不会引入一个全新类型的执行节点,而是把标量结果包成一张内存表,让后续逻辑/物理计划完全复用标准TableScan通路:
// 将 uwheel 聚合结果转换为以 MemTable 为源的 TableScan fn agg_to_table_scan(result: f64, schema: SchemaRef) -> Result<LogicalPlan> { let data = Float64Array::from(vec![result]); let record_batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(data)])?; let df_schema = Arc::new(DFSchema::try_from(schema.clone())?); let mem_table = MemTable::try_new(schema, vec![vec![record_batch]])?; mem_table_as_table_scan(mem_table, df_schema) }这个技巧值得借鉴:当你想让优化器“短路”一段昂贵计算时,不一定要造新算子——把结果封装为MemTable并返回LogicalPlan::TableScan,下游投影、join、limit 等所有既有逻辑都能无缝消费该结果,也省去了自定义物理节点的注册与代码生成工作。若查询结果本身就是明细数据(而非标量),µWheel 同样可以用MemTable承载从索引展开的记录批次。
µWheel 项目由 Max Meldrum 主导开发并发布了多篇技术博客深入讲解其与 DataFusion 的集成(项目与博客为仓库外部资源,本仓库文档中不再保留外部链接);其价值在于示范了“外部索引 + 计划期重写”这一类扩展在 DataFusion 优化器框架内的落法。
五、物理层扩展:自定义 ExecutionPlan 与 End-to-End 验证
文档标题中的“operators”不止逻辑重写这一半。当重写产出的节点需要在执行期做真正的算子级行为(如 Top-K 流式维护、外部引擎执行),就需要实现ExecutionPlantrait,定义见 execution_plan.rs。从 trait 的签名看,自定义物理节点通常需要处理:name()/static_name()提供节点短名、schema()与properties()描述输出模式与分区/排序等物理属性、check_invariants()校验节点不变量(默认实现调用check_default_invariants),以及执行期驱动execute()产出SendableRecordBatchStream等职责。
仓库内提供了一个完整的端到端示例:user_defined_plan.rs。该测试文件头部注释自述其演示内容:
- 定义一个
TopKNode,实现扩展节点接口; - 编写一条
OptimizerRule,把SELECT ... ORDER BY revenue DESC LIMIT 3这类“全排序后丢弃”的逻辑计划重写为使用TopKNode的计划(朴素计划会先完整排序再只保留 3 行,Top-K 只需维护一个大小为 K 的缓冲,显著减少内存占用); - 实现从逻辑节点创建
ExecutionPlan的物理规划路径,并真正跑通产出结果。
这个测试覆盖了本文主题的关键闭环:逻辑层用OptimizerRule做模式匹配与替换,物理层用自定义ExecutionPlan承接执行,可以作为自研算子(例如接列式存储直读、外部执行引擎、Top-K 算子)的参照实现直接阅读。
六、相关扩展能力与延伸阅读
DataFusion 的扩展体系是分工明确的几个正交入口,可结合仓库文档按需组合:
- extensions.md:扩展机制总览(UDF、扩展节点等);
- extending-sql.md:扩展 SQL 语法与
SqlToRel规划阶段(ExprPlanner / TypePlanner / RelationPlanner),适合在“解析/规划”更早的阶段介入; - custom-table-providers.md:自定义表提供方,覆盖数据源侧的扩展;
- query-optimizer.md:理解内置优化器规则的组织方式,有助于决定自定义规则插入的位置。
适用前提与限制:本文所述接口以当前仓库源码为准,其中OptimizerRule::rewrite返回Result<Transformed<LogicalPlan>, DataFusionError>、apply_order()缺省表示规则自行处理递归等细节均来自 optimizer.rs;supports_rewrite方法在 47.0.0 起已标记弃用,新实现不必依赖它。自定义规则必须遵守“不改语义、无变化即透传”的约定,否则可能破坏优化器的不动点迭代或产生错误结果;OptimizerRule规则追加在内置规则之后,意味着你的规则看到的是内置简化/谓词下推之后的计划形态,这既是便利(模式更规整)也是约束(某些中间形态可能已被改写)。
- 大数据
- 数据分析
- 后端
【免费下载链接】datafusion
Apache DataFusion SQL Query Engine
相关推荐
解决结构化数据多语言转换的技术挑战:jsontt 架构设计与实战指南
解决结构化数据多语言转换的技术挑战:jsontt 架构设计与实战指南 在全球化软件开发的背景下,多语言支持已成为现代应用的基本要求。然而,处理复杂的 JSON
开发工具CLIAI 应用SwiftUI-2048核心组件解析:BlockView与游戏界面构建
SwiftUI 2048核心组件解析:BlockView与游戏界面构建 想要学习如何使用SwiftUI构建优雅的2048游戏吗?🎮 本文将深入解析SwiftU
DeepSeek-R1-Distill-Qwen-14B模型更新与迁移指南:从旧版本到新版本的平滑过渡
DeepSeek R1 Distill Qwen 14B模型更新与迁移指南:从旧版本到新版本的平滑过渡 DeepSeek R1 Distill Qwen 14B
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考