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

异步下载多个Blob - Azure Function运行事件循环检测问题

解决Azure Function中异步Blob下载的事件循环冲突问题

问题核心

Azure Function运行时本身已启动事件循环,直接调用asyncio.run()会触发asyncio.run() cannot be called from a running event loop错误——因为asyncio.run()会尝试创建新循环,与当前线程的活跃循环冲突。

代码修改方案

核心逻辑是检测当前线程是否存在活跃事件循环,根据结果选择不同的任务执行方式:

修改后的完整代码:

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):
    connect_str = "connection-string"
    blob_serv_client = BlobServiceClient.from_connection_string(connect_str)

    async with blob_serv_client as blob_service_client:
        await download_blob_to_file(blob_service_client, "sample-container", transaction_date, customer_id)

if __name__ == '__main__':
    transaction_date = '20240409'
    customer_id = '001'
    # customer_id_list = ['001', '002', '003', '004']
    
    # 事件循环检测与处理逻辑
    try:
        loop = asyncio.get_running_loop()
        if loop.is_running():
            # 已有活跃循环,将任务提交到现有循环执行
            loop.create_task(main(transaction_date, customer_id))
    except RuntimeError:
        # 无活跃循环时,启动新循环执行任务
        asyncio.run(main(transaction_date, customer_id))

关键修改说明

  1. 检测位置:将循环检测逻辑放在if __name__ == '__main__':代码块中,替换原有的asyncio.run()调用。
  2. 分支处理:
    • 捕获RuntimeError:当无活跃循环时抛出此异常,此时用asyncio.run()启动任务。
    • 存在活跃循环时:调用loop.create_task()将异步任务加入现有循环,避免重复创建循环。
  3. Azure Function适配提示:如果是在Azure Function的异步触发器中,无需手动处理循环,直接将函数入口改为异步并返回await main(...)即可,示例:
# Azure Function异步触发器示例
async def main(req: func.HttpRequest) -> func.HttpResponse:
    transaction_date = req.params.get('transaction_date')
    customer_id = req.params.get('customer_id')
    await main(transaction_date, customer_id)
    return func.HttpResponse("下载完成", status_code=200)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 22:36:07