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

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策略

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"
    }
  ]
}

排查方向建议

  1. Worker权限验证:确认Beanstalk Worker实例使用的IAM角色是否包含在SQS访问策略中。当前策略仅授权了rvm-jenkins-service-roll角色和rvm-jenkins-user用户,若Worker使用的是其他角色(比如Beanstalk默认服务角色),则会缺少sqs:ReceiveMessage等权限,导致无法消费消息。
  2. 队列账户ID匹配:SQS队列地址中的账户ID是718854804674,但访问策略中的ARN使用的是718854674,存在位数不一致,需确认是否为笔误,这会导致权限不匹配。
  3. Worker进程状态:登录Beanstalk实例执行ps aux | grep celery,确认Worker进程是否正常运行;查看Beanstalk控制台的进程日志,排查是否有启动失败或崩溃记录。
  4. 消息序列化问题:在Worker日志中查找是否有任务反序列化错误,可尝试在任务中添加logger.info("Task executed")等日志输出,验证任务是否能被触发。
  5. SQS消息状态检查:在AWS控制台查看队列的ApproximateNumberOfMessagesNotVisible指标,确认是否有消息被Worker获取但未确认,导致长期处于不可见状态。
  6. Broker连接测试:在Worker实例上执行celery -A rvm inspect ping,验证Worker是否能正常连接到SQS Broker;执行celery -A rvm list queues,确认Worker能识别目标队列。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 05:40:55