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

Celery结合Amazon SQS问题:消息滞留在可用队列未被Worker消费

Django + Celery + Amazon SQS 异步任务无法消费问题排查

问题描述

tasks.py中的verify_mail函数直接调用正常,但使用delay()异步调用时,消息已成功发送至SQS(控制台「Messages Available」可见),但Celery Worker完全不接收消息,消息始终未进入「Messages in Flight」区域。

相关代码与配置

1. views.py 任务调用代码

verify_mail.delay(abcdef@gmail.com)

2. tasks.py 任务代码

from celery import shared_task
from time import sleep
from django.shortcuts import render, redirect, HttpResponse
import boto3

@shared_task
def verify_mail(new_email):
    ses = boto3.client('ses')
    response = ses.verify_email_identity(
        EmailAddress = new_email
    )
    return None

3. Celery Worker启动命令

celery -A myProject worker -l INFO --without-gossip --without-mingle --without-heartbeat -Ofair --pool=solo

4. celery.py 配置

import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'myProject.settings')

app = Celery('myProject')
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()

@app.task(bind=True, ignore_result=True)
def debug_task(self):
    print(f'Request: {self.request!r}')

5. settings.py Celery相关配置

CELERY_BROKER_URL = "sqs://{aws_access_key}:{aws_secret_key}@".format(
    aws_access_key=AWS_ACCESS_KEY_ID, aws_secret_key=AWS_SECRET_ACCESS_KEY
)

CELERY_ACCEPT_CONTENT = ['application/json']
CELERY_RESULT_SERIALIZER = 'json'
CELERY_TASK_SERIALIZER = 'json'
CELERY_BROKER_TRANSPORT_OPTIONS = {
    'region': 'eu-central-1',
}
CELERY_RESULT_BACKEND = None
CELERY_ENABLE_REMOTE_CONTROL = False
CELERY_SEND_EVENTS = False

排查方向与解决方案

1. 队列名称匹配检查

Celery默认使用celery作为队列名称,确认SQS控制台中存在对应区域(eu-central-1)的celery队列;若使用自定义队列,需在settings中显式指定:

CELERY_TASK_DEFAULT_QUEUE = '你的自定义队列名'

2. AWS权限验证

确保Worker所用AWS账号拥有以下SQS权限:

  • sqs:ReceiveMessage
  • sqs:DeleteMessage
  • sqs:GetQueueAttributes
  • sqs:ChangeMessageVisibility
    同时确认账号有权限访问eu-central-1区域的SQS资源。

3. Worker显式指定监听队列

启动Worker时明确指定监听队列,避免默认队列不匹配:

celery -A myProject worker -l INFO --without-gossip --without-mingle --without-heartbeat -Ofair --pool=solo -Q celery

若使用自定义队列,将celery替换为你的队列名称。

4. SQS可见性超时设置

检查队列的可见性超时时间,需大于任务最大执行时长,防止消息未处理完成就被重新放回队列。

5. 任务序列化验证

在任务中添加日志,确认Worker是否能接收并解析参数:

import logging
logger = logging.getLogger(__name__)

@shared_task
def verify_mail(new_email):
    logger.info(f"接收到待验证邮箱: {new_email}")
    # 原有业务代码

启动Worker后查看日志是否有对应输出,排查序列化/反序列化问题。

6. 版本兼容性检查

确认Celery与Django、boto3版本兼容,参考官方文档核对SQS配置格式是否符合当前Celery版本要求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 06:05:27