Module: Resque::Plugins::FairShare
- Defined in:
- lib/resque/plugins/fair_share.rb
Overview
Constant Summary collapse
- ENQUEUE_SCRIPT =
File.read(File.('../../../lua/fair_enqueue.lua', __dir__))
Instance Method Summary collapse
- #after_enqueue_fair_share(*args) ⇒ Object
- #after_perform_fair_share(*args) ⇒ Object
- #before_perform_fair_share(*args) ⇒ Object
- #fair_share_on(key = nil, &block) ⇒ Object
- #on_failure_fair_share(_exception, *args) ⇒ Object
Instance Method Details
#after_enqueue_fair_share(*args) ⇒ Object
24 25 26 27 28 29 30 31 |
# File 'lib/resque/plugins/fair_share.rb', line 24 def after_enqueue_fair_share(*args) partition_value = extract_partition_value(args) redis.incr(redis_key(partition_value, 'pending')) emit_gauge(partition_value, 'pending') enqueue_to_sub_queue(partition_value, args) end |
#after_perform_fair_share(*args) ⇒ Object
53 54 55 56 57 58 59 |
# File 'lib/resque/plugins/fair_share.rb', line 53 def after_perform_fair_share(*args) partition_value = extract_partition_value(args) redis.decr(redis_key(partition_value, 'in_flight')) update_partition_score(partition_value) emit_gauge(partition_value, 'in_flight') end |
#before_perform_fair_share(*args) ⇒ Object
33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 |
# File 'lib/resque/plugins/fair_share.rb', line 33 def before_perform_fair_share(*args) partition_value = extract_partition_value(args) in_flight_key = redis_key(partition_value, 'in_flight') in_flight = redis.get(in_flight_key).to_i unless in_flight < Resque::FairShare.configuration.max_in_flight re_enqueue(partition_value, args) raise Resque::Job::DontPerform, "Max in-flight limit reached for partition #{partition_value}. " \ "Current: #{in_flight}, max: #{Resque::FairShare.configuration.max_in_flight}" end redis.incr(in_flight_key) update_partition_score(partition_value) emit_gauge(partition_value, 'in_flight') redis.decr(redis_key(partition_value, 'pending')) emit_gauge(partition_value, 'pending') end |
#fair_share_on(key = nil, &block) ⇒ Object
16 17 18 19 20 21 22 |
# File 'lib/resque/plugins/fair_share.rb', line 16 def fair_share_on(key = nil, &block) if block @_fair_share_resolver = block elsif key @_fair_share_key = key end end |
#on_failure_fair_share(_exception, *args) ⇒ Object
61 62 63 64 65 66 67 |
# File 'lib/resque/plugins/fair_share.rb', line 61 def on_failure_fair_share(_exception, *args) partition_value = extract_partition_value(args) redis.decr(redis_key(partition_value, 'in_flight')) update_partition_score(partition_value) emit_gauge(partition_value, 'in_flight') end |