库用户指南
扩展运算符
登录后可跨设备保存划线和私人笔记登录
扩展算子
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 撰写的博客文章。
评论
登录后参与评论
正在加载评论…
KnowForge