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 分布式执行、InfluxDB IOx、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>Partitioned Hash Join</strong>:当 build side 太大无法装入内存时分区 hash
  • <strong>Collect Left Join</strong>:流式收集 + 延迟 probe
  • <strong>LeftSemi/LeftAnti Join</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:

  1. Phase 1:将两表按 hash(key) % N 分区写入 N 个 spill 文件
  2. Phase 2:逐对读取分区文件,在内存中完成 Join
  3. 这保证即使在 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>:

    • 不是 ClickHouse 的替代,而是它们之上的查询层
    • 不是 DuckDB 的竞争,而是 Rust 生态发展的产物
    • 核心价值:允许 Rust 应用内嵌完整 SQL 引擎,0 外部依赖

    十、6.x 版本新特性与未来方向

    DataFusion 保持每月发布节奏,近期重大更新:

    • <strong>GroupedHashAggregate</strong>:减少 Hash Join 后的二次扫描
    • <strong>ReadUnionExec 并行</strong>:UNION ALL 各分支并行执行
    • <strong>TopK 内存优化</strong>:ORDER BY LIMIT 的内存从 O(N) 降到 O(K)
    • <strong>Logical Plan 稳定 v2</strong>:跨版本兼容性保证
    • <strong>Substrait 计划互操作</strong>:与其他引擎交换序列化计划

    未来路线图重点:

    1. <strong>分布式执行</strong>:Ballista 已整合 Feldera 的分布式调度器
    2. <strong>流处理增强</strong>:低延迟 watermark + state backend
    3. <strong>Cost-based Optimizer (CBO)</strong>:基于统计信息的 Join 重排序
    4. <strong>WASM 支持</strong>:浏览器中运行精简版 DataFusion

    5. 十一、总结

      DataFusion 证明了用 Rust 重写数据库核心组件的可行性和优势:

      1. <strong>零成本抽象</strong>:编译期排除 GC 和数据竞争
      2. <strong>Arrow 原生达成零拷贝</strong>:与 Python/JS 生态共享内存
      3. <strong>嵌入式友好</strong>:无需独立部署,只需 `cargo add datafusion`
      4. 对于需要在 Rust 应用中集成 SQL 能力的团队——无论是构建时序数据库(GreptimeDB)、日志分析(InfluxDB IOx)、还是湖仓查询层(delta-rs),DataFusion 都是当前最优选择。


        参考资源

        • DataFusion GitHub: https://github.com/apache/datafusion
        • DataFusion 中文文档: https://datafusion.apache.org/
        • Book: "Query Engine Design with DataFusion"
        • SIMD in Arrow: https://arrow.apache.org/blog/2022/11/07/simd-arrow-internals/
        • Ballista: https://github.com/apache/datafusion-ballista
点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部