【问题标题】:EventMachine with em-synchrony I need to correctly throttle my http requestsEventMachine 与 em-synchrony 我需要正确限制我的 http 请求
【发布时间】:2012-08-29 03:29:51
【问题描述】:

我有一个消费者,它通过事件订阅从队列中提取消息。它接收这些消息,然后连接到一个相当慢的 http 接口。我有一个 8 个工作池,一旦这些都填满,我需要停止从队列中拉取请求,并让正在处理 http 作业的纤程继续工作。这是我整理的一个例子。

def send_request(callback)
  EM.synchrony do

    while $available <= 0

      sleep 2 

      puts "sleeping"
    end 
    url = 'http://example.com/api/Restaurant/11111/images/?image%5Bremote_url%5D=https%3A%2F%2Firs2.4sqi.net%2Fimg%2Fgeneral%2Foriginal%2F8NMM4yhwsLfxF-wgW0GA8IJRJO8pY4qbmCXuOPEsUTU.jpg&image%5Bsource_type_enum%5D=3'
    result = EM::Synchrony.sync EventMachine::HttpRequest.new(url, :inactivity_timeout => 0).send("apost", :head => {:Accept => 'services.v1'})

    callback.call(result.response) 
  end 
end

def display(value)
  $available += 1
  puts value.inspect
end

$available = 8 

EM.run do
  EM.add_periodic_timer(0.001) do
    $available -= 1
    puts "Available: #{$available}"

    puts "Tick ..." 
    puts send_request(method(:display))
  end 

end

我发现,如果我在同步块中的 while 循环内调用 sleep,反应器循环就会卡住。如果我在 if 语句中调用 sleep (只睡一次),那么大多数情况下请求完成的时间是足够的,但充其量是不可靠的。如果我使用 EM::Synchrony.sleep,那么主反应器循环将不断创建新请求。

有没有办法暂停主循环但让纤维完成它们的执行?

【问题讨论】:

    标签: ruby http asynchronous eventmachine fibers


    【解决方案1】:
    sleep 2
    

    ...

    add_periodic_timer(0.001)
    

    你是认真的吗?

    你有没有想过有多少send_request 在循环中休眠?而且每秒增加 1000 个。

    这个呢:

    require 'eventmachine'
    require 'em-http'
    require 'fiber'
    
    class Worker
      URL = 'http://example.com/api/whatever'
    
      def initialize callback
        @callback = callback
      end
    
      def work
        f = Fiber.current
        loop do
          http = EventMachine::HttpRequest.new(URL).get :timeout => 20
    
          http.callback do
            @callback.call http.response
            f.resume
          end
          http.errback do
            f.resume
          end
    
          Fiber.yield
        end
      end
    end
    
    def display(value)
      puts "Done: #{value.size}"
    end
    
    EventMachine.run do
      8.times do
        Fiber.new do
          Worker.new(method(:display)).work
        end.resume
      end
    end
    

    【讨论】:

    • 感谢您的回复。可能是我说的不够清楚。我在示例中使用 add_periodic_timer 来模拟我收到的消息量。实际代码通过 subscribe 连接到 rabbitmq。 AFAIK 我无法限制在消费者端接收消息。这就是我决定用一个工人池限制我的消费者并在工人忙时暂停主反应器循环的原因。由于我的消费者作为守护进程运行,因此简单地执行“8.times do”并不能解决我的问题,因为会有无限量的消息。
    猜你喜欢
    • 2011-12-12
    • 1970-01-01
    • 2013-11-12
    • 2012-06-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-02
    相关资源
    最近更新 更多