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

Django+DRF+Celery:按数据库配置为不同分支调度定时任务

解决AppRegistryNotReady错误

直接在celery.py顶层导入Django模型会触发该错误,因为Celery启动时Django的应用注册表尚未加载完成。推荐两种解决方式:

  • 在任务函数内部延迟导入模型,而非文件顶部导入
  • 若需在Celery配置阶段访问模型,先手动初始化Django:
# celery.py 中
import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', '你的项目.settings')

# 手动初始化Django
import django
django.setup()

# 此时可安全导入模型
from 你的应用.models import WorkingDaysPolicy

实现分支专属定时任务

要实现每个分支在absence_Starts_at时间、仅工作日执行任务,推荐两种可行方案:

方案1:自定义Celery Beat调度器(动态生成定时任务)

通过自定义调度器从数据库读取配置,动态生成对应分支的定时任务条目。

  1. 创建自定义调度器类:
# 你的应用/celery_scheduler.py
from celery.beat import Scheduler, ScheduleEntry
from celery.utils.log import get_logger
from django.utils.timezone import now
from 你的应用.models import WorkingDaysPolicy
import datetime

logger = get_logger(__name__)

class BranchWorkingPolicyScheduler(Scheduler):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self._last_refresh = now()
        self._schedule = self._generate_schedule()

    def _generate_schedule(self):
        schedule = {}
        # 遍历所有有效分支配置
        for policy in WorkingDaysPolicy.objects.filter(branch__isnull=False, absence_Starts_at__isnull=False):
            branch_id = policy.branch.id
            check_time = policy.absence_Starts_at
            # 生成cron规则:匹配指定时间,排除休息日
            cron_minute = check_time.minute
            cron_hour = check_time.hour
            allowed_weekdays = [str(d) for d in range(7) if str(d) not in policy.weekend_days]
            cron_weekday = ','.join(allowed_weekdays) if allowed_weekdays else '*'

            # 构建定时任务
            schedule_key = f'absence_check_branch_{branch_id}'
            schedule[schedule_key] = ScheduleEntry(
                name=schedule_key,
                task='你的应用.tasks.check_branch_absences',
                schedule=f'{cron_minute} {cron_hour} * * {cron_weekday}',
                args=(branch_id,),
                last_run_at=now()
            )
        return schedule

    def get_schedule(self):
        # 每小时刷新一次调度,适配配置变更
        if (now() - self._last_refresh).total_seconds() > 3600:
            self._schedule = self._generate_schedule()
            self._last_refresh = now()
        return self._schedule
  1. 在celery.py中配置使用自定义调度器:
# celery.py
app.conf.beat_scheduler = '你的应用.celery_scheduler.BranchWorkingPolicyScheduler'
  1. 编写实际执行的任务:
# 你的应用/tasks.py
from celery import shared_task
from django.utils.timezone import now

@shared_task
def check_branch_absences(branch_id):
    # 延迟导入模型,避免启动时的AppRegistry问题
    from 你的应用.models import Branch, Employee
    try:
        branch = Branch.objects.get(id=branch_id)
        policy = branch.branch_working_days.first()
        # 额外校验当天是否为工作日(防止配置变更后调度未及时刷新)
        current_weekday = now().weekday()
        if str(current_weekday) not in policy.weekend_days:
            # 此处编写你的缺勤判定业务逻辑
            print(f"正在检查分支{branch.name}的缺勤情况:{now()}")
    except (Branch.DoesNotExist, WorkingDaysPolicy.DoesNotExist):
        # 分支或配置不存在,直接跳过
        pass

方案2:周期性调度器任务(更简单易维护)

设置一个高频运行的调度任务,检查当前时间是否符合分支配置,再触发具体任务。

  1. 编写调度与执行任务:
# 你的应用/tasks.py
from celery import shared_task
from django.utils.timezone import now
import datetime

@shared_task
def schedule_absence_checks():
    # 延迟导入模型
    from 你的应用.models import WorkingDaysPolicy
    current_time = now().time()
    current_weekday = now().weekday()  # 0=周一,对应模型中的WeekendDays值

    # 遍历所有有效配置
    for policy in WorkingDaysPolicy.objects.filter(branch__isnull=False, absence_Starts_at__isnull=False):
        # 匹配时间(允许±1分钟误差)
        time_diff = abs(
            datetime.timedelta(hours=current_time.hour, minutes=current_time.minute) -
            datetime.timedelta(hours=policy.absence_Starts_at.hour, minutes=policy.absence_Starts_at.minute)
        ).total_seconds()
        if time_diff <= 60 and str(current_weekday) not in policy.weekend_days:
            # 触发分支专属检查任务
            check_branch_absences.delay(policy.branch.id)

@shared_task
def check_branch_absences(branch_id):
    # 具体缺勤检查逻辑,同方案1
    from 你的应用.models import Branch
    try:
        branch = Branch.objects.get(id=branch_id)
        print(f"正在处理分支{branch.name}的缺勤检查")
        # 业务代码...
    except Branch.DoesNotExist:
        pass
  1. 在celery.py中配置调度任务每分钟运行:
# celery.py
app.conf.beat_schedule = {
    'schedule-absence-checks-every-minute': {
        'task': '你的应用.tasks.schedule_absence_checks',
        'schedule': 60.0,  # 每分钟执行一次
    },
}

注意事项

  • 方案1中配置变更后需等待调度器刷新周期(如1小时)生效,若需实时生效,可在WorkingDaysPolicy保存信号中触发调度器刷新;方案2则会在下一分钟自动适配。
  • 务必在任务内部延迟导入模型,规避启动阶段的AppRegistryNotReady错误。
  • 可给check_branch_absences任务添加幂等性处理,防止重复执行。

内容的提问来源于stack exchange,提问作者Mahmoud Said

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 12:25:31