Class: Omnizip::Parallel::WorkerPool
- Inherits:
-
Object
- Object
- Omnizip::Parallel::WorkerPool
- Defined in:
- lib/omnizip/parallel/worker_pool.rb
Overview
Worker pool wrapper for Fractor-based parallel processing
Manages a pool of Fractor workers for parallel compression/extraction. Handles job distribution, result collection, and graceful shutdown.
Instance Attribute Summary collapse
-
#errors ⇒ Array
readonly
Collected errors.
-
#results ⇒ Array
readonly
Collected results.
-
#running ⇒ Boolean
readonly
Whether pool is running.
-
#supervisor ⇒ Fractor::Supervisor
readonly
Underlying Fractor supervisor.
Instance Method Summary collapse
-
#complete? ⇒ Boolean
Check if pool has completed all work.
-
#failed_results ⇒ Array
Get failed results.
-
#initialize(worker_class:, num_workers: nil, continuous: false) ⇒ WorkerPool
constructor
Initialize worker pool.
-
#run ⇒ void
Run the pool in batch mode and wait for completion.
-
#shutdown(timeout: 30) ⇒ void
Shutdown the worker pool.
-
#start ⇒ void
Start the worker pool.
-
#stats ⇒ Hash
Get pool statistics.
-
#submit(work) ⇒ void
Submit work item to the pool.
-
#submit_batch(works) ⇒ void
Submit multiple work items.
-
#successful_results ⇒ Array
Get successful results.
Constructor Details
#initialize(worker_class:, num_workers: nil, continuous: false) ⇒ WorkerPool
Initialize worker pool
39 40 41 42 43 44 45 46 47 48 49 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 39 def initialize(worker_class:, num_workers: nil, continuous: false) @worker_class = worker_class @num_workers = num_workers || detect_cpu_count @continuous = continuous @supervisor = nil @results = [] @errors = [] @running = false @result_mutex = Mutex.new @work_queue = nil end |
Instance Attribute Details
#errors ⇒ Array (readonly)
Returns collected errors.
29 30 31 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 29 def errors @errors end |
#results ⇒ Array (readonly)
Returns collected results.
26 27 28 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 26 def results @results end |
#running ⇒ Boolean (readonly)
Returns whether pool is running.
32 33 34 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 32 def running @running end |
#supervisor ⇒ Fractor::Supervisor (readonly)
Returns underlying Fractor supervisor.
23 24 25 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 23 def supervisor @supervisor end |
Instance Method Details
#complete? ⇒ Boolean
Check if pool has completed all work
179 180 181 182 183 184 185 186 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 179 def complete? return false if @continuous return false unless @supervisor # In batch mode, check if all work is processed result_aggregator = @supervisor.results result_aggregator && !@supervisor.work_queue.empty? end |
#failed_results ⇒ Array
Get failed results
156 157 158 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 156 def failed_results @result_mutex.synchronize { @errors.dup } end |
#run ⇒ void
This method returns an undefined value.
Run the pool in batch mode and wait for completion
118 119 120 121 122 123 124 125 126 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 118 def run raise "Can only run in batch mode" if @continuous raise "Worker pool not started" unless @running @supervisor.run # Collect results collect_results end |
#shutdown(timeout: 30) ⇒ void
This method returns an undefined value.
Shutdown the worker pool
132 133 134 135 136 137 138 139 140 141 142 143 144 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 132 def shutdown(timeout: 30) return unless @running if @continuous @supervisor.stop @supervisor_thread&.join(timeout) end # Collect final results collect_results @running = false end |
#start ⇒ void
This method returns an undefined value.
Start the worker pool
54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 54 def start return if @running # Create Fractor supervisor with worker pool configuration @supervisor = Fractor::Supervisor.new( worker_pools: [ { worker_class: @worker_class, num_workers: @num_workers, }, ], continuous_mode: @continuous, ) # For continuous mode, set up work queue if @continuous @work_queue = Fractor::WorkQueue.new @work_queue.register_with_supervisor(@supervisor) end @running = true # Start supervisor in background thread for continuous mode if @continuous @supervisor_thread = Thread.new do @supervisor.run rescue StandardError => e @result_mutex.synchronize do @errors << { error: e, message: "Supervisor error: #{e.}" } end end else # For batch mode, don't start yet - wait for work items @supervisor.start_workers end end |
#stats ⇒ Hash
Get pool statistics
163 164 165 166 167 168 169 170 171 172 173 174 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 163 def stats return {} unless @supervisor { workers: @num_workers, running: @running, continuous: @continuous, results: @results.size, errors: @errors.size, total_processed: @results.size + @errors.size, } end |
#submit(work) ⇒ void
This method returns an undefined value.
Submit work item to the pool
95 96 97 98 99 100 101 102 103 104 105 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 95 def submit(work) raise "Worker pool not started" unless @running if @continuous # In continuous mode, add to work queue @work_queue << work else # In batch mode, add to supervisor @supervisor.add_work_item(work) end end |
#submit_batch(works) ⇒ void
This method returns an undefined value.
Submit multiple work items
111 112 113 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 111 def submit_batch(works) works.each { |work| submit(work) } end |
#successful_results ⇒ Array
Get successful results
149 150 151 |
# File 'lib/omnizip/parallel/worker_pool.rb', line 149 def successful_results @result_mutex.synchronize { @results.dup } end |