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

如何消费RabbitMQ中除指定routing_key外的所有消息?

解决RabbitMQ Topic交换机排除特定Routing Key的问题

首先得明确:RabbitMQ的Topic交换机本身不支持否定匹配,没法直接通过绑定一个routing key来排除A.B.A这类特定值。不过我们有两种靠谱的方案来实现你的需求,下面分别给出Python+pika的代码示例。

方案一:使用Alternate Exchange(备用交换机)推荐

这个方案是在路由层面解决问题,所有不匹配A.B.A的消息会自动被转发到备用交换机,再路由到你的第二个队列,完全不需要在消费端做过滤,效率更高。

步骤说明

  1. 创建一个Fanout类型的备用交换机(Fanout会把所有消息转发给绑定的队列,刚好符合我们接收所有无法路由消息的需求)
  2. 创建第二个队列,绑定到备用交换机
  3. 创建主Topic交换机,设置它的alternate-exchange属性为备用交换机的名称
  4. 创建第一个队列,绑定到主交换机,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:24:50