A queue lets worker threads take jobs from a shared input without manually choosing which worker receives each job.

queue workers Start workers first, push one job for each worker, then join after every queued job has been claimed.

Queue Worker Threads

worker_count
queue_worker_threads.rb
Replay: real traced execution (multi-file project)
worker_count = 2
queue = Queue.new
completed = []
lock = Mutex.new

workers = worker_count.times.map do |worker_id|
  Thread.new do
    job = queue.pop
    lock.synchronize do
      completed << "#{worker_id}:#{job}"
    end
  end
end

worker_count.times do |index|
  queue.push("job#{index + 1}")
end

workers.each(&:join)

puts "workers=#{worker_count}"
puts "completed=#{completed.length}"
puts "worker_entries=#{completed.sort.join(",")}"
worker_count = 1
queue = Queue.new
completed = []
lock = Mutex.new

workers = worker_count.times.map do |worker_id|
  Thread.new do
    job = queue.pop
    lock.synchronize do
      completed << "#{worker_id}:#{job}"
    end
  end
end

worker_count.times do |index|
  queue.push("job#{index + 1}")
end

workers.each(&:join)

puts "workers=#{worker_count}"
puts "completed=#{completed.length}"
puts "worker_entries=#{completed.sort.join(",")}"
worker_count = 3
queue = Queue.new
completed = []
lock = Mutex.new

workers = worker_count.times.map do |worker_id|
  Thread.new do
    job = queue.pop
    lock.synchronize do
      completed << "#{worker_id}:#{job}"
    end
  end
end

worker_count.times do |index|
  queue.push("job#{index + 1}")
end

workers.each(&:join)

