Celery使用SQS作为Broker时收到事件未触发consume任务怎么解决
问题原因
- Celery worker默认只识别符合Celery任务协议的消息,你直接向SQS发送普通JSON字符串,worker收到后无法匹配到对应的任务函数,因此不会触发
consume执行 - 其次你当前的
celery.py没有加载celeryconfig配置文件,自定义的队列、路由规则都没有生效
解决步骤
第一步:修复配置加载
在celery.py文件中,初始化Celery实例后添加配置加载逻辑:
from celery.utils.log import get_task_logger from celery import Celery app = Celery(__name__) # 新增这行,加载celeryconfig.py的配置 app.config_from_object('celeryconfig') logger = get_task_logger(__name__) @app.task(routing_key='consume', name="consume", bind=True, acks_late=True, ignore_result=True) def consume(self, msg): print('Message received') logger.info('Message received') # 对接收的消息做对应处理 # print('this is the new message', msg) return True
第二步:二选一选择适配方案
方案1:调整发送消息的格式(推荐,适配Celery原生逻辑)
Celery要求消息必须包含任务名、参数等字段,你需要将AWS CLI发送的消息体修改为符合Celery协议的格式,修改后的命令如下:
aws --endpoint-url http://localhost:9324 sqs send-message --queue-url http://localhost:9324/queue/re.fifo --message-group-id owais --message-deduplication-id test18 --message-body '{"task": "consume", "args": [{"test":"test"}], "kwargs": {}, "id": "test-task-id"}'
格式说明:
task字段值要和你定义的任务的name属性完全一致,这里就是consumeargs是传递给任务的位置参数列表,你要传的JSON对象放在数组里即可id是任务的唯一标识,可自定义任意字符串
方案2:兼容原生SQS消息(适用于无法修改消息发送端的场景)
如果你无法调整发送到SQS的消息格式,可以自定义消息处理器,把所有收到的原生消息直接转发给consume任务:
- 首先修改
celeryconfig.py,添加适配配置:
from kombu import ( Exchange, Queue ) broker_transport = 'sqs' broker_transport_options = { 'region': 'us-east-1', 'wait_time_seconds': 20, 'visibility_timeout': 3600 } worker_concurrency = 10 # 放开内容类型限制 accept_content = ['application/json', 'text/plain', 'json'] result_serializer = 'json' content_encoding = 'utf-8' task_serializer = 'json' worker_enable_remote_control = False worker_send_task_events = True result_backend = None task_queues = ( Queue('re.fifo', exchange=Exchange('consume', type='direct'), routing_key='consume'), ) task_routes = {'consume': {'queue': 're.fifo'}}
- 在
celery.py末尾添加自定义消费逻辑:
from kombu import Consumer from celery import bootsteps class SQSMessageConsumer(bootsteps.ConsumerStep): def get_consumers(self, channel): return [Consumer(channel, queues=[app.conf.task_queues[0]], callbacks=[self.handle_message], accept=['json', 'text/plain'])] def handle_message(self, body, message): # 直接将原生消息体传给consume任务执行 consume.delay(body) # 手动确认消息 message.ack() app.steps['consumer'].add(SQSMessageConsumer)
验证
启动Celery worker时添加日志级别参数,查看任务注册和消费情况:
celery -A celery worker --loglevel=info
启动后确认日志中显示已注册consume任务,且已监听re.fifo队列即可。
内容的提问来源于stack exchange,提问作者Muhammad Owais
相关产品推荐
相关产品推荐

