Django应用通过Celery发任务至SQS,Worker未执行求助
问题:Celery任务发送到SQS队列后无法被Worker执行
我在AWS上创建了名为celery-celery的SQS队列,调用API时通过Celery向该队列发送任务。消息已成功到达队列(队列地址:https://sqs.us-east-1.amazonaws.com/718854804674/celery-celery),每次调用API都会新增可用消息,但这些消息始终停留在队列中,不会被Worker执行。
部署环境与API配置
Django应用部署在AWS Beanstalk上,API路径为/api/print-hello-world/,URL配置如下:
path('print-hello-world/', print_hello_world_task, name='print_hello_world')
视图代码
def print_hello_world_task(request): print("Received request to start timer task") logger.info("Received request to start timer task") task = print_hello_world.delay() # 异步触发Celery任务 print(f"Task started with ID: {task.id}") return JsonResponse({'status': 'Task started', 'task_id': task.id})
Celery任务代码(tasks.py)
@shared_task(queue='celery-celery') def print_hello_world(): return "Executed"
Procfile配置
web: gunicorn --workers 1 rvm.wsgi:application worker: celery -A rvm worker -l debug --loglevel=DEBUG --queues=celery-celery --logfile=/var/log/celery/celery.log -E
Django settings.py中的Celery配置
AWS_ACCESS_KEY = os.getenv('AWS_ACCESS_KEY') AWS_SECRET_KEY = os.getenv('AWS_SECRET_KEY') CELERY_RESULT_BACKEND = 'django-db' CELERY_CACHE_BACKEND = 'django-cache' CELERY_ACCEPT_CONTENT = ['json'] CELERY_TASK_SERIALIZER = 'json' CELERY_TASK_DEFAULT_QUEUE = 'celery-celery' CELERY_RESULT_SERIALIZER = 'json' CELERY_TIMEZONE = 'UTC' CELERY_LOG_FILE = '/var/log/celery/celery.log' CELERY_BROKER_URL = 'sqs://{0}:{1}@'.format(AWS_ACCESS_KEY, AWS_SECRET_KEY) logger.info(f'CELERY_BROKER_URL: {CELERY_BROKER_URL}') CELERY_BROKER_TRANSPORT_OPTIONS = { 'region': 'us-east-1', # AWS区域 'polling_interval': 1, # 消息轮询间隔 'visibility_timeout': 3600, # 消息可见性超时 'broker_connection_retry_on_startup': True, # 启动时重试连接 'use_ssl': True, # 启用SSL 'sqs': { 'signature_version': 'v4', # 使用SigV4签名 } }
celery.py配置
import os from celery import Celery import warnings from celery.utils.log import get_logger from kombu import Queue, Exchange # 设置Celery默认的Django配置模块 os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'rvm.settings') app = Celery('rvm') app.conf.task_default_queue = 'celery-celery' # 从Django配置加载Celery设置,前缀为CELERY_ app.config_from_object('django.conf:settings', namespace='CELERY') # 自动加载所有已注册Django应用的任务模块 app.autodiscover_tasks() app.conf.update( broker_connection_retry_on_startup=True, ) # 配置任务队列,确保Worker监听正确队列 app.conf.task_queues = ( Queue('celery-celery', Exchange('celery-celery'), routing_key='celery-celery'), ) # 设置日志 logger = get_logger(__name__) # 忽略PendingDeprecationWarning警告 warnings.filterwarnings('ignore', category=PendingDeprecationWarning)
日志确认消息已到达队列
celery.log片段显示队列中有消息,但未被处理:
[2024-08-29 04:19:43,559: DEBUG/MainProcess] Response body: b'<?xml version="1.0"?><GetQueueAttributesResponse xmlns="http://queue.amazonaws.com/doc/2012-11-05/"><GetQueueAttributesResult><Attribute><Name>ApproximateNumberOfMessages</Name><Value>3</Value></Attribute></GetQueueAttributesResult><ResponseMetadata><RequestId>7c6389e1-ea52-593a-a589-d6a8134d2a9c</RequestId></ResponseMetadata></GetQueueAttributesResponse>' [2024-08-29 04:19:43,560: DEBUG/MainProcess] Event needs-retry.sqs.GetQueueAttributes: calling handler <botocore.retryhandler.RetryHandler object at 0x7fb8ff92ebd0> [2024-08-29 04:19:43,560: DEBUG/MainProcess] No retry needed. [2024-08-29 04:19:43,563: DEBUG/MainProcess] basic.qos: prefetch_count->8
IAM与SQS访问策略
IAM策略
SQS访问策略
{ "Version": "2012-10-17", "Statement": [ { "Sid": "__owner_statement", "Effect": "Allow", "Principal": { "AWS": "arn:aws:iam::718854674:root" }, "Action": "SQS:*", "Resource": "arn:aws:sqs:us-east-1:718854674:*" }, { "Sid": "AllowRvmJenkinsServiceRoleAccess_Part1", "Effect": "Allow", "Principal": { "AWS": "arn:aws:iam::718854674:role/service-role/rvm-jenkins-service-roll" }, "Action": [ "sqs:GetQueueAttributes", "sqs:SendMessage", "sqs:ReceiveMessage", "sqs:DeleteMessage" ], "Resource": "arn:aws:sqs:us-east-1:718854674:celery-celery" }, { "Sid": "AllowRvmJenkinsServiceRoleAccess_Part2", "Effect": "Allow", "Principal": { "AWS": "arn:aws:iam::718854674:role/service-role/rvm-jenkins-service-roll" }, "Action": [ "sqs:ChangeMessageVisibility", "sqs:GetQueueUrl", "sqs:ListQueues" ], "Resource": "arn:aws:sqs:us-east-1:718854674:celery-celery" }, { "Sid": "AllowRvmJenkinsUserAccess_Part1", "Effect": "Allow", "Principal": { "AWS": "arn:aws:iam::718854674:user/rvm-jenkins-user" }, "Action": [ "sqs:GetQueueAttributes", "sqs:SendMessage", "sqs:ReceiveMessage", "sqs:DeleteMessage" ], "Resource": "arn:aws:sqs:us-east-1:718854674:celery-celery" }, { "Sid": "AllowRvmJenkinsUserAccess_Part2", "Effect": "Allow", "Principal": { "AWS": "arn:aws:iam::718854674:user/rvm-jenkins-user" }, "Action": [ "sqs:ChangeMessageVisibility", "sqs:GetQueueUrl", "sqs:ListQueues" ], "Resource": "arn:aws:sqs:us-east-1:718854674:celery-celery" } ] }
排查方向建议
- Worker权限验证:确认Beanstalk Worker实例使用的IAM角色是否包含在SQS访问策略中。当前策略仅授权了
rvm-jenkins-service-roll角色和rvm-jenkins-user用户,若Worker使用的是其他角色(比如Beanstalk默认服务角色),则会缺少sqs:ReceiveMessage等权限,导致无法消费消息。 - 队列账户ID匹配:SQS队列地址中的账户ID是
718854804674,但访问策略中的ARN使用的是718854674,存在位数不一致,需确认是否为笔误,这会导致权限不匹配。 - Worker进程状态:登录Beanstalk实例执行
ps aux | grep celery,确认Worker进程是否正常运行;查看Beanstalk控制台的进程日志,排查是否有启动失败或崩溃记录。 - 消息序列化问题:在Worker日志中查找是否有任务反序列化错误,可尝试在任务中添加
logger.info("Task executed")等日志输出,验证任务是否能被触发。 - SQS消息状态检查:在AWS控制台查看队列的
ApproximateNumberOfMessagesNotVisible指标,确认是否有消息被Worker获取但未确认,导致长期处于不可见状态。 - Broker连接测试:在Worker实例上执行
celery -A rvm inspect ping,验证Worker是否能正常连接到SQS Broker;执行celery -A rvm list queues,确认Worker能识别目标队列。
内容的提问来源于stack exchange,提问作者Daniel Johnson
相关产品推荐
相关产品推荐

