如何在同步Celery任务中获取协程对象的返回值?
解决Celery同步任务中获取异步函数返回值的问题
要在Celery的@shared_task同步函数里拿到异步函数的返回值,核心是运行协程对象并获取其执行结果,可以通过Python标准库的asyncio模块实现:
步骤1:导入asyncio模块
在任务文件顶部添加:
import asyncio
步骤2:修改同步任务中调用异步函数的逻辑
把原来直接赋值协程对象的代码,改成用asyncio.run()执行异步函数并获取结果:
修改后的完整任务代码:
import asyncio from celery import shared_task from django.utils import timezone # 导入你的模型和异步函数 from .models import Club, ClubLoadModel from .your_module import club_count # 替换成异步函数所在的实际模块 @shared_task def get_club_load(): all_clubs = Club.objects.all() for club in all_clubs: if club.monitoring_club_load: # 简化布尔值判断,无需写==True try: # 运行异步函数,获取返回的整数值 club_load = asyncio.run(club_count(club.id)) except Exception as e: # 捕获异步执行异常,避免单个俱乐部的错误导致整个任务崩溃 print(f"获取俱乐部{club.id}人数失败: {str(e)}") continue # 跳过当前俱乐部,继续处理下一个 club_capacity = club.club_load procent_club_load = None if club_capacity and club_capacity > 0: # 避免除以0的运算错误 procent_club_load = int((club_load / club_capacity)*100) club_load_record = ClubLoadModel( club_id=club.id, created_at=timezone.now(), weekday=timezone.localtime().weekday() + 1, club_load=club_load, procent_club_load=procent_club_load ) club_load_record.save()
关键说明
asyncio.run()会创建新的事件循环,运行传入的协程直到完成并返回结果,这是Python 3.7+官方推荐的协程执行方式。- 添加
try-except块是为了处理异步函数可能抛出的异常(比如数据库查询失败、外部接口调用错误等),保证任务的稳定性。 - 新增
club_capacity > 0的判断,避免出现除以0的运算错误,提升代码健壮性。
内容的提问来源于stack exchange,提问作者JackOMan69
相关产品推荐
相关产品推荐

