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

如何扩展Python异步代码实现多Azure Blobs批量下载?

实现Azure Blob异步多文件下载

要实现多文件同时异步下载,核心是利用asyncio.gather并发执行多个下载任务,以下是修改后的完整代码:

import asyncio
from azure.storage.blob.aio import BlobServiceClient

async def download_blob_to_file(blob_service_client: BlobServiceClient, container_name, transaction_date, customer_id):
    blob_client = blob_service_client.get_blob_client(container=container_name, blob=f"{transaction_date}/{customer_id}.csv")
    with open(file=f'{customer_id}.csv', mode="wb") as sample_blob:
        download_stream = await blob_client.download_blob()
        data = await download_stream.readall()
        sample_blob.write(data)


async def main(transaction_date, customer_id_list):
    connect_str = "connection-string"
    blob_serv_client = BlobServiceClient.from_connection_string(connect_str)

    async with blob_serv_client as blob_service_client:
        # 批量创建所有下载任务
        download_tasks = [
            download_blob_to_file(blob_service_client, "sample-container", transaction_date, customer_id)
            for customer_id in customer_id_list
        ]
        # 并发执行所有下载任务
        await asyncio.gather(*download_tasks)

if __name__ == '__main__':
    transaction_date = '20240409'
    customer_id_list = ['001', '002', '003', '004']
    asyncio.run(main(transaction_date, customer_id_list))

关键修改说明

  1. 调整入口参数:将main函数的参数从单个customer_id改为接收customer_id_list,适配多文件场景
  2. 批量创建任务:通过列表推导式为每个客户ID生成对应的下载异步任务
  3. 并发执行任务:使用asyncio.gather(*download_tasks)一次性启动所有任务,实现多文件同时下载

可选优化:限制并发数

如果需要避免同时发起过多请求导致资源占用过高,可以用asyncio.Semaphore控制并发数量,示例如下:

import asyncio
from azure.storage.blob.aio import BlobServiceClient

async def download_blob_to_file(blob_service_client: BlobServiceClient, container_name, transaction_date, customer_id, semaphore):
    # 用信号量限制并发数
    async with semaphore:
        blob_client = blob_service_client.get_blob_client(container=container_name, blob=f"{transaction_date}/{customer_id}.csv")
        with open(file=f'{customer_id}.csv', mode="wb") as sample_blob:
            download_stream = await blob_client.download_blob()
            data = await download_stream.readall()
            sample_blob.write(data)


async def main(transaction_date, customer_id_list):
    connect_str = "connection-string"
    blob_serv_client = BlobServiceClient.from_connection_string(connect_str)
    # 设置最大并发数为3,可根据实际情况调整
    semaphore = asyncio.Semaphore(3)

    async with blob_serv_client as blob_service_client:
        download_tasks = [
            download_blob_to_file(blob_service_client, "sample-container", transaction_date, customer_id, semaphore)
            for customer_id in customer_id_list
        ]
        await asyncio.gather(*download_tasks)

if __name__ == '__main__':
    transaction_date = '20240409'
    customer_id_list = ['001', '002', '003', '004']
    asyncio.run(main(transaction_date, customer_id_list))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 03:58:10