Celery任务函数不能直接await,因其worker基于同步模型且不原生支持async/await;正确做法是用asyncio.to_thread()或预置事件循环复用,或升级至Celery 5.3+并启用--pool=threads。

asyncio任务不能直接交给Celery执行
Celery默认运行在同步线程中,不支持await语法或async def函数。如果你把一个async def my_task()直接注册为Celery任务,调用时会报RuntimeWarning: coroutine 'my_task' was never awaited,或者更隐蔽地卡住、返回None——它根本没运行协程体。
根本原因:Celery worker启动的是普通Python线程,事件循环(event loop)不在其上下文中;而asyncio.run()每次调用都会新建并关闭loop,无法复用,也不适合高并发场景。
- 别写
@app.task async def task(): ...——语法错误,Celery不识别 - 别在task函数里直接
await asyncio.sleep(1)——会阻塞整个worker线程 - 不要试图在
celery -A tasks worker启动后手动asyncio.set_event_loop()——loop绑定失败或被覆盖
用asyncio.to_thread()包装异步逻辑(Python 3.9+)
最轻量、最安全的做法:把异步函数包裹进同步接口,再交给Celery。Python 3.9引入的asyncio.to_thread()能在线程池中安全运行协程并等待结果,避免阻塞主线程。
关键点是“运行协程”不是“定义协程”——你要在同步函数体内显式await它:
立即学习“Python免费学习笔记(深入)”;
import asyncio
from celery import Celery
<p>app = Celery('tasks', broker='redis://localhost')</p><p>@app.task
def run_async_job():</p><h1>同步入口,但内部驱动异步逻辑</h1><pre class="brush:php;toolbar:false;">return asyncio.run(run_my_async_logic())async def run_my_async_logic(): await asyncio.sleep(1) return {"status": "done", "data": 42}
⚠️注意:asyncio.run()虽可用,但每次新建loop开销大,不适合高频任务。生产环境建议改用asyncio.get_event_loop().run_until_complete(),前提是确保loop已存在且未关闭(见下一条)。
在Celery worker启动时预置并复用event loop
Celery提供worker_process_init信号,可在每个worker子进程初始化时运行一次代码。这是注入和复用loop的最佳时机。
- 必须使用
asyncio.new_event_loop()+asyncio.set_event_loop(),不能依赖默认loop(可能已被关闭) - loop必须设为全局可访问变量(如模块级
_loop),供task函数调用 - 不要在task里调用
asyncio.run(),改用_loop.run_until_complete(coro)
import asyncio
from celery import Celery
<p>app = Celery('tasks', broker='redis://localhost')
_loop = None</p><p>@app.task
def fetch_data_async():
global _loop
coro = do_http_fetch()
return _loop.run_until_complete(coro)</p><p>async def do_http_fetch():
import aiohttp
async with aiohttp.ClientSession() as session:
async with session.get('<a href="https://www.php.cn/link/5f69e19efaba426d62faeab93c308f5c">https://www.php.cn/link/5f69e19efaba426d62faeab93c308f5c</a>') as resp:
return await resp.json()</p><p>@app.task
def worker_init():
pass</p><h1>在worker启动时设置loop</h1><p>from celery.signals import worker_process_init</p><p>@worker_process_init.connect
def init_worker(**kwargs):
global _loop
_loop = asyncio.new_event_loop()
asyncio.set_event_loop(_loop)
用Celery 5.3+ 的native asyncio support(推荐但需升级)
Celery 5.3起正式支持async def任务,但有硬性前提:worker必须用--pool=threads或--pool=gevent,且不能用prefork(默认)。否则仍会报错。
启用方式很简单,但容易漏掉配置:
- 任务函数必须是
async def,且装饰器写法为@app.task(bind=True, autoretry_for=(Exception,), retry_kwargs={'max_retries': 3}) - 启动worker时加
celery -A tasks worker --pool=threads -c 4(-c指定并发数,对应线程数) - 所有await调用必须在该协程内完成,不能跨task传递未完成的Future
这种模式下,Celery会为每个task分配独立的协程上下文,无需手动管理loop,但牺牲了prefork模型的内存隔离优势——所有task共享同一进程的内存空间,异常可能互相影响。
真正难处理的不是怎么写async task,而是当你的异步任务里要调用另一个也带async的第三方库(比如aioredis或httpx.AsyncClient),而这个库又依赖特定loop策略或SSL上下文时——这时候得仔细看它的__init__是否支持传入loop,以及是否允许在非主线程loop中初始化。


















