如何通过Celery指定RabbitMQ的自定义队列与交换机
解决方案
Celery默认内置了名为celery的队列、直连交换机和对应绑定规则,未显式配置自定义规则时就会自动使用这套默认配置,按以下步骤修改即可实现自定义队列、交换机投递:
1. 代码配置修改
需要先导入kombu的队列、交换机类,再声明自定义绑定规则:
from celery import Celery from kombu import Queue, Exchange import time broker_uri='amqp://xxxx:xxxx@xxxx:5672;amqp://xxxx:xxxx@xxxx:5672;amqp://xxxx:xxxx@xxxx:5672/' backend_uri="mongodb+srv://xxxxxx.mongodb.net/celery_test?retryWrites=true&w=majority" app = Celery('TestApp', broker=broker_uri,backend=backend_uri) # 声明自定义交换机,type可按需修改为direct/topic/fanout custom_exchange = Exchange('自定义交换机名', type='direct', durable=True) # 配置队列与交换机、路由键的绑定关系 app.conf.task_queues = ( Queue( '自定义队列名', exchange=custom_exchange, routing_key='自定义路由键', durable=True # 不需要队列持久化可设为False ), ) # 覆盖默认配置,避免自动生成默认celery队列 app.conf.task_default_queue = '自定义队列名' app.conf.task_default_exchange = '自定义交换机名' app.conf.task_default_routing_key = '自定义路由键' @app.task def reverse(text): time.sleep(10) return text[:-1]
2. 调整Worker启动命令
启动Worker时指定消费自定义队列:
celery -A 你的项目模块文件名 worker --queues=自定义队列名 -l info
启动后查看日志就能看到队列、交换机已经替换为你自定义的配置。
3. 任务投递方式
- 配置默认规则后直接调用即可投递到自定义队列:
reverse.delay("测试文本") - 多队列场景下可显式指定目标队列:
reverse.apply_async(args=["测试文本"], queue="自定义队列名")
内容的提问来源于stack exchange,提问作者Sandeepa Kariyawasam
相关产品推荐
相关产品推荐

