Django Celery任务重复执行:send_emails多次发送重复邮件
Django + Celery 重复发送邮件问题排查与解决
问题描述
在Django项目中开发邮件发送应用,通过Celery Beat每日00:30调度执行schedule_emails函数,该函数运行正常,但send_emails任务会被重复执行,导致特定客户收到多份相同邮件。查看Celery Worker日志发现,同一个task_id被多次接收并执行成功。当前使用Supervisor后台托管Celery Worker和Celery Beat服务,相关配置及日志如下:
配置文件与日志
settings.py 配置
# Celery Settings CELERY_BROKER_URL = 'redis://127.0.0.1:6379' CELERY_ACCEPT_CONTENT = ['application/json'] CELERY_RESULT_SERIALIZER = 'json' CELERY_TASK_SERIALIZER = 'json' CELERY_TIMEZONE = 'Asia/Kolkata' CELERY_TASK_ACKS_LATE = True CELERY_RESULT_BACKEND = 'django-db' CELERY_BEAT_SCHEDULE_FILENAME = '.celery/beat-schedule' CELERYD_LOG_FILE = '.celery/celery.log' CELERYBEAT_LOG_FILE = '.celery/celerybeat.log' CELERY_BEAT_SCHEDULE = { 'schedule_emails': { 'task': 'myapp.tasks.schedule_emails', 'schedule': crontab(hour=0, minute=30), }, }
celery.py 配置
from __future__ import absolute_import, unicode_literals import os from celery import Celery from django.conf import settings from decouple import config if config('ENVIRONMENT') == 'development': os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'MyProject.settings.development') else: os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'MyProject.settings.production') app = Celery('MyProject') app.conf.enable_utc = False app.conf.update(timezone='Asia/Kolkata') app.config_from_object(settings, namespace='CELERY') app.autodiscover_tasks() @app.task(bind=True) def debug_task(self): print(f'Request : {self.request!r}')
tasks.py 代码
from celery import shared_task @shared_task def send_emails(client_email): # 邮件发送代码 @shared_task def schedule_emails(): client_data = [] # 包含客户邮箱和发送时间的字典列表 for data in client_data : send_emails.apply_async(args=[data.get('email')],eta=data.get('time'))
Supervisor 配置
[program:celery_worker] command=path-to-enviroment/bin/celery -A MyProject worker --loglevel=info directory=project-path user=admin autostart=true autorestart=true redirect_stderr=true stdout_logfile=project-path/.celery/celery_worker.log [program:celery_beat] command=path-to-enviroment/bin/celery -A MyProject beat --loglevel=info directory=project-path user=username autostart=true autorestart=true redirect_stderr=true stdout_logfile=project-path/MyProject/.celery/celery_beat.log
celery_beat.log 日志
[2023-07-01 19:23:51,519: INFO/MainProcess] beat: Starting... [2023-07-02 00:30:00,056: INFO/MainProcess] Scheduler: Sending due task schedule_emails (auto_email.tasks.schedule_emails) [2023-07-02 04:00:00,091: INFO/MainProcess] Scheduler: Sending due task celery.backend_cleanup (celery.backend_cleanup) [2023-07-03 00:30:00,091: INFO/MainProcess] Scheduler: Sending due task schedule_emails (auto_email.tasks.schedule_emails) [2023-07-03 04:00:00,092: INFO/MainProcess] Scheduler: Sending due task celery.backend_cleanup (celery.backend_cleanup) [2023-07-04 00:30:00,091: INFO/MainProcess] Scheduler: Sending due task schedule_emails (auto_email.tasks.schedule_emails) [2023-07-04 04:00:00,090: INFO/MainProcess] Scheduler: Sending due task celery.backend_cleanup (celery.backend_cleanup)
celery_worker.log 日志
[2023-07-04 00:30:37,760: INFO/MainProcess] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] received [2023-07-04 01:32:16,692: INFO/MainProcess] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] received [2023-07-04 01:35:23,098: INFO/MainProcess] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] received [2023-07-04 02:36:54,644: INFO/MainProcess] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] received [2023-07-04 03:38:36,868: INFO/MainProcess] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] received [2023-07-04 03:38:44,074: INFO/ForkPoolWorker-1] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] succeeded in 3.069899449998047s: None [2023-07-04 03:38:44,081: INFO/ForkPoolWorker-1] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] succeeded in 0.004782014002557844s: None [2023-07-04 03:38:44,190: INFO/ForkPoolWorker-2] Task auto_email.tasks.send_emails[66c8ef59-34e2-48be-85e4-ed14cd9b56cf] succeeded in 3.1878137569874525s: None
问题原因与解决办法
核心原因
CELERY_TASK_ACKS_LATE = True的副作用:该配置让Worker在任务执行完成后才向Redis Broker发送确认信号。如果Worker在执行过程中意外重启、崩溃,或者任务执行时间过长,Broker会判定任务未完成,重新分发给其他Worker,导致重复执行。- 任务无唯一标识:
send_emails任务没有自定义唯一ID,Celery自动生成的ID可能被重复推送,加上Broker的重发机制,就会出现同一个任务被多次执行的情况。 - 多Worker实例冲突:如果Supervisor误启动了多个Worker进程,或者手动启动的Worker未关闭,多个Worker会同时监听任务队列,抢同一个任务执行。
具体解决步骤
1. 调整ACK配置(优先推荐)
如果业务不需要延迟确认,直接关闭CELERY_TASK_ACKS_LATE:
# settings.py CELERY_TASK_ACKS_LATE = False # 默认值就是False,可直接删除原配置
这样Worker收到任务后立即发送ACK,Broker不会重复推送该任务。
2. 给任务添加唯一ID
如果必须保留ACKS_LATE,可以为每个send_emails任务生成唯一ID,基于客户邮箱和发送时间,避免重复任务:
# tasks.py from celery.utils import uuid from django.utils import timezone @shared_task def schedule_emails(): client_data = [] # 你的客户数据列表 for data in client_data : # 生成唯一task_id task_id = f"send_email_{data.get('email').replace('@', '_')}_{data.get('time').timestamp()}" send_emails.apply_async( args=[data.get('email')], eta=data.get('time'), task_id=task_id )
3. 确保Worker实例唯一
检查当前运行的Worker进程:
ps aux | grep celery
如果发现多余的Worker,手动杀掉,然后调整Supervisor配置,指定Worker并发数(比如--concurrency=2),避免自动生成过多进程:
# Supervisor celery_worker配置 command=path-to-enviroment/bin/celery -A MyProject worker --loglevel=info --concurrency=2
4. 给任务设置过期时间
防止任务被无限重发,添加expires参数:
send_emails.apply_async( args=[data.get('email')], eta=data.get('time'), expires=3600 # 任务1小时后过期,不再被执行 )
5. 实现幂等任务(终极保障)
在send_emails函数中添加重复发送校验,比如记录已发送的邮件日志,避免重复执行:
# tasks.py from django.utils import timezone from myapp.models import EmailLog # 需提前创建该模型,包含email、sent_at字段 @shared_task def send_emails(client_email): # 检查今日是否已发送过该邮件 today = timezone.now().date() if EmailLog.objects.filter(email=client_email, sent_at__date=today).exists(): return "邮件已发送,跳过执行" # 执行邮件发送代码 # ... 发送逻辑 ... # 发送成功后记录日志 EmailLog.objects.create(email=client_email, sent_at=timezone.now()) return "邮件发送成功"
内容的提问来源于stack exchange,提问作者Manoj Kamble
相关产品推荐
相关产品推荐

