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

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

问题原因与解决办法

核心原因

  1. CELERY_TASK_ACKS_LATE = True的副作用:该配置让Worker在任务执行完成后才向Redis Broker发送确认信号。如果Worker在执行过程中意外重启、崩溃,或者任务执行时间过长,Broker会判定任务未完成,重新分发给其他Worker,导致重复执行。
  2. 任务无唯一标识:send_emails任务没有自定义唯一ID,Celery自动生成的ID可能被重复推送,加上Broker的重发机制,就会出现同一个任务被多次执行的情况。
  3. 多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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 05:44:52