如何用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
相关产品推荐
相关产品推荐

