【问题标题】:How do I manage ruby threads so they finish all their work?如何管理 ruby​​ 线程以便它们完成所有工作?
【发布时间】:2023-03-08 12:10:01
【问题描述】:

我有一个可以划分为独立单元的计算,我现在处理它的方式是创建固定数量的线程,然后交出要在每个线程中完成的工作块。所以在伪代码中是这样的

# main thread
work_units.take(10).each {|work_unit| spawn_thread_for work_unit}

def spawn_thread_for(work)
  Thread.new do
    do_some work
    more_work = work_units.pop
    spawn_thread_for more_work unless more_work.nil?
  end
end

基本上,一旦创建了初始数量的线程,每个线程都会做一些工作,然后继续从工作堆栈中获取要完成的工作,直到没有留下任何东西。当我在 irb 中运行时一切正常,但是当我使用解释器执行脚本时,一切都没有那么好。我不确定如何让主线程等到所有工作完成。有没有一种很好的方法可以做到这一点,还是我坚持在主线程中执行sleep 10 until work_units.empty?

【问题讨论】:

  • take(10) 不是意味着只会处理前 10 个work_units 吗?
  • @Andrew Grimm:实际上没有,但是重读这个问题,看起来work_units 将超出spawn_thread_for 的范围more_work = work_units.pop
  • 在实际代码中,它是一个实例变量,因此不存在范围问题。
  • 注意:depending on which ruby implementation you use,线程可能不会真正并行工作(例如,使用多个 CPU 内核/通过拆分多个线程来更快地完成工作)。

标签: ruby multithreading threadpool


【解决方案1】:

您可以使用Thread#join

加入(p1 = v1)公开

调用线程将暂停执行并运行thr。在 thr 退出或经过限制秒数之前不会返回。如果超时,则返回nil,否则返回thr。

您也可以使用Enumerable#each_slice 批量迭代工作单元

work_units.each_slice(10) do |batch|
  # handle each work unit in a thread
  threads = batch.map do |work_unit|
    spawn_thread_for work_unit
  end

  # wait until current batch work units finish before handling the next batch
  threads.each(&:join)
end

【讨论】:

    【解决方案2】:
    Thread.list.each{ |t| t.join unless t == Thread.current }
    

    【讨论】:

      【解决方案3】:

      在 ruby​​ 1.9(和 2.0)中,您可以为此目的使用 stdlib 中的 ThreadsWait

      require 'thread'
      require 'thwait'
      
      threads = []
      threads << Thread.new { }
      threads << Thread.new { }
      ThreadsWait.all_waits(*threads)
      

      【讨论】:

      • 你能解释一下ThreadsWait.all_waits(*threads)threads.each { |thr| thr.join }之间的区别吗?
      • 第二种方式是现代API。
      【解决方案4】:

      如果您修改spawn_thread_for 以保存对您创建的Thread 的引用,那么您可以在线程上调用Thread#join 以等待完成:

      x = Thread.new { sleep 0.1; print "x"; print "y"; print "z" }
      a = Thread.new { print "a"; print "b"; sleep 0.2; print "c" }
      x.join # Let the threads finish before
      a.join # main thread exits...
      

      产生:

      abxyzc
      

      (从ri Thread.new 文档中窃取。有关更多详细信息,请参阅ri Thread.join 文档。)

      因此,如果您修改 spawn_thread_for 以保存 Thr​​ead 引用,您可以将它们全部加入:

      (未经测试,但应该有味道)

      # main thread
      work_units = Queue.new # and fill the queue...
      
      threads = []
      10.downto(1) do
        threads << Thread.new do
          loop do
            w = work_units.pop
            Thread::exit() if w.nil?
            do_some_work(w)
          end
        end
      end
      
      # main thread continues while work threads devour work
      
      threads.each(&:join)
      

      【讨论】:

      • 但这可能会使threads 相当大,因为当线程完成时,它的引用仍在threads 数组中。此外,threads 正在修改,而 threads.each(&amp;:join) 正在执行,所以这仍然不能解决问题。
      • @davidk,也许不是在每个线程中调用spawn_thread_for,而是简单地寻找更多工作并直接开始?
      • @sarnold,我只想要固定数量的线程处于活动状态。如果我寻找更多的工作,然后启动一个线程,我可能会超过 10 个线程的限制,因为调度发生的方式和每个线程正在执行的工作量的不确定性。
      • @davidk01,为什么要启动另一个线程,当一个“完成”时?只需在已经运行的线程中获取一个新的工作单元。启动十个线程,然后让它们开始从您的工作队列中吸出工作单元。 (有关线程安全队列的实现,请参见 Queue 类。)
      • @sarnold,好点子。没想到,谢谢指点。
      【解决方案5】:

      您似乎在复制 Parallel Each (Peach) 库提供的内容。

      【讨论】:

      • 对于我想到的简单用例,我不需要整个库,但感谢您提供的链接。
      猜你喜欢
      • 2011-12-17
      • 2011-01-31
      • 2010-10-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多