You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Rust按分区读取ADLS v2上的Delta Lake数据

Rust实现ADLS v2上Delta Lake分区数据读取的方案

针对ADLS v2中3TB+的分区Delta表,你可以通过两种高效方式实现按分区读取,避免全量加载内存:


方案一:delta-rs获取分区文件 + DataFusion 流式读取

如果你已经通过delta-rs拿到了目标分区的文件URI,后续可以直接用DataFusion加载这些Parquet格式的文件并流式处理,无需全量加载:

步骤示例

  1. 依赖配置
    在Cargo.toml中添加必要依赖:
[dependencies]
delta-rs = { version = "0.16", features = ["azure", "parquet"] }
datafusion = "32.0"
azure_storage_blobs = "0.18"
tokio = { version = "1.0", features = ["full"] }
  1. 初始化ADLS存储并获取分区文件
    如果之前获取分区文件的逻辑有问题,可参考以下代码:
use delta_rs::storage::azure::AzureStorageBackend;
use delta_rs::DeltaTable;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 配置ADLS v2存储信息
    let account = "your-adls-account";
    let container = "your-container";
    let table_path = "delta-table-path";
    // 认证方式:支持SAS token、Azure AD服务主体等,这里以SAS为例
    let sas_token = "your-sas-token";
    
    // 初始化Azure存储后端
    let mut backend = AzureStorageBackend::new(account, container, sas_token)?;
    backend.set_use_https(true);
    
    // 加载Delta表
    let table = DeltaTable::new(backend, table_path).await?;
    
    // 指定分区过滤条件
    let partition_filters = vec![
        ("event_date", "=", "2024-05-20"),
        ("event_type", "=", "click"),
        ("sample_freq", "=", "1h"),
    ];
    
    // 获取目标分区的所有文件
    let files = table.get_files(Some(partition_filters)).await?;
    // 转换为ADLS完整URI列表
    let file_uris: Vec<String> = files
        .into_iter()
        .map(|f| format!("https://{}.dfs.core.windows.net/{}/{}", account, container, f.path))
        .collect();
    
    // 后续用DataFusion读取这些文件
    read_files_with_datafusion(file_uris).await?;
    
    Ok(())
}
  1. DataFusion流式读取文件
    用DataFusion加载Parquet文件并逐批处理,避免内存溢出:
use datafusion::prelude::*;

async fn read_files_with_datafusion(file_uris: Vec<String>) -> Result<(), Box<dyn std::error::Error>> {
    // 创建DataFusion会话上下文
    let ctx = SessionContext::new();
    
    // 注册Parquet文件为临时视图
    ctx.register_parquet("partition_data", &file_uris, ParquetReadOptions::default()).await?;
    
    // 执行自定义查询
    let df = ctx.sql("SELECT * FROM partition_data").await?;
    
    // 流式处理结果,逐批输出
    let stream = df.execute_stream().await?;
    tokio::pin!(stream);
    
    while let Some(batch) = stream.next().await {
        let batch = batch?;
        // 处理当前批次数据,比如写入下游存储、统计计算等
        println!("处理批次:{}行数据", batch.num_rows());
    }
    
    Ok(())
}

方案二:直接用DataFusion结合delta-rs扩展实现分区下推

更高效的方式是让DataFusion自动识别Delta表的分区信息,通过SQL的WHERE子句直接过滤分区,无需手动获取文件URI:

步骤示例

  1. 依赖配置
    额外添加delta-rs的datafusion扩展:
[dependencies]
delta-rs = { version = "0.16", features = ["azure", "datafusion"] }
datafusion = "32.0"
azure_storage_blobs = "0.18"
tokio = { version = "1.0", features = ["full"] }
  1. 注册Delta表并执行分区过滤查询
use delta_rs::datafusion::DeltaTableProvider;
use delta_rs::storage::azure::AzureStorageBackend;
use datafusion::prelude::*;
use std::sync::Arc;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 配置ADLS存储信息
    let account = "your-adls-account";
    let container = "your-container";
    let table_path = "delta-table-path";
    let sas_token = "your-sas-token";
    
    // 初始化Azure存储后端
    let mut backend = AzureStorageBackend::new(account, container, sas_token)?;
    backend.set_use_https(true);
    
    // 创建DeltaTableProvider
    let table_provider = DeltaTableProvider::new(backend, table_path).await?;
    
    // 创建DataFusion会话上下文并注册Delta表
    let ctx = SessionContext::new();
    ctx.register_table("delta_events", Arc::new(table_provider))?;
    
    // 执行带分区过滤的SQL查询,DataFusion会自动下推过滤到分区层,仅加载目标分区文件
    let sql = r#"
        SELECT * FROM delta_events 
        WHERE event_date = '2024-05-20' 
          AND event_type = 'click' 
          AND sample_freq = '1h'
    "#;
    let df = ctx.sql(sql).await?;
    
    // 流式处理结果
    let stream = df.execute_stream().await?;
    tokio::pin!(stream);
    
    while let Some(batch) = stream.next().await {
        let batch = batch?;
        println!("处理批次:{}行数据", batch.num_rows());
        // 自定义业务逻辑
    }
    
    Ok(())
}

关键注意事项

  • ADLS认证:除SAS token外,还支持Azure AD服务主体认证,需调整AzureStorageBackend的初始化逻辑适配。
  • 内存控制:始终使用execute_stream进行流式处理,避免调用collect等会全量加载数据的方法。
  • 分区列类型匹配:确保过滤条件中的值类型与Delta表定义的分区列类型一致(比如event_date是字符串还是日期类型)。

内容的提问来源于stack exchange,提问作者vekeras

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 13:20:45