【问题标题】:Ruby 3 collecting results from multiple scheduled fibersRuby 3 从多个预定纤程中收集结果
【发布时间】:2022-12-11 11:06:13
【问题描述】:

Ruby 3 引入了Fiber.schedule 来同时分派异步任务。

与this question(关于线程并发)中的问题类似,我想要一种在纤程调度程序上启动多个并发任务的方法,一旦它们都被安排好等待它们的组合结果,有点等同于Promise.all在 JavaScript 中。

我可以想出这种天真的方法:

require 'async'

def io_work(t)
  sleep t
  :ok
end

Async do
  results = []

  [0.1, 0.3, 'cow'].each_with_index do |t, i|
    n = i + 1
    Fiber.schedule do
      puts "Starting fiber #{n}\n"
      result = io_work t
      puts "Done working for #{t} seconds in fiber #{n}"
      results << [n, result]
    rescue
      puts "Execution failed in fiber #{n}"
      results << [n, :error]
    end
  end

  # await combined results
  sleep 0.1 until results.size >= 3

  puts "Results: #{results}"
end

有没有更简单的结构可以做同样的事情?

【问题讨论】:

    标签: ruby asynchronous concurrency fibers ruby-3


    【解决方案1】:

    由于 Async 任务已经安排好了,我不确定你是否需要所有这些。

    如果您只想等待所有项目完成,您可以使用Async::Barrier

    例子:

    require 'async'
    require 'async/barrier'
    
    
    def io_work(t)
      sleep t
      :ok
    end
    
    Async do
      barrier = Async::Barrier.new
      results = []
      [1, 0.3, 'cow'].each.with_index(1) do |data, idx|
        barrier.async do 
          results << begin
            puts "Starting task #{idx}
    "
            result = io_work data
            puts "Done working for #{data} seconds in task #{idx}"
            [idx,result]
          rescue
            puts "Execution failed in task #{idx}"
            [idx, :error]
          end          
        end 
      end
      barrier.wait
      puts "Results: #{results}"
    end
    

    基于 sleep 值,这将输出

    Starting task 1
    Starting task 2
    Starting task 3
    Execution failed in task 3
    Done working for 0.3 seconds in task 2
    Done working for 1 seconds in task 1
    Results: [[3, :error], [2, :ok], [1, :ok]]
    

    barrier.wait 将等到所有异步任务完成,没有它输出看起来像

    Starting fiber 1
    Starting fiber 2
    Starting fiber 3
    Execution failed in fiber 3
    Results: [[3, :error]]
    Done working for 0.3 seconds in fiber 2
    Done working for 1 seconds in fiber 1
    

    【讨论】:

    • 使用Async::Barrier 是个好主意,因为它会专门跟踪任务执行情况,而Async::Barrier#wait 会等到所有任务完成。这与使用 Thread::Queue 不同,您需要以某种方式通知父任务所有工作都已完成。
    【解决方案2】:

    我对解决方案的人体工程学不满意,所以我创建了 the gem fiber-collector 来解决它。

    免责声明:我正在描述一个我是其作者的图书馆

    问题场景中的示例用法:

    require 'async'
    require 'fiber/collector'
    
    def io_work(t)
      sleep t
      :ok
    end
    
    Async do
      Fiber::Collector.schedule { io_work(1) }.and { io_work(0.3) }.all
    end.wait
    # => [:ok, :ok]
    
    
    Async do
      Fiber::Collector.schedule { io_work(1) }.and { io_work(0.3) }.and { io_work('cow') }.all
    end.wait
    # => raises error
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-04-26
      • 2018-04-12
      • 1970-01-01
      • 2017-02-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-05-30
      相关资源
      最近更新 更多