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调度器(动态生成定时任务)
通过自定义调度器从数据库读取配置,动态生成对应分支的定时任务条目。
- 创建自定义调度器类:
# 你的应用/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
- 在
celery.py中配置使用自定义调度器:
# celery.py app.conf.beat_scheduler = '你的应用.celery_scheduler.BranchWorkingPolicyScheduler'
- 编写实际执行的任务:
# 你的应用/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:周期性调度器任务(更简单易维护)
设置一个高频运行的调度任务,检查当前时间是否符合分支配置,再触发具体任务。
- 编写调度与执行任务:
# 你的应用/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
- 在
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
相关产品推荐
相关产品推荐

