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:ReceiveMessagesqs:DeleteMessagesqs:GetQueueAttributessqs: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
相关产品推荐
相关产品推荐

