如何用asyncio及Python3.11 TaskGroup从Azure Blobs并发读取Parquet至Pandas DataFrame
异步并发下载Azure Blob Parquet文件到Pandas DataFrame(Python 3.11 TaskGroup实现)
原代码的核心问题
download_blobs_async方法中,await asyncio.gather(*tasks)未返回任务结果,导致调用后得到None,无法通过result[0]访问下载流main函数错误地在异步函数内部调用asyncio.run,且直接执行main()而非通过asyncio.run(main())启动事件循环- 未利用Python 3.11的
TaskGroup特性简化并发任务管理 list_blobs_in_container_async返回的是BlobProperties对象,而非可直接用于下载的blob名称
修正后的完整代码
import logging import asyncio import pandas as pd from azure.storage.blob.aio import ContainerClient from io import BytesIO class BlobStorageAsync: def __init__(self, connection_string, container_name, logging_enable): self.connection_string = connection_string self.container_name = container_name self.container_client = ContainerClient.from_connection_string( conn_str=connection_string, container_name=container_name, logging_enable=logging_enable ) async def list_blobs_in_container_async(self, name_starts_with): blobs_list = [] async for blob in self.container_client.list_blobs(name_starts_with=name_starts_with): # 提取blob的名称属性,用于后续下载 blobs_list.append(blob.name) return blobs_list async def download_blob_async(self, blob_name): async with self.container_client.get_blob_client(blob=blob_name) as blob_client: stream = await blob_client.download_blob() data = await stream.readall() return BytesIO(data) async def download_blobs_async(self, blobs_list): # 使用Python 3.11+的TaskGroup管理并发任务 downloaded_streams = [] async with asyncio.TaskGroup() as tg: for blob_name in blobs_list: # 提交任务到TaskGroup并保存任务对象 task = tg.create_task(self.download_blob_async(blob_name)) downloaded_streams.append(task) # 取出所有任务的返回结果 return [task.result() for task in downloaded_streams] async def main(): # 替换为你的实际配置 connection_string = "你的Azure存储连接字符串" container_name = "你的容器名称" logging_enable = True # 方式1:手动指定要下载的blob名称 blobs_list = ["file1.parquet", "file2.parquet"] # 方式2:从容器中自动获取符合前缀的blob # BSA = BlobStorageAsync(connection_string, container_name, logging_enable) # blobs_list = await BSA.list_blobs_in_container_async(name_starts_with="parquet/") BSA = BlobStorageAsync(connection_string, container_name, logging_enable) # 获取所有下载完成的流列表 result_streams = await BSA.download_blobs_async(blobs_list) # 处理第一个下载文件的示例 if result_streams: df = pd.read_parquet(result_streams[0]) print(df.head()) if __name__ == '__main__': try: # 正确启动异步事件循环 asyncio.run(main()) except Exception as ex: print(f"错误: {ex}")
关键改动说明
- TaskGroup简化任务管理:Python 3.11引入的
TaskGroup自动处理任务的创建、调度和生命周期,无需手动维护任务列表和调用asyncio.gather,代码更简洁且能自动处理任务异常 - 结果收集逻辑:在TaskGroup中提交任务后,通过
task.result()获取每个任务的返回值,最终返回所有下载完成的BytesIO流列表 - Blob名称修正:
list_blobs_in_container_async中提取blob的name属性,确保后续下载时传入正确的标识 - 事件循环启动修正:在
__main__块中通过asyncio.run(main())正确启动异步事件循环,避免嵌套调用的错误 - 资源管理优化:
download_blob_async中直接用async with管理blob_client的生命周期,简化资源释放逻辑
内容的提问来源于stack exchange,提问作者LucaM
相关产品推荐
相关产品推荐

