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

如何用asyncio及Python3.11 TaskGroup从Azure Blobs并发读取Parquet至Pandas DataFrame

异步并发下载Azure Blob Parquet文件到Pandas DataFrame(Python 3.11 TaskGroup实现)

原代码的核心问题

  1. download_blobs_async方法中,await asyncio.gather(*tasks)未返回任务结果,导致调用后得到None,无法通过result[0]访问下载流
  2. main函数错误地在异步函数内部调用asyncio.run,且直接执行main()而非通过asyncio.run(main())启动事件循环
  3. 未利用Python 3.11的TaskGroup特性简化并发任务管理
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 12:20:41