Module: Resque::Plugins::FairShare

Defined in:
lib/resque/plugins/fair_share.rb

Overview

Resque job hooks for fair-share scheduling. Extend this module in your job class to track per-partition pending and in-flight counters in Redis.

class MyJob
extend Resque::Plugins::FairShare
@queue = :default
fair_share_on :account_id
def self.perform(params); end
end

Constant Summary collapse

ENQUEUE_SCRIPT =
File.read(File.expand_path('../../../lua/fair_enqueue.lua', __dir__))

Instance Method Summary collapse

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