diff --git a/lib/resque/scheduler/delaying_extensions.rb b/lib/resque/scheduler/delaying_extensions.rb index 06093118..9af21dd2 100644 --- a/lib/resque/scheduler/delaying_extensions.rb +++ b/lib/resque/scheduler/delaying_extensions.rb @@ -261,12 +261,18 @@ def remove_delayed_job_from_timestamp(timestamp, klass, *args) count end + # Counts all jobs across all delayed timestamps. + # Uses pipelined LLEN calls in batches to avoid blocking the Redis thread + # for too long while being significantly faster than sequential calls. def count_all_scheduled_jobs - total_jobs = 0 - Array(redis.zrange(:delayed_queue_schedule, 0, -1)).each do |ts| - total_jobs += redis.llen("delayed:#{ts}").to_i + timestamps = Array(redis.zrange(:delayed_queue_schedule, 0, -1)) + return 0 if timestamps.empty? + + timestamps.each_slice(10_000).sum do |batch| + redis.pipelined do |pipeline| + batch.each { |ts| pipeline.llen("delayed:#{ts}") } + end.sum end - total_jobs end # Discover if a job has been delayed. diff --git a/test/delayed_queue_test.rb b/test/delayed_queue_test.rb index abd4c06d..01afe8ee 100644 --- a/test/delayed_queue_test.rb +++ b/test/delayed_queue_test.rb @@ -1100,6 +1100,10 @@ end end + test 'count_all_scheduled_jobs returns 0 when no jobs are scheduled' do + assert_equal(0, Resque.count_all_scheduled_jobs) + end + test 'delayed?' do Resque.enqueue_at Time.now + 1, SomeIvarJob Resque.enqueue_at Time.now + 1, SomeIvarJob, id: 1