【问题标题】:High priority queue over a lower one in ruby AMQP with RabbitMQ?使用 RabbitMQ 的 ruby​​ AMQP 中的高优先级队列优于低优先级队列?
【发布时间】:2013-06-05 11:02:33
【问题描述】:

鉴于我有一个工作人员订阅了两个队列“低”和“高”,我希望工作人员仅在高优先级队列为空时处理来自低优先级队列的消息。

我正在尝试通过定义两个通道并将预取设置为更高优先级队列上的更高值来做到这一点,如下所示:http://dougbarth.github.io/2011/07/01/approximating-priority-with-rabbitmq.html

这是我的工人代码:

require "rubygems"
require "amqp"

EventMachine.run do
  connection = AMQP.connect(:host => '127.0.0.1')
  channel_low  = AMQP::Channel.new(connection)
  channel_high  = AMQP::Channel.new(connection)

  # Attempting to set the prefetch higher on the high priority queue
  channel_low.prefetch(10)
  channel_high.prefetch(20)

  low_queue    = channel_low.queue("low", :auto_delete => false)
  high_queue    = channel_high.queue("high", :auto_delete => false)

  low_queue.subscribe do |payload|
    puts "#{payload}"
    slow_task
  end

  high_queue.subscribe do |payload|
    puts "#{payload}"
    slow_task
  end

  def slow_task
    # Do some slow work
    sleep(1)
  end
end

当我针对它运行此客户端时,我没有看到首先处理的高优先级消息:

require "rubygems"
require "amqp"

EventMachine.run do
  connection = AMQP.connect(:host => '127.0.0.1')
  channel  = AMQP::Channel.new(connection)

  low_queue    = channel.queue("low")
  high_queue    = channel.queue("high")
  exchange = channel.direct("")

  10.times do |i| 
    message = "LOW #{i}"
    puts "sending: #{message}"
    exchange.publish message, :routing_key => low_queue.name
  end

  # EventMachine.add_periodic_timer(0.0001) do
  10.times do |i|
    message = "HIGH #{i}"
    puts "sending: #{message}"
    exchange.publish message, :routing_key => high_queue.name
  end

end

输出:

Client >>>
sending: LOW 0
sending: LOW 1
sending: LOW 2
sending: LOW 3
sending: LOW 4
sending: LOW 5
sending: LOW 6
sending: LOW 7
sending: LOW 8
sending: LOW 9
sending: HIGH 0
sending: HIGH 1
sending: HIGH 2
sending: HIGH 3
sending: HIGH 4
sending: HIGH 5
sending: HIGH 6
sending: HIGH 7
sending: HIGH 8
sending: HIGH 9

Server >>>

HIGH 0
HIGH 1
LOW 0
LOW 1
LOW 2
HIGH 2
LOW 3
LOW 4
LOW 5
LOW 6
LOW 7
HIGH 3
LOW 8
LOW 9
HIGH 4
HIGH 5
HIGH 6
HIGH 7
HIGH 8
HIGH 9

【问题讨论】:

  • 永远不要在事件机器中使用sleep,它会暂停整个事件循环,你的并发性会消失。请改用EM.add_timer。请记住,只有在您完成消息后才能手动确认消息。
  • 预取不会规定任何优先级,只有多少消息将被“飞行”,即。已将消息传递给客户端但尚未确认消息,因此它更多的是关于并发而不是优先级..
  • 我正在使用睡眠延迟来模拟真正的工作人员,它运行一个阻塞操作,大约需要 0.3 秒才能完成。这就是这段代码被分解为工人的原因。

标签: ruby rabbitmq priority-queue amqp


【解决方案1】:

正如迈克尔所说,您当前的方法存在一些问题:

  • 不启用显式确认意味着 RabbitMQ 在发送消息时考虑消息传递,而不是在您处理消息时考虑
  • 您的消息非常小,可以通过网络快速传递
  • 您的 subscribe 块在 EventMachine 读取网络数据时被调用,对于它从套接字读取的每个完整数据帧调用一次
  • 最后,阻塞反应器线程(使用睡眠)将阻止 EM 将确认发送到套接字,因此无法实现正确的行为。

