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

为何在同一Celery任务函数中两次调用asyncio.run会报错?

问题分析与解决:Celery任务中asyncio.run重复调用报错"Event loop is closed"

报错原因

  • asyncio.run() 的设计逻辑是每次调用都会新建事件循环,执行完毕后自动关闭该循环。try块里第一次调用asyncio.run(execute_match_search())时,循环已经被关闭。到finally块再次调用asyncio.run(emit_notification())时,它会尝试复用已关闭的循环(或者在当前进程环境中,默认循环已失效),直接触发"Event loop is closed"错误。
  • Celery的worker进程默认是同步运行的,没有内置的异步事件循环管理机制,重复调用asyncio.run()会导致循环资源的复用/释放冲突。

解决办法

1. 手动复用单个事件循环

创建一个事件循环实例,在try-finally流程中复用它,手动控制循环的启动和关闭:

import asyncio
from celery import Celery

app = Celery('tasks', broker='pyamqp://guest@localhost//')

@app.task
def background_task():
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        loop.run_until_complete(execute_match_search())
    finally:
        loop.run_until_complete(emit_notification())
        loop.close()

2. 合并异步任务,单次调用asyncio.run()

如果两个异步任务不需要严格分步执行,可以把它们打包成一个异步函数,只调用一次asyncio.run():

import asyncio
from celery import Celery

app = Celery('tasks', broker='pyamqp://guest@localhost//')

async def full_task_flow():
    await execute_match_search()
    await emit_notification()

@app.task
def background_task():
    asyncio.run(full_task_flow())

这种方式从根源上避免了重复创建/关闭循环的问题,也是最推荐的方案。

3. 使用Celery异步扩展池

如果你的Celery任务大量依赖异步操作,可以用celery-aio-pool让Celery直接管理异步事件循环:
首先安装扩展:

pip install celery-aio-pool

启动worker时指定异步池:

celery -A tasks worker -P aio -l info

之后可以直接定义异步Celery任务,无需手动处理asyncio.run():

@app.task(acks_late=True)
async def background_task():
    try:
        await execute_match_search()
    finally:
        await emit_notification()

4. 捕获异常后重建循环(临时方案)

如果必须在finally中单独调用异步函数,可以捕获循环关闭异常后重建循环:

import asyncio
from celery import Celery

app = Celery('tasks', broker='pyamqp://guest@localhost//')

@app.task
def background_task():
    try:
        asyncio.run(execute_match_search())
    finally:
        try:
            asyncio.run(emit_notification())
        except RuntimeError as e:
            if "Event loop is closed" in str(e):
                loop = asyncio.new_event_loop()
                asyncio.set_event_loop(loop)
                loop.run_until_complete(emit_notification())
                loop.close()

这种方式属于临时 workaround,不建议作为长期方案使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 06:41:17