Class: WorkQueue
- Inherits:
-
Object
show all
- Defined in:
- lib/workqueue.rb,
lib/workqueue/version.rb
Defined Under Namespace
Classes: ThreadsafeCounter
Instance Attribute Summary collapse
Class Method Summary
collapse
Instance Method Summary
collapse
Constructor Details
#initialize(init_queue = [], opts = {}, &job) ⇒ WorkQueue
Returns a new instance of WorkQueue.
17
18
19
20
21
22
23
24
25
|
# File 'lib/workqueue.rb', line 17
def initialize(init_queue=[], opts={}, &job)
@job = job
@queue = Queue.new
@mutex = Mutex.new
opts.each { |k, v| send(:"#{k}=", v) }
concat(init_queue)
end
|
Instance Attribute Details
#job ⇒ Object
Returns the value of attribute job.
16
17
18
|
# File 'lib/workqueue.rb', line 16
def job
@job
end
|
#queue ⇒ Object
Returns the value of attribute queue.
15
16
17
|
# File 'lib/workqueue.rb', line 15
def queue
@queue
end
|
#size ⇒ Object
28
29
30
|
# File 'lib/workqueue.rb', line 28
def size
@size ||= 2
end
|
Class Method Details
.version ⇒ Object
2
3
4
|
# File 'lib/workqueue/version.rb', line 2
def self.version
'0.1.2'
end
|
Instance Method Details
#abort! ⇒ Object
55
56
57
|
# File 'lib/workqueue.rb', line 55
def abort!
@mutex.synchronize { @aborted = true }
end
|
#concat(arr) ⇒ Object
74
75
76
|
# File 'lib/workqueue.rb', line 74
def concat(arr)
arr.each { |x| push(x) }
end
|
#join ⇒ Object
78
79
80
81
82
83
84
|
# File 'lib/workqueue.rb', line 78
def join
@mutex.synchronize { @joined = true }
workers.each(&:join)
self
end
|
#push(e) ⇒ Object
Also known as:
<<
67
68
69
70
71
|
# File 'lib/workqueue.rb', line 67
def push(e)
queue.push([e, cursor.incr])
self
end
|
#results ⇒ Object
86
87
88
89
|
# File 'lib/workqueue.rb', line 86
def results
join
aggregate
end
|
#run ⇒ Object
59
60
61
62
63
64
65
|
# File 'lib/workqueue.rb', line 59
def run
@workers = (1..size).map do
Thread.new { work! }
end
self
end
|
#work! ⇒ Object
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
|
# File 'lib/workqueue.rb', line 32
def work!
loop do
begin
payload, index = []
@mutex.synchronize {
unless @aborted or (@joined and queue.empty?)
payload, index = queue.shift
end
}
break if payload.nil?
aggregate[index] = job.call(payload)
rescue Exception
abort!
raise
end
end
end
|