Class: Omnizip::Parallel::JobQueue
- Inherits:
-
Object
- Object
- Omnizip::Parallel::JobQueue
- Defined in:
- lib/omnizip/parallel/job_queue.rb
Overview
Thread-safe job queue for parallel compression/extraction
Manages a queue of compression or extraction jobs with priority support. Jobs are ordered by priority (large files first for better load balancing).
Defined Under Namespace
Classes: Job
Instance Attribute Summary collapse
-
#max_size ⇒ Integer
readonly
Maximum queue size.
-
#size ⇒ Integer
readonly
Current queue size.
Instance Method Summary collapse
-
#clear ⇒ Integer
Clear all jobs from queue.
-
#close ⇒ Object
Close the queue.
-
#closed? ⇒ Boolean
Check if queue is closed.
-
#empty? ⇒ Boolean
Check if queue is empty.
-
#initialize(max_size: 1000) ⇒ JobQueue
constructor
Initialize job queue.
-
#pop(timeout: nil) ⇒ Job?
Pop a job from the queue.
-
#pop_batch(count, timeout: nil) ⇒ Array<Job>
Pop multiple jobs in batch.
-
#push(file:, data: nil, size: 0, priority: :normal, metadata: {}) ⇒ Job
Push a job onto the queue.
-
#push_with_size(file:, size:, data: nil, metadata: {}) ⇒ Job
Push a job with automatic priority based on file size.
-
#stats ⇒ Hash
Get queue statistics.
Constructor Details
#initialize(max_size: 1000) ⇒ JobQueue
Initialize job queue
41 42 43 44 45 46 47 48 |
# File 'lib/omnizip/parallel/job_queue.rb', line 41 def initialize(max_size: 1000) @max_size = max_size @queue = [] @mutex = Mutex.new @cond = ConditionVariable.new @closed = false @size = 0 end |
Instance Attribute Details
#max_size ⇒ Integer (readonly)
Returns maximum queue size.
33 34 35 |
# File 'lib/omnizip/parallel/job_queue.rb', line 33 def max_size @max_size end |
#size ⇒ Integer (readonly)
Returns current queue size.
36 37 38 |
# File 'lib/omnizip/parallel/job_queue.rb', line 36 def size @size end |
Instance Method Details
#clear ⇒ Integer
Clear all jobs from queue
180 181 182 183 184 185 186 187 188 |
# File 'lib/omnizip/parallel/job_queue.rb', line 180 def clear @mutex.synchronize do count = @queue.size @queue.clear @size = 0 @cond.broadcast count end end |
#close ⇒ Object
Close the queue
No more jobs can be pushed after closing. Pending pops will return nil.
170 171 172 173 174 175 |
# File 'lib/omnizip/parallel/job_queue.rb', line 170 def close @mutex.synchronize do @closed = true @cond.broadcast # Wake up all waiting threads end end |
#closed? ⇒ Boolean
Check if queue is closed
162 163 164 |
# File 'lib/omnizip/parallel/job_queue.rb', line 162 def closed? @mutex.synchronize { @closed } end |
#empty? ⇒ Boolean
Check if queue is empty
155 156 157 |
# File 'lib/omnizip/parallel/job_queue.rb', line 155 def empty? @mutex.synchronize { @queue.empty? } end |
#pop(timeout: nil) ⇒ Job?
Pop a job from the queue
113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 |
# File 'lib/omnizip/parallel/job_queue.rb', line 113 def pop(timeout: nil) @mutex.synchronize do if timeout deadline = Time.now + timeout while @queue.empty? && !@closed remaining = deadline - Time.now return nil if remaining <= 0 @cond.wait(@mutex, remaining) end else @cond.wait(@mutex) while @queue.empty? && !@closed end return nil if @queue.empty? job = @queue.shift @size -= 1 @cond.signal # Signal waiting pushers job end end |
#pop_batch(count, timeout: nil) ⇒ Array<Job>
Pop multiple jobs in batch
141 142 143 144 145 146 147 148 149 150 |
# File 'lib/omnizip/parallel/job_queue.rb', line 141 def pop_batch(count, timeout: nil) jobs = [] count.times do job = pop(timeout: timeout) break unless job jobs << job end jobs end |
#push(file:, data: nil, size: 0, priority: :normal, metadata: {}) ⇒ Job
Push a job onto the queue
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 |
# File 'lib/omnizip/parallel/job_queue.rb', line 59 def push(file:, data: nil, size: 0, priority: :normal, metadata: {}) @mutex.synchronize do raise ClosedQueueError, "Queue is closed" if @closed # Wait if queue is full @cond.wait(@mutex) while @size >= @max_size && !@closed raise ClosedQueueError, "Queue is closed" if @closed job = Job.new( file: file, data: data, size: size, priority: priority, metadata: , ) @queue << job @size += 1 # Keep queue sorted by priority @queue.sort! @cond.signal job end end |
#push_with_size(file:, size:, data: nil, metadata: {}) ⇒ Job
Push a job with automatic priority based on file size
94 95 96 97 98 99 100 101 102 103 104 105 106 107 |
# File 'lib/omnizip/parallel/job_queue.rb', line 94 def push_with_size(file:, size:, data: nil, metadata: {}) # Determine priority based on size # Large files (>10MB) get high priority for better load balancing priority = if size > 10 * 1024 * 1024 :high elsif size > 1024 * 1024 :normal else :low end push(file: file, data: data, size: size, priority: priority, metadata: ) end |
#stats ⇒ Hash
Get queue statistics
193 194 195 196 197 198 199 200 201 202 203 |
# File 'lib/omnizip/parallel/job_queue.rb', line 193 def stats @mutex.synchronize do { size: @size, max_size: @max_size, closed: @closed, utilization: @max_size.zero? ? 0.0 : @size.to_f / @max_size, priority_counts: @queue.group_by(&:priority).transform_values(&:count), } end end |