【问题标题】:Delayed messages loop with RabbitMQ使用 RabbitMQ 的延迟消息循环
【发布时间】:2017-05-02 06:55:07
【问题描述】:

我正在尝试使用 Rabbit 的操作来实现拒绝/延迟循环,即:

我有:

  1. 主队列绑定了 Main Exchange,DLX 绑定到 StandBy Exchange。
  2. StandBy Queue 与 StandBy Exchange 绑定到它的 60s TTL 和 DLX 到 Main Exchange

基本上我想:

  1. 从主队列消费
  2. 拒绝消息(在某​​些情况下)
  3. 会因为拒绝而将其重定向到备用队列
  4. 当 TTL 到期时,将消息重新排队到主队列。

步骤 1、2 和 3 都可以,但最后一步丢弃消息而不是重新排队。

来自 RabbitMQ 文档的一些理论是我用来设计的:

来自队列的消息可能是“死信”;也就是说,当发生以下任何事件时,重新发布到另一个交易所:

  1. 邮件被拒绝(basic.reject 或 basic.nack),requeue=false,
  2. 消息的 TTL 过期;或
  3. 超出队列长度限制。

...

有可能形成消息死信的循环。例如,当队列死信消息发送到默认交换时,可能会发生这种情况而没有指定死信路由键。如果在整个周期中没有拒绝,则此类周期中的消息(即两次到达同一队列的消息)将被丢弃。

理论说它应该重新排队,因为它在第 2 步的循环中有一个 rejection,所以,你能帮我弄清楚为什么它会丢弃消息而不是 re - 排队?

更新:

我的目标版本是2.8.4,似乎在那一刻if there was no rejections in the entire cycle不在用例中,无论如何你可以自己检查RabbitMQ 2.8.x Docs

我会接受@george 的回答,因为此代码可以实现最初的目标。

【问题讨论】:

    标签: rabbitmq


    【解决方案1】:

    Rafael,我不确定您使用的是什么客户端,但是使用 Python 中的 Pika 客户端,您可以实现类似的功能。为简单起见,我只使用一次交换。你确定你正确设置了交换和路由密钥吗?

    发件人.py

    import sys
    import pika
    connection = pika.BlockingConnection(pika.ConnectionParameters(
                   'localhost'))
    channel = connection.channel()
    channel.exchange_declare(exchange='cycle', type='direct')
    channel.queue_declare(queue='standby_queue',
                          arguments={
                              'x-message-ttl': 10000,
                              'x-dead-letter-exchange': 'cycle',
                              'x-dead-letter-routing-key': 'main_queue'})
    channel.queue_declare(queue='main_queue',
                          arguments={
                              'x-dead-letter-exchange': 'cycle',
                              'x-dead-letter-routing-key': 'standby_queue'})
    channel.queue_bind(queue='main_queue', exchange='cycle')
    channel.queue_bind(queue='standby_queue', exchange='cycle')
    channel.basic_publish(exchange='cycle',
                          routing_key='main_queue',
                          body="message body")
    connection.close()
    

    receiver.py

    import sys
    import pika
    def callback(ch, method, properties, body):
        print "Processing message: {}".format(body)
        # replace with condition for rejection
        if True:
            print "Rejecting message"
            ch.basic_nack(method.delivery_tag, False, False)
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()
    channel.basic_consume(callback, queue='main_queue')
    channel.start_consuming()
    

    【讨论】:

    • 嗨@george,谢谢你的回复,我也在使用Python和Pika,我有一个非常相似的代码,但我实际上运行你的代码时得到了相同的结果,这让我想到了什么,你使用的是哪个版本的 RabbitMQ?,我用 2.8.4 和 3.6.6 版本尝试了这段代码
    • 我会接受你的问题,但我已经再次尝试针对这两个服务器并且在最新版本中有效。然后我进一步研究,旧文档没有提到“拒绝”异常,所以我猜在那种情况下它会丢弃任何试图重新排队的消息
    猜你喜欢
    • 2011-05-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多