为了实现优先级概念,您需要将接收网络数据与确认网络数据分开。在我们的应用程序(我写过博客文章)中,我们使用了一个后台线程和一个优先级队列来重新排序进入的工作。这在每个工作人员中引入了一个小的消息缓冲区。其中一些消息可能是低优先级的消息,在没有更高优先级的消息可以处理之前,它们不会被处理。

这是一个稍作修改的工作代码,它使用工作线程和优先级队列来获得所需的结果。

require "rubygems"
require "amqp"
require "pqueue"

EventMachine.run do
  connection = AMQP.connect(:host => '127.0.0.1')
  channel_low  = AMQP::Channel.new(connection)
  channel_high  = AMQP::Channel.new(connection)

  # Attempting to set the prefetch higher on the high priority queue
  channel_low.prefetch(10)
  channel_high.prefetch(20)

  low_queue    = channel_low.queue("low", :auto_delete => false)
  high_queue    = channel_high.queue("high", :auto_delete => false)

  # Our priority queue for buffering messages in the worker's memory
  to_process = PQueue.new {|a,b| a[0] > b[0] }

  # The pqueue gem isn't thread safe
  mutex = Mutex.new

  # Background thread for working blocking operation. We can spin up more of
  # these to increase concurrency.
  Thread.new do
    loop do
      _, header, payload = mutex.synchronize { to_process.pop }

      if payload
        puts "#{payload}"
        slow_task
        # We need to call ack on the EM thread.
        EM.next_tick { header.ack }
      else
        sleep(0.1)
      end
    end
  end

  low_queue.subscribe(:ack => true) do |header, payload|
    mutex.synchronize { to_process << [0, header, payload] }
  end

  high_queue.subscribe(:ack => true) do |header, payload|
    mutex.synchronize { to_process << [10, header, payload] }
  end

  def slow_task
    # Do some slow work
    sleep(1)
  end
end

如果您需要增加并发性,您可以生成多个后台线程。

【讨论】:

    【解决方案2】:

    您的方法是最常用的解决方法之一。然而,有几个 发布的具体示例中的问题。

    Prefetch 不控制哪个通道具有交付优先级。它控制有多少消息 可以“进行中”(未确认)。

    这可以用作穷人的优先级技术,但是,您使用自动确认模式, 所以预取不会发挥作用(RabbitMQ 立即认为您的消息已确认 一旦发送出去)。

    如果您只发布少量消息,然后您的示例完成运行, 排序很可能很大程度上取决于您发布消息的顺序。

    要通过手动确认查看预取设置的效果,您需要运行更长的时间(这取决于您的消息速率,但至少一分钟。

    【讨论】:

      【解决方案3】:

      你能试试这样的吗?我不确定这是否能满足您的需求,但值得一试。

      10.times do 
         10.times do 
            exchange.publish "HIGH", :routing_key => high_queue.name
          end  
         exchange.publish "LOW", :routing_key => low_queue.name
      end
      

      【讨论】:

      • 抱歉,这可能不清楚,所有消息都会立即生成,我正在通过先创建低级消息进行测试。我想要工作的是工人在低位之前先在高位上工作。 (一低就可以了)
      • 你能睡一两分钟然后检查一下。睡眠(2.分钟)
      • 将其设置为睡眠(10)并生成:`LOW LOW LOW LOW LOW LOW LOW LOW LOW LOW HIGH HIGH HIGH HIGH HIGH HIGH
      【解决方案4】:

      您正在向处理器发布小消息,该处理器可能会在它们到达后的几毫秒内处理它们。我怀疑你在这里看到的是样本噪音,而不是你真正认为你看到的。

      您能否验证消息发布的顺序是否符合您的预期?

      【讨论】:

      • 我在消息中添加了一些额外的细节以显示它们的顺序并编辑了上面的测试和结果
      • 您仍在显示它们的消费顺序,而不是它们的生产顺序。您需要做的是在发布时打印正文,而不是在消费时打印。您可能会发现发布顺序与您假设的不同。
      • 试一试,在它发送消息时也可以输出。它以预期的顺序发送它们,但另一端仍然很疯狂。
      猜你喜欢
      • 2015-01-06
      • 2013-02-13
      • 2011-12-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-02-24
      • 2011-03-20
      相关资源
      最近更新 更多