如何消费RabbitMQ中除指定routing_key外的所有消息?
解决RabbitMQ Topic交换机排除特定Routing Key的问题
首先得明确:RabbitMQ的Topic交换机本身不支持否定匹配,没法直接通过绑定一个routing key来排除A.B.A这类特定值。不过我们有两种靠谱的方案来实现你的需求,下面分别给出Python+pika的代码示例。
方案一:使用Alternate Exchange(备用交换机)推荐
这个方案是在路由层面解决问题,所有不匹配A.B.A的消息会自动被转发到备用交换机,再路由到你的第二个队列,完全不需要在消费端做过滤,效率更高。
步骤说明
- 创建一个Fanout类型的备用交换机(Fanout会把所有消息转发给绑定的队列,刚好符合我们接收所有无法路由消息的需求)
- 创建第二个队列,绑定到备用交换机
- 创建主Topic交换机,设置它的
alternate-exchange属性为备用交换机的名称 - 创建第一个队列,绑定到主交换机,routing key设为
A.B.A
这样,当消息的routing key是A.B.A时,会被主交换机路由到第一个队列;其他所有消息因为无法匹配主交换机的绑定规则,会被自动转发到备用交换机,最终进入第二个队列。
生产者代码
import pika # 建立连接 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 1. 声明备用交换机(Fanout类型) alternate_exchange_name = 'alternate_exchange' channel.exchange_declare(exchange=alternate_exchange_name, exchange_type='fanout') # 2. 声明第二个队列并绑定到备用交换机 queue_name_2 = 'non_aba_queue' channel.queue_declare(queue=queue_name_2, durable=True) channel.queue_bind(exchange=alternate_exchange_name, queue=queue_name_2) # 3. 声明主Topic交换机,指定备用交换机 main_exchange_name = 'main_topic_exchange' channel.exchange_declare( exchange=main_exchange_name, exchange_type='topic', durable=True, arguments={'alternate-exchange': alternate_exchange_name} # 关键:设置备用交换机 ) # 4. 声明第一个队列并绑定到主交换机,routing key为A.B.A queue_name_1 = 'aba_queue' channel.queue_declare(queue=queue_name_1, durable=True) channel.queue_bind(exchange=main_exchange_name, queue=queue_name_1, routing_key='A.B.A') # 发送测试消息 # 测试1:会进入aba_queue channel.basic_publish( exchange=main_exchange_name, routing_key='A.B.A', body=b'Message for A.B.A', properties=pika.BasicProperties(delivery_mode=2) # 持久化消息 ) # 测试2:会进入non_aba_queue channel.basic_publish( exchange=main_exchange_name, routing_key='A.B.B', body=b'Message for A.B.B', properties=pika.BasicProperties(delivery_mode=2) ) # 测试3:会进入non_aba_queue(routing key长度不同) channel.basic_publish( exchange=main_exchange_name, routing_key='A.B', body=b'Message for A.B', properties=pika.BasicProperties(delivery_mode=2) ) print("消息发送完成") connection.close()
第一个队列的消费者代码(接收A.B.A消息)
import pika def callback(ch, method, properties, body): print(f"收到A.B.A队列的消息: {body.decode()}") ch.basic_ack(delivery_tag=method.delivery_tag) connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() channel.queue_declare(queue='aba_queue', durable=True) channel.basic_consume(queue='aba_queue', on_message_callback=callback) print("等待A.B.A队列的消息...") channel.start_consuming()
第二个队列的消费者代码(接收非A.B.A消息)
import pika def callback(ch, method, properties, body): print(f"收到非A.B.A队列的消息: {body.decode()},routing key: {method.routing_key}") ch.basic_ack(delivery_tag=method.delivery_tag) connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() channel.queue_declare(queue='non_aba_queue', durable=True) channel.basic_consume(queue='non_aba_queue', on_message_callback=callback) print("等待非A.B.A队列的消息...") channel.start_consuming()
方案二:消费端过滤(简单直接)
如果不想配置备用交换机,也可以直接给第二个队列绑定#(匹配所有routing key),然后在消费者中过滤掉A.B.A的消息,把这类消息直接拒绝并丢弃(或者根据需求处理)。
第二个队列的消费者代码(带过滤逻辑)
import pika def callback(ch, method, properties, body): if method.routing_key == 'A.B.A': # 拒绝该消息,并且不重新入队(避免死循环) ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False) print(f"过滤掉A.B.A的消息: {body.decode()}") else: print(f"收到非A.B.A队列的消息: {body.decode()},routing key: {method.routing_key}") ch.basic_ack(delivery_tag=method.delivery_tag) connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() main_exchange_name = 'main_topic_exchange' channel.exchange_declare(exchange=main_exchange_name, exchange_type='topic', durable=True) queue_name_2 = 'non_aba_queue_filtered' channel.queue_declare(queue=queue_name_2, durable=True) # 绑定#,匹配所有消息 channel.queue_bind(exchange=main_exchange_name, queue=queue_name_2, routing_key='#') channel.basic_consume(queue=queue_name_2, on_message_callback=callback) print("等待非A.B.A队列的消息(消费端过滤)...") channel.start_consuming()
这个方案的缺点是A.B.A的消息会先路由到第二个队列,再被过滤掉,稍微浪费一点资源,但胜在配置简单,适合快速实现。
内容的提问来源于stack exchange,提问作者Evgeny
相关产品推荐
相关产品推荐

