Concurrency and Modules
Queue Worker Threads
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
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(",")}"
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 end13enddo |worker_id|
pass 1 of 26workers = worker_count.times.map do |worker_id0|7 Thread.new do8 job = queue.pop9 lock.synchronize do10 completed << "#{worker_id}:#{job}"11 end12 end13endworkers ← [⟨Thread C /tmp/execution/⟨tmp E⟩.rb:65 run⟩, ⟨Thread D /tmp/execution/⟨tmp E⟩.rb:65 run⟩]
pass 2 of 26workers→ [⟨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 end13endworker_count.times do |index|
15worker_count2.times do |index|16 queue.push("job#{index + 1}")17enddo |index|
pass 1 of 215worker_count.times do |index0|16 queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17enddo |index|
pass 2 of 215worker_count2.times do |index1|16 queue⟨Thread::Queue A⟩.push("job#{index1 + 1}")17endjob ← job1
pass 1 of 26workers = 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 end13enddo
pass 1 of 28job = queue.pop9lock.synchronize do10 completed[] << "#{worker_id0}:#{jobjob1}"11endjob ← job2
pass 2 of 26workers = 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 end13enddo
pass 2 of 28job = queue.pop9lock.synchronize do10 completed["0:job1"] << "#{worker_id1}:#{jobjob2}"11endworkers.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
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 end13endworkers ← [⟨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 end13endworker_count.times do |index|
15worker_count1.times do |index|16 queue.push("job#{index + 1}")17enddo |index|
15worker_count1.times do |index0|16 queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17endjob ← 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 end13enddo
8job = queue.pop9lock.synchronize do10 completed[] << "#{worker_id0}:#{jobjob1}"11endworkers.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
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 end13enddo |worker_id|
pass 1 of 36workers = worker_count.times.map do |worker_id0|7 Thread.new do8 job = queue.pop9 lock.synchronize do10 completed << "#{worker_id}:#{job}"11 end12 end13endAll 3 passes — pass 1 is the card above pass worker_idworker_countworkers1 0 — — 2 1 — — 3 2 3 [⟨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⟩] worker_count.times do |index|
15worker_count3.times do |index|16 queue.push("job#{index + 1}")17enddo |index|
pass 1 of 315worker_count.times do |index0|16 queue⟨Thread::Queue A⟩.push("job#{index0 + 1}")17endAll 3 passes — pass 1 is the card above pass indexworker_count1 0 — 2 1 — 3 2 3 job ← job1
pass 1 of 36workers = 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 end13endAll 3 passes — pass 1 is the card above pass job1 job1 2 job2 3 job3 do
pass 1 of 38job = queue.pop9lock.synchronize do10 completed[] << "#{worker_id0}:#{jobjob1}"11endAll 3 passes — pass 1 is the card above pass completedworker_idjob1 [] 0 job1 2 ["0:job1"] 1 job2 3 ["0:job1", "1:job2"] 2 job3 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