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

如何在同步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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 14:10:35