如何用Rust按分区读取ADLS v2上的Delta Lake数据
Rust实现ADLS v2上Delta Lake分区数据读取的方案
针对ADLS v2中3TB+的分区Delta表,你可以通过两种高效方式实现按分区读取,避免全量加载内存:
方案一:delta-rs获取分区文件 + DataFusion 流式读取
如果你已经通过delta-rs拿到了目标分区的文件URI,后续可以直接用DataFusion加载这些Parquet格式的文件并流式处理,无需全量加载:
步骤示例
- 依赖配置
在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"] }
- 初始化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(()) }
- 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:
步骤示例
- 依赖配置
额外添加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"] }
- 注册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
相关产品推荐
相关产品推荐

