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

使用AsyncIO的Django API仍同步运行问题排查与解决

API调用异步函数仍同步阻塞,如何用AsyncIO实现立即返回?

问题描述

我创建了APIDummyWait5 API,核心目标有两个:

  1. API能在毫秒级返回响应
  2. 内部调用的get_country_details函数异步执行——该函数会先休眠随机生成的1-25秒wait_time,再调用第三方API获取国家信息并写入CSV。

但实际测试时,当wait_time为11秒,API要等11秒才返回,完全不符合毫秒级返回的预期。请问为什么代码还是同步运行?仅用AsyncIO该怎么解决?

相关代码

class APIDummyWait5(APIView):
    def get(self, request):
        print("Coming Here")
        worker = os.getpid()
        thread = threading.get_ident()

        country = request.GET.get('country')

        print(f"API wait 5 - for worker - {worker} at Thread {thread} starting at: {datetime.datetime.now()}")

        wait_time = rd.randint(1, 25)
        response = {'wait_time': wait_time, 'country': country}

        asyncio.run(main(wait_time=wait_time, country=country))

        print(f"API Wait 5 - for worker - {worker} at Thread {thread} ending at: {datetime.datetime.now()}")

        return Response({
            "data": response
        })
async def main(wait_time, country):
    result = await get_country_details(wait_time, country)
async def get_country_details(wait_time, country):
    worker = os.getpid()

    print(f"API Quick @ {wait_time} - for worker - {worker} starting at: {datetime.datetime.now()}")

    await asyncio.sleep(wait_time)  # Simulate waiting time

    url = f"https://restcountries.com/v3.1/name/{country}"
    
    async with aiohttp.ClientSession() as session:
        async with session.get(url) as response:
            if response.status == 200:
                data = await response.json()
                
                file = pd.read_csv('country_data.csv')
                name_list = list(file['name'])
                capital_list = list(file['capital'])

                country_data = data[0]
                country = country_data['name']['common']
                capital = country_data['capital'][0]

                name_list.append(country)
                capital_list.append(capital)

                new_data = pd.DataFrame({'name': name_list, 'capital': capital_list})
                new_data.to_csv('country_data.csv', index=False)

                print(f"Asynchronously Ran Function for {wait_time} seconds.")
                return 1
            else:
                print(f"Error: Unable to fetch details for {country}")
                return None

为什么会同步阻塞?

核心问题在于**asyncio.run()是阻塞式调用**:

  • 你在同步的API视图函数(get方法是同步实现)里调用asyncio.run(main(...))时,这个方法会启动一个全新的AsyncIO事件循环,并且会等待整个异步任务链(main→get_country_details)完全执行完毕才会退出。
  • 也就是说,API的主线程会被asyncio.run()彻底卡住,直到休眠、第三方API调用、CSV写入全部完成,自然无法实现毫秒级返回。

另外还有隐藏问题:get_country_details里的Pandas读写CSV是同步阻塞IO操作,即使解决了调用方式的问题,这部分代码也会阻塞事件循环,影响其他异步任务的执行效率。


仅用AsyncIO的解决方案

要实现API立即返回,需要把异步任务丢进后台事件循环执行,而非等待它完成。具体步骤如下:

1. 将API视图改为异步视图

DRF支持异步视图,改成异步后可以复用已有的事件循环,避免手动创建循环的开销和问题:

class APIDummyWait5(APIView):
    async def get(self, request):  # 改为async方法
        print("Coming Here")
        worker = os.getpid()
        thread = threading.get_ident()

        country = request.GET.get('country')

        print(f"API wait 5 - for worker - {worker} at Thread {thread} starting at: {datetime.datetime.now()}")

        wait_time = rd.randint(1, 25)
        response = {'wait_time': wait_time, 'country': country}

        # 获取当前事件循环,提交任务到后台执行,不等待结果
        loop = asyncio.get_event_loop()
        loop.create_task(get_country_details(wait_time=wait_time, country=country))

        print(f"API Wait 5 - for worker - {worker} at Thread {thread} ending at: {datetime.datetime.now()}")

        return Response({
            "data": response
        })

2. 移除冗余的main函数

原来的main函数只是单纯等待get_country_details,没有额外逻辑,直接调用目标函数即可:

# 删掉这个无意义的中间层函数
# async def main(wait_time, country):
#     result = await get_country_details(wait_time, country)

3. 修复同步IO阻塞问题(可选但推荐)

get_country_details里的Pandas读写CSV是同步操作,会阻塞事件循环。可以用asyncio.to_thread()把这部分代码丢到线程池执行,避免卡主整个事件循环:

async def get_country_details(wait_time, country):
    worker = os.getpid()

    print(f"API Quick @ {wait_time} - for worker - {worker} starting at: {datetime.datetime.now()}")

    await asyncio.sleep(wait_time)  # 异步休眠,不会阻塞事件循环

    url = f"https://restcountries.com/v3.1/name/{country}"
    
    async with aiohttp.ClientSession() as session:
        async with session.get(url) as response:
            if response.status == 200:
                data = await response.json()
                
                # 把同步的CSV操作丢到线程池执行
                await asyncio.to_thread(write_country_to_csv, data)

                print(f"Asynchronously Ran Function for {wait_time} seconds.")
                return 1
            else:
                print(f"Error: Unable to fetch details for {country}")
                return None

# 把CSV操作抽成独立的同步函数
def write_country_to_csv(data):
    file = pd.read_csv('country_data.csv')
    name_list = list(file['name'])
    capital_list = list(file['capital'])

    country_data = data[0]
    country = country_data['name']['common']
    capital = country_data['capital'][0]

    name_list.append(country)
    capital_list.append(capital)

    new_data = pd.DataFrame({'name': name_list, 'capital': capital_list})
    new_data.to_csv('country_data.csv', index=False)

关键逻辑说明

  • 异步视图:DRF的异步视图会自动在AsyncIO事件循环中运行,无需手动创建循环,避免了循环嵌套的问题。
  • loop.create_task():将异步任务提交到当前事件循环的后台队列,立即返回,不会阻塞当前视图的执行,从而实现API毫秒级响应。
  • asyncio.to_thread():将同步IO操作转移到线程池执行,保证事件循环不会被阻塞,确保其他异步任务能正常调度。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 20:05:22