库用户指南

扩展运算符

师成师成· 更新于 2026-09-28· 阅读 5 分钟· 0 次阅读

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

扩展算子

DataFusion 支持通过自定义的优化器规则来转换 LogicalPlan 和 ExecutionPlan,从而扩展算子能力。本节将使用 µWheel 项目来说明这些能力。

关于 DataFusion µWheel

DataFusion µWheel 是一个原生的 DataFusion 优化器,它通过自定义索引实现快速的时间维度聚合与剪枝,从而提升基于时间的分析查询性能。µWheel 与 DataFusion 的集成是与 DataFusion 社区共同完成的。

优化逻辑计划

rewrite 函数通过识别与已存储的 wheel 索引相匹配的时间模式和聚合函数来转换逻辑计划。当找到匹配项时,它会查询相应的索引以获取预计算的聚合值,将这些结果存入 MemTable,并将其作为新的 LogicalPlan::TableScan 返回。如果没有找到匹配项,原始计划将通过 DataFusion 的标准执行路径继续处理,不做任何改动。

fn rewrite(
  &self,
  plan: LogicalPlan,
  _config: &dyn OptimizerConfig,
) -> Result<Transformed<LogicalPlan>> {
    // Attempts to rewrite a logical plan to a uwheel-based plan that either provides
    // plan-time aggregates or skips execution based on min/max pruning.
    if let Some(rewritten) = self.try_rewrite(&plan) {
        Ok(Transformed::yes(rewritten))
    } else {
        Ok(Transformed::no(plan))
    }
}
// Converts a uwheel aggregate result to a TableScan with a MemTable as source
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)
}

想要更深入了解 µWheel 项目的用法,请访问 Max Meldrum 撰写的博客文章。

评论

登录后参与评论

正在加载评论…