DataFusion: Rust 原生查询引擎深入工程
现代数据分析系统中,查询引擎是连接存储与计算的桥梁。本文深入剖析 Apache DataFusion——用 Rust 编写的高性能查询引擎,探讨其架构设计、优化器实现、执行模型以及在生产环境中的工程实践。
一、为什么需要一个新的查询引擎
在大数据生态中,Spark SQL、Presto、ClickHouse 各有优势,但存在一些共性问题:
<strong>JVM负担</strong>:Spark/Presto依赖GC,内存管理颗粒度粗<strong>扩展困难</strong>:UDF/UDAF需要跨语言桥接<strong>启动延迟</strong>:JVM冷启动秒级,不适合Serverless场景<strong>资源隔离差</strong>:多租户场景下内存/CPU隔离依赖外部机制
DataFusion 用 Rust 从头构建,实现了:
零GC停顿,内存由Rust所有权系统管理原生Arrow内存格式,列式计算零拷贝核心库可嵌入任何Rust应用(如Ballista分布式执行、InfluxDBIOx、GreptimeDB)亚毫秒级启动,适合边缘计算与Serverless
二、核心架构分层
DataFusion 采用经典数据库分层架构,但每一层都做了 Rust-idiomatic 的重构:
┌─────────────────────────────────────────────────┐
│ SQL / DataFrame API │
├─────────────────────────────────────────────────┤
│ Logical Plan (抽象语法树) │
├────────────┬────────────────────────────────────┤
│ Analyzer │ Optimizer (规则+代价) │
├────────────┴────────────────────────────────────┤
│ Physical Plan (算子执行图) │
├─────────────────────────────────────────────────┤
│ Execution Engine (向量化/SIMD) │
├─────────────────────────────────────────────────┤
│ Memory Management (Arrow + Spill) │
├─────────────────────────────────────────────────┤
│ Object Store / Table Provider │
└─────────────────────────────────────────────────┘
<strong>关键设计点</strong>:
<strong>向量化执行</strong>:所有算子操作Arrow`RecordBatch`,利用SIMD加速<strong>流式执行</strong>:顶层是`SendableRecordBatchStream`,从构建之初就支持异步流<strong>物化中间结果</strong>:必要时spill到磁盘,保证OOM安全
三、从 SQL 到 RecordBatch:完整流程
看一个实际的查询流程,理解 DataFusion 如何处理一条 SQL:
SELECT department, AVG(salary) as avg_sal
FROM employees
WHERE hire_date > '2020-01-01'
GROUP BY department
HAVING AVG(salary) > 50000
ORDER BY avg_sal DESC
LIMIT 10;
3.1 解析与逻辑计划生成
DataFusion 使用 `sqlparser-rs` 解析 SQL,生成初始 Logical Plan:
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> Result<()> {
let ctx = SessionContext::new();
// 注册 Parquet 文件作为表
ctx.register_parquet("employees", "data/employees.parquet",
ParquetReadOptions::default()).await?;
let df = ctx.sql("SELECT department, AVG(salary) FROM employees GROUP BY department").await?;
// 查看逻辑计划
df.logical_plan().indent();
Ok(())
}
3.2 优化器规则链
DataFusion 的优化器采用 Cascades 风格的规则引擎,按阶段应用转换:
// 核心优化规则(简化版)
let rules: Vec<Arc<dyn OptimizerRule>> = vec![
Arc::new(SimplifyExpressions::new()), // 常量折叠/化简
Arc::new(EliminateFilter::new()), // 去除恒真 Filter
Arc::new(ReduceCrossJoin::new()), // Join 条件提取
Arc::new(ProjectionPushdown::new()), // 列裁剪下推
Arc::new(FilterPushdown::new()), // 谓词下推
Arc::new(LimitPushdown::new()), // Limit 下推
Arc::new(AggregateStatistics::new()), // 聚合统计信息推断
Arc::new(HashBuildProbeOrder::new()), // Hash Join 大小表选择
Arc::new(JoinSelection::new()), // Join 算法选择
];
每条规则都是幂等的:应用一次无变化则跳过。优化器迭代应用直到 Logical Plan 收敛。
3.3 物理计划生成
Logical Plan 经过代价模型转换为可执行的 Physical Plan:
GlobalLimitExec: fetch=10
SortExec: [avg_sal DESC]
AggregateExec: mode=Final, gby=[department], aggr=[AVG(salary)]
CoalesceBatchesExec: target_batch_size=8192
AggregateExec: mode=Partial, gby=[department], aggr=[AVG(salary)]
FilterExec: hire_date > 2020-01-01
ParquetExec: groups={20 batches}, projection=[department, salary, hire_date]
关键算子:
<strong>ParquetExec</strong>:配合`object_store`crate实现谓词下推、行组裁剪、列投影<strong>AggregateExec</strong>:两阶段聚合(Partial→Final),Partial阶段可分布式预聚合<strong>SortExec</strong>:TopN优化,避免全量排序
四、向量化执行与 SIMD 加速
DataFusion 的执行核心是 `RecordBatch`——Arrow 格式的列式数据块。对它的操作天然适合 SIMD:
// 示例:Arrow 数组的向量化比较(底层调用 SIMD 指令)
use arrow::array::*;
use arrow::compute::kernels::cmp::gt;
// arr: Int64Array = [10, 20, 30, 40, 50]
// 返回 BooleanArray: [false, false, true, true, true]
let result = gt(&arr, &Value::Int64(25)).unwrap();
// 聚合类的 SIMD 加速:sum/min/max 自动向量化
use arrow::compute::sum;
let total = sum(&arr).unwrap(); // 底层 AVX2/NEON 加速
更复杂的例子——带 Filter 的聚合:
use arrow::compute::{filter, sum};
use arrow::array::{Int64Array, BooleanArray};
// 只对 salary > 50000 的行求和
let salary: &Int64Array = batch.column(2).as_any().downcast_ref().unwrap();
let mask: &BooleanArray = batch.column(3).as_any().downcast_ref().unwrap();
let filtered = filter(salary, mask).unwrap();
let result = sum(filtered.as_ref()).unwrap();
在生产环境中,对 10 亿行数据执行 `SELECT SUM(x) WHERE y > 0.5`,DataFusion 的向量化执行比逐行解释快 <strong>50-100x</strong>。
五、内存管理与 Spill 机制
Rust 的所有权系统让内存泄漏几乎不可能,但查询引擎需要额外的策略:
5.1 MemoryPool
DataFusion 内置两种内存池:
use datafusion::execution::memory_pool::*;
// 1. 公平池:所有算子按请求顺序分配,先到先得
let pool = FairPool::new(4 * 1024 * 1024 * 1024u64); // 4GB
// 2. 贪婪池:HashJoin 等需要大内存的算子优先
let pool = GreedyPool::new(4 * 1024 * 1024 * 1024u64);
5.2 Spill to Disk
当聚合/排序的中间结果超过内存限制,DataFusion 自动将其序列化到临时文件:
use datafusion::prelude::*;
let mut config = SessionConfig::new();
config = config.with_repartition_aggregations(true); // Enable spill-friendly repartition
config.options_mut().execution.memory_limit = Some(2 * 1024 * 1024 * 1024); // 2GB
config.options_mut().execution.spill_file_prefix = Some(std::path::PathBuf::from("/tmp/df_spill"));
let ctx = SessionContext::with_config(config);
// 即使查询需要 100GB 中间结果,也会自动 spill 成 10 个磁盘文件后归并
let result = ctx.sql("SELECT user_id, count(*) FROM clicks GROUP BY user_id LIMIT 1M").await?;
Spill 采用 Arrow IPC 格式序列化,反序列化代价极低(本质是指针操作)。
六、Join 算法实现与优化
Join 是 OLAP 中最昂贵的操作,DataFusion 实现了多种策略:
6.1 Hash Join
use datafusion::prelude::*;
// DataFusion 自动选择 build side
// 小表构建 Hash Table,大表 Probe
let df = ctx.sql("
SELECT o.order_id, o.amount, c.name
FROM orders o
JOIN customers c ON o.customer_id = c.id
").await?;
底层使用 `hashbrown` crate(Rust 标准库 HashMap 的高性能 fork),支持:
<strong>PartitionedHashJoin</strong>:当buildside太大无法装入内存时分区hash<strong>CollectLeftJoin</strong>:流式收集+延迟probe<strong>LeftSemi/LeftAntiJoin</strong>:只输出匹配/不匹配的左表行
6.2 Merge Join
对已排序的输入,Merge Join 是 O(N+M) 的最优选择:
// 对两个已经按 join key 排序的输入执行 stream merge join
// 内存占用极低,适合超大数据集
let df = ctx.sql("
SELECT a.*, b.*
FROM (SELECT * FROM a ORDER BY id) a
JOIN (SELECT * FROM b ORDER BY id) b ON a.id = b.id
").await?;
6.3 Grace Hash Join
当 Hash Table 超过内存限制时,DataFusion 退化为 Grace Hash Join:
Phase1:将两表按hash(key)%N分区写入N个spill文件Phase2:逐对读取分区文件,在内存中完成Join不是ClickHouse的替代,而是它们之上的查询层不是DuckDB的竞争,而是Rust生态发展的产物核心价值:允许Rust应用内嵌完整SQL引擎,0外部依赖<strong>GroupedHashAggregate</strong>:减少HashJoin后的二次扫描<strong>ReadUnionExec并行</strong>:UNIONALL各分支并行执行<strong>TopK内存优化</strong>:ORDERBYLIMIT的内存从O(N)降到O(K)<strong>LogicalPlan稳定v2</strong>:跨版本兼容性保证<strong>Substrait计划互操作</strong>:与其他引擎交换序列化计划<strong>分布式执行</strong>:Ballista已整合Feldera的分布式调度器<strong>流处理增强</strong>:低延迟watermark+statebackend<strong>Cost-basedOptimizer(CBO)</strong>:基于统计信息的Join重排序<strong>WASM支持</strong>:浏览器中运行精简版DataFusion<strong>零成本抽象</strong>:编译期排除GC和数据竞争<strong>Arrow原生达成零拷贝</strong>:与Python/JS生态共享内存<strong>嵌入式友好</strong>:无需独立部署,只需`cargoadddatafusion`DataFusionGitHub:https://github.com/apache/datafusionDataFusion中文文档:https://datafusion.apache.org/Book:"QueryEngineDesignwithDataFusion"SIMDinArrow:https://arrow.apache.org/blog/2022/11/07/simd-arrow-internals/Ballista:https://github.com/apache/datafusion-ballista
这保证即使在 1GB 内存限制下也能完成 1TB 级别的 Join。
七、Table Provider 与自定义数据源
DataFusion 的核心抽象是 `TableProvider` trait,允许嵌入任意数据源:
use datafusion::datasource::TableProvider;
use datafusion::arrow::datatypes::{Schema, Field, DataType};
use async_trait::async_trait;
struct MyCustomTable {
schema: SchemaRef,
object_store: Arc<dyn ObjectStore>,
path: String,
}
#[async_trait]
impl TableProvider for MyCustomTable {
fn as_any(&self) -> &dyn Any { self }
fn schema(&self) -> SchemaRef { self.schema.clone() }
async fn scan(
&self,
state: &SessionState,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
// 返回自定义的 ExecutionPlan
Ok(Arc::new(MyExecPlan::new(
self.object_store.clone(),
self.path.clone(),
projection.cloned(),
filters.to_vec(),
limit,
)))
}
}
实战案例:接入 LakeFS / Delta Lake / Iceberg:
// delta-rs 已内置 DataFusion 集成
use deltalake::open_table;
let table = open_table("s3://lake/orders").await?;
ctx.register_table("orders", Arc::new(table))?;
// Delta Log 的查询也能被 DataFusion 优化
let df = ctx.sql("SELECT * FROM orders WHERE date >= '2024-01-01'").await?;
八、生产环境工程实践
8.1 性能调优参数
use datafusion::prelude::*;
let config = SessionConfig::new()
.with_target_partitions(8) // 并行度
.with_repartition_joins(true) // Join 前自动重分区
.with_repartition_aggregations(true) // 聚合前自动重分区
.with_parquet_pruning(true) // Parquet 行组裁剪
.with_collect_statistics(true); // 收集统计信息供优化器使用
// 针对大查询的内存控制
let runtime = RuntimeEnv::new(
RuntimeConfig::new()
.with_memory_limit(8 * 1024 * 1024 * 1024, 1.0) // 8GB, 使用比例 100%
.with_disk_manager(DiskManagerConfig::NewTemp {
prefix: Some("df_spill".into()),
})
)?;
let ctx = SessionContext::new_with_config_rt(config, Arc::new(runtime));
8.2 查询计划可视化
// 用 DOT 格式输出查询计划,方便优化分析
let df = ctx.sql("SELECT * FROM t WHERE x > 10 ORDER BY y LIMIT 5").await?;
let plan = df.create_physical_plan().await?;
println!("{}", displayable(plan.as_ref()).indent());
输出示例:
GlobalLimitExec: fetch=5
SortExec: [y ASC NULLS LAST]
CoalesceBatchesExec: target_batch_size=8192
FilterExec: x > 10
ParquetExec: groups={50 batches}, projection=[x, y]
8.3 异步执行与背压
DataFusion 的执行基于 Tokio 异步运行时,天然支持背压:
let df = ctx.sql("SELECT * FROM large_table").await?;
let mut stream = df.execute_stream().await?;
while let Some(batch) = stream.next().await {
let batch = batch?;
// 消费端缓慢处理不会导致内存积压
process(batch).await;
}
消费速度慢时,上游算子会自动降速(Tokio 的 Semaphore 机制)。
8.4 与 Arrow Flight SQL 集成
DataFusion 可以作为 Arrow Flight SQL 服务端:
use datafusion::execution::context::SessionContext;
use arrow_flight::flight_service_server::FlightServiceServer;
let ctx = SessionContext::new();
let service = FlightServiceExt::new(ctx);
// 客户端通过 gRPC 发送查询,返回 Arrow IPC 流
// 零序列化开销,适合微服务间通信
九、DataFusion vs 其他引擎对比
特性 |
DataFusion |
Spark SQL |
DuckDB |
ClickHouse |
|---|---|---|---|---|
运行时 |
Rust (无 GC) |
JVM |
C++ |
C++ |
内存格式 |
Arrow 原生 |
Arrow 2.0 |
自有列存 |
自有列存 |
启动方式 |
库嵌入 |
独立进程 |
库嵌入 |
独立服务 |
优化器 |
Cascades |
Cascades |
规则+启发式 |
规则+启发式 |
扩展性 |
UDF/UDAF Rust |
JVM UDF |
C++ 扩展 |
C++ 扩展 |
分布式 |
Ballista/Sail |
原生 |
无 |
原生 |
适用场景 |
嵌入式 SQL 层、流处理 |
批处理集群 |
本地分析 |
OLAP 实时查询 |
<strong>DataFusion 的独特定位</strong>:
十、6.x 版本新特性与未来方向
DataFusion 保持每月发布节奏,近期重大更新:
未来路线图重点:
十一、总结
DataFusion 证明了用 Rust 重写数据库核心组件的可行性和优势:
对于需要在 Rust 应用中集成 SQL 能力的团队——无论是构建时序数据库(GreptimeDB)、日志分析(InfluxDB IOx)、还是湖仓查询层(delta-rs),DataFusion 都是当前最优选择。

发表评论 取消回复