puts "workers=#{worker_count}"
puts "completed=#{completed.length}"
puts "worker_entries=#{completed.sort.join(",")}"
  1. worker_count ← 2, queue ← ⟨Thread::Queue A⟩, completed ← [], lock ← ⟨Thread::Mutex B⟩

    1worker_count→ 2 = 2  #@worker_count=1, 32queue→ ⟨Thread::Queue A⟩ = Queue.new3completed→ [] = []4lock→ ⟨Thread::Mutex B⟩ = Mutex.new56workers = worker_count2.times.map do |worker_id|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  2. do |worker_id|

    pass 1 of 2
    6workers = worker_count.times.map do |worker_id0|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  3. workers ← [⟨Thread C /tmp/execution/⟨tmp E⟩.rb:65 run⟩, ⟨Thread D /tmp/execution/⟨tmp E⟩.rb:65 run⟩]

    pass 2 of 2
    6workers→ [⟨Thread C /tmp/execution/⟨tmp E⟩.rb:65 run⟩, ⟨Thread D /tmp/execution/⟨tmp E⟩.rb:65 run⟩] = worker_count2.times.map do |worker_id1|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  4. worker_count.times do |index|

    15worker_count2.times do |index|16  queue.push("job#{index + 1}")17end
  5. do |index|

    pass 1 of 2
    15worker_count.times do |index0|16  queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17end
  6. do |index|

    pass 2 of 2
    15worker_count2.times do |index1|16  queue⟨Thread::Queue A⟩.push("job#{index1 + 1}")17end
  7. job ← job1

    pass 1 of 2
    6workers = worker_count.times.map do |worker_id|7  Thread.new do8    job→ job1 = queue⟨Thread::Queue A⟩.pop9    lock⟨Thread::Mutex B⟩.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  8. do

    pass 1 of 2
    8job = queue.pop9lock.synchronize do10  completed[] << "#{worker_id0}:#{jobjob1}"11end
  9. job ← job2

    pass 2 of 2
    6workers = worker_count.times.map do |worker_id|7  Thread.new do8    job→ job2 = queue⟨Thread::Queue A⟩.pop9    lock⟨Thread::Mutex B⟩.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  10. do

    pass 2 of 2
    8job = queue.pop9lock.synchronize do10  completed["0:job1"] << "#{worker_id1}:#{jobjob2}"11end
  11. workers.each(&:join)

    19workers.each(&:join)[⟨Thread C /tmp/execution/⟨tmp E⟩.rb:65 dead⟩, ⟨Thread D /tmp/execution/⟨tmp E⟩.rb:65 dead⟩]2021puts "workers=#{worker_count2}"22puts "completed=#{completed.length2}"23puts "worker_entries=#{completed.sort.join(",")0:job1,1:job2}"
    outputworkers=2
    completed=2
    worker_entries=0:job1,1:job2
  1. worker_count ← 1, queue ← ⟨Thread::Queue A⟩, completed ← [], lock ← ⟨Thread::Mutex B⟩

    1worker_count→ 1 = 12queue→ ⟨Thread::Queue A⟩ = Queue.new3completed→ [] = []4lock→ ⟨Thread::Mutex B⟩ = Mutex.new56workers = worker_count1.times.map do |worker_id|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  2. workers ← [⟨Thread C /tmp/execution/⟨tmp D⟩.rb:65 run⟩]

    6workers→ [⟨Thread C /tmp/execution/⟨tmp D⟩.rb:65 run⟩] = worker_count1.times.map do |worker_id0|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  3. worker_count.times do |index|

    15worker_count1.times do |index|16  queue.push("job#{index + 1}")17end
  4. do |index|

    15worker_count1.times do |index0|16  queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17end
  5. job ← job1

    6workers = worker_count.times.map do |worker_id|7  Thread.new do8    job→ job1 = queue⟨Thread::Queue A⟩.pop9    lock⟨Thread::Mutex B⟩.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  6. do

    8job = queue.pop9lock.synchronize do10  completed[] << "#{worker_id0}:#{jobjob1}"11end
  7. workers.each(&:join)

    19workers.each(&:join)[⟨Thread C /tmp/execution/⟨tmp D⟩.rb:65 dead⟩]2021puts "workers=#{worker_count1}"22puts "completed=#{completed.length1}"23puts "worker_entries=#{completed.sort.join(",")0:job1}"
    outputworkers=1
    completed=1
    worker_entries=0:job1
  1. worker_count ← 3, queue ← ⟨Thread::Queue A⟩, completed ← [], lock ← ⟨Thread::Mutex B⟩

    1worker_count→ 3 = 32queue→ ⟨Thread::Queue A⟩ = Queue.new3completed→ [] = []4lock→ ⟨Thread::Mutex B⟩ = Mutex.new56workers = worker_count3.times.map do |worker_id|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
  2. do |worker_id|

    pass 1 of 3
    6workers = worker_count.times.map do |worker_id0|7  Thread.new do8    job = queue.pop9    lock.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
    All 3 passes — pass 1 is the card above
    passworker_idworker_countworkers
    10
    21
    323[⟨Thread C /tmp/execution/⟨tmp F⟩.rb:65 run⟩, ⟨Thread D /tmp/execution/⟨tmp F⟩.rb:65 run⟩, ⟨Thread E /tmp/execution/⟨tmp F⟩.rb:65 run⟩]
  3. worker_count.times do |index|

    15worker_count3.times do |index|16  queue.push("job#{index + 1}")17end
  4. do |index|

    pass 1 of 3
    15worker_count.times do |index0|16  queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17end
    All 3 passes — pass 1 is the card above
    passindexworker_count
    10
    21
    323
  5. job ← job1

    pass 1 of 3
    6workers = worker_count.times.map do |worker_id|7  Thread.new do8    job→ job1 = queue⟨Thread::Queue A⟩.pop9    lock⟨Thread::Mutex B⟩.synchronize do10      completed << "#{worker_id}:#{job}"11    end12  end13end
    All 3 passes — pass 1 is the card above
    passjob
    1job1
    2job2
    3job3
  6. do

    pass 1 of 3
    8job = queue.pop9lock.synchronize do10  completed[] << "#{worker_id0}:#{jobjob1}"11end
    All 3 passes — pass 1 is the card above
    passcompletedworker_idjob
    1[]0job1
    2["0:job1"]1job2
    3["0:job1", "1:job2"]2job3
  7. workers.each(&:join)

    19workers.each(&:join)[⟨Thread C /tmp/execution/⟨tmp F⟩.rb:65 dead⟩, ⟨Thread D /tmp/execution/⟨tmp F⟩.rb:65 dead⟩, ⟨Thread E /tmp/execution/⟨tmp F⟩.rb:65 dead⟩]2021puts "workers=#{worker_count3}"22puts "completed=#{completed.length3}"23puts "worker_entries=#{completed.sort.join(",")0:job1,1:job2,2:job3}"
    outputworkers=3
    completed=3
    worker_entries=0:job1,1:job2,2:job3