如何扩展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))
关键修改说明
- 调整入口参数:将
main函数的参数从单个customer_id改为接收customer_id_list,适配多文件场景 - 批量创建任务:通过列表推导式为每个客户ID生成对应的下载异步任务
- 并发执行任务:使用
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
相关产品推荐
相关产品推荐

