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

如何用adlfs(fsspec)、pyarrow/pandas并发下载Blob并实现列投影与行过滤

高性能并发下载多个Azure Blob Parquet文件(带列投影和行过滤)

问题背景

我已经能用pandas实现单个Azure Blob上Parquet文件的列投影和行过滤读取,现在想借助pyarrow、fsspec/adlfs,用async/await异步代码实现多个Blob文件的并发下载,同时保留列投影和行过滤的优化,节省带宽、降低成本并提升速度。

单个文件的实现代码如下:

import pandas as pd
storage_options={"connection_string": "MY_CONNECTION_STRING_EXAMPLE"}
CONTAINER = "MY_CONTAINER_NAME"
FILE = "PATH/TO/FILE.parquet"
FILEPATH = f"az://{CONTAINER}/{FILE}"
# 列投影和行过滤示例
columns = ["col1", "col2"]
rows = [[("myindex", "<=", pd.Timestamp("2023-11-08 08:00:07")), ("myindex", ">=", pd.Timestamp("2023-11-08 08:00:02"))]]
df = pd.read_parquet(FILEPATH, storage_options=storage_options, columns=columns, filters=rows)

实现方案

1. 安装依赖

先确保装好所需工具包:

pip install pandas pyarrow fsspec adlfs aiohttp

2. 异步并发读取代码

用asyncio配合pyarrow的异步读取能力,结合adlfs的异步文件系统支持,实现多文件并发处理,同时保留列投影和行过滤:

import asyncio
import pyarrow.parquet as pq
import adlfs
import pandas as pd

async def read_single_parquet_async(file_path, storage_options, columns, filters):
    # 初始化异步Azure文件系统
    fs = adlfs.AzureBlobFileSystem(**storage_options)
    # 异步打开文件流
    async with fs.open(file_path, 'rb') as f:
        # 读取Parquet并应用列、行过滤
        parquet_file = pq.ParquetFile(f)
        table = parquet_file.read(columns=columns, filters=filters)
        return table.to_pandas()

async def main():
    storage_options={"connection_string": "MY_CONNECTION_STRING_EXAMPLE"}
    CONTAINER = "MY_CONTAINER_NAME"
    # 要处理的多个文件列表
    file_list = [
        f"az://{CONTAINER}/PATH/TO/FILE1.parquet",
        f"az://{CONTAINER}/PATH/TO/FILE2.parquet",
        f"az://{CONTAINER}/PATH/TO/FILE3.parquet"
    ]
    # 统一的列投影和行过滤参数
    columns = ["col1", "col2"]
    filters = [[("myindex", "<=", pd.Timestamp("2023-11-08 08:00:07")), ("myindex", ">=", pd.Timestamp("2023-11-08 08:00:02"))]]
    
    # 创建所有异步任务
    tasks = [read_single_parquet_async(file, storage_options, columns, filters) for file in file_list]
    # 等待所有任务完成,拿到结果
    results = await asyncio.gather(*tasks)
    # 合并所有结果DataFrame(按需选择)
    combined_df = pd.concat(results, ignore_index=True)
    print(combined_df.head())

if __name__ == "__main__":
    asyncio.run(main())

3. 并发数控制(可选)

如果文件数量特别多,建议加个并发数限制,避免触发Azure的请求频率阈值:

async def read_single_parquet_async(file_path, storage_options, columns, filters, semaphore):
    async with semaphore:
        fs = adlfs.AzureBlobFileSystem(**storage_options)
        async with fs.open(file_path, 'rb') as f:
            parquet_file = pq.ParquetFile(f)
            table = parquet_file.read(columns=columns, filters=filters)
            return table.to_pandas()

async def main():
    # ... 其他代码不变 ...
    # 限制同时并发5个任务
    semaphore = asyncio.Semaphore(5)
    tasks = [read_single_parquet_async(file, storage_options, columns, filters, semaphore) for file in file_list]
    # ... 其他代码不变 ...

4. 更简洁的批量处理方案(适合分区数据集)

如果你的Parquet文件是按规则分区存储的数据集,直接用pyarrow.dataset更省心,内部自动优化并发读取:

import pyarrow.dataset as ds
import adlfs
import pandas as pd

storage_options={"connection_string": "MY_CONNECTION_STRING_EXAMPLE"}
CONTAINER = "MY_CONTAINER_NAME"

fs = adlfs.AzureBlobFileSystem(**storage_options)
# 指向数据集根目录
dataset = ds.dataset(f"az://{CONTAINER}/PATH/TO/DATASET", filesystem=fs, format="parquet")
# 读取时直接应用列投影和行过滤
table = dataset.read(columns=["col1", "col2"], filters=[[("myindex", "<=", pd.Timestamp("2023-11-08 08:00:07")), ("myindex", ">=", pd.Timestamp("2023-11-08 08:00:02"))]])
combined_df = table.to_pandas()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 08:05:27