ruby / job queues
I used a small Ruby and Postgres job queue in a Rack-based web framework:
- Each queue runs 1 job at a time.
- A worker takes jobs First In, First Out.
- A job is any object that responds to
Job.new(db).call, with optional argsJob.new(db).call(foo: 1, bar: "baz"). - The only dependencies are Ruby, Postgres,
and a custom DB wrapper around the
pgdriver.
Modest needs
I had about 20 queues. About 80% called third-party APIs with rate limits, such as GitHub, Discord, Slack, and Postmark. One job at a time was enough.
How
Create a jobs table in Postgres:
CREATE TABLE jobs (
id SERIAL,
queue text NOT NULL,
name text NOT NULL,
args jsonb DEFAULT '{}' NOT NULL,
status text DEFAULT 'pending'::text NOT NULL,
callsite text,
created_at timestamp DEFAULT now() NOT NULL,
started_at timestamp,
finished_at timestamp
);
Run one process:
bundle exec ruby queues/poll.rb
queues/poll.rb forks one child per worker:
require_relative "../lib/db"
require_relative "discord_worker"
require_relative "github_worker"
require_relative "postmark_worker"
require_relative "slack_worker"
$stdout.sync = true
module Queues
WORKERS = [
Queues::DiscordWorker,
Queues::GithubWorker,
Queues::PostmarkWorker,
Queues::SlackWorker
].freeze
end
# Ensure all workers implement the interface.
Queues::WORKERS.each(&:validate!)
# Ensure queues are only worked on by one worker.
dup_queues = Queues::WORKERS.map(&:queue).tally.select { |_, v| v > 1 }.keys
if dup_queues.any?
raise "duplicate queues: #{dup_queues.join(", ")}"
end
children = Queues::WORKERS.map do |worker|
fork do
worker.new(DB.new).poll
rescue SignalException
end
end
begin
children.each { |pid| Process.wait(pid) }
rescue SignalException => sig
if Signal.list.values_at("HUP", "INT", "KILL", "QUIT", "TERM").include?(sig.signo)
children.each do |pid|
Process.kill("TERM", pid)
rescue Errno::ESRCH
end
# Give children time to finish in-flight jobs
deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 5
remaining = children.dup
while remaining.any? && Process.clock_gettime(Process::CLOCK_MONOTONIC) < deadline
remaining.reject! do |pid|
Process.wait(pid, Process::WNOHANG)
rescue Errno::ECHILD
true
end
if remaining.any?
sleep 0.5
end
end
remaining.each do |pid|
Process.kill("KILL", pid)
rescue Errno::ESRCH
end
remaining.each do |pid|
Process.wait(pid)
rescue Errno::ECHILD
end
end
end
Create a base worker class:
# queues/poll_worker.rb
module Queues
class PollWorker
class << self
attr_accessor :queue, :jobs
end
def self.validate!
if queue.to_s.strip.empty?
raise NotImplementedError, "#{name} does not specify a queue"
end
if @jobs.nil? || @jobs.empty?
raise NotImplementedError, "#{name} does not define any jobs"
end
end
attr_reader :db
def initialize(db)
@db = db
end
def poll
puts "queue=#{queue} poll=#{poll_interval}s"
loop do
sleep poll_interval
pending_jobs.each do |job|
status = "err: interrupted"
latency = 0
result = db.exec(<<~SQL, [job["id"]]).first
UPDATE
jobs
SET
started_at = now(),
status = 'started'
WHERE
id = $1
RETURNING
EXTRACT(EPOCH FROM now() - created_at) AS latency
SQL
latency =
if result && result["latency"]
result["latency"].round(2).to_f
else
0
end
status = work(job_name: job["name"], job_args: job["args"])
rescue => err
status = "err: #{err}"
Sentry.capture_exception(err)
ensure
if job && job["id"]
result = db.exec(<<~SQL, [status, job["id"]]).first
UPDATE
jobs
SET
finished_at = now(),
status = $1
WHERE
id = $2
RETURNING
EXTRACT(EPOCH FROM now() - started_at) AS elapsed
SQL
elapsed =
if result && result["elapsed"]
result["elapsed"].round(2).to_f
else
0
end
puts %(queue=#{queue} job=#{job["name"]} id=#{job["id"]} status="#{status}" latency=#{latency}s duration=#{elapsed}s)
min_job_time = 1.0 / max_jobs_per_second_for_status(status)
sleep [min_job_time - elapsed, 0].max
end
end
end
end
private def pending_jobs
db.exec(<<~SQL, [queue])
SELECT
id,
name,
args
FROM
jobs
WHERE
queue = $1
AND started_at IS NULL
AND status = 'pending'
ORDER BY
created_at ASC
SQL
end
private def poll_interval
10
end
private def max_jobs_per_second_for_status(status)
max_jobs_per_second
end
private def max_jobs_per_second
Float::INFINITY
end
private def queue
self.class.queue
end
private def work(job_name:, job_args:)
worker = self.class.jobs.find { |job| job.name == job_name }
if !worker
msg = "unknown job `#{job_name}` for queue `#{queue}`"
Sentry.capture_message(msg, extra: {job_args: job_args})
return "err: #{msg}"
end
if worker.instance_method(:call).arity == 0
worker.new(db).call
else
worker.new(db).call(**job_args.transform_keys(&:to_sym))
end
end
end
end
Implement a specific worker:
# queues/github_worker.rb
require_relative "poll_worker"
require_relative "../lib/github/job_one"
require_relative "../lib/github/job_two"
module Queues
class GithubWorker < PollWorker
@queue = "github"
@jobs = [Github::JobOne, Github::JobTwo]
# Override max_jobs_per_second_for_status to slow down
# further when a job returns a rate-limit error.
private def max_jobs_per_second_for_status(status)
if status&.include?("rate limited")
1 / 300.0 # wait 5 minutes on rate limit
else
10
end
end
end
end
A worker with no external API skips the rate limit:
# queues/no_throttle_worker.rb
module Queues
class NoThrottleWorker < PollWorker
@queue = "no_throttle"
@jobs = [Caches::Refresh, DataReviews::CleanUp, ...]
end
end
Enqueuing jobs
A helper inserts jobs:
# lib/jobs/insert.rb
module Jobs
class Insert
attr_reader :db
def initialize(db)
@db = db
end
def call(queue:, name:, args: {}, args_params: [])
loc = caller_locations(1, 1).first
callsite = "#{loc.path}:#{loc.lineno}"
param_offset = args_params.size
case args
when Hash
params = args_params + [queue, name, callsite, args]
sql = "SELECT $#{params.size} AS args"
when Array
params = args_params + [queue, name, callsite, args.to_json]
sql = "SELECT json_array_elements($#{params.size}) AS args"
when String
params = args_params + [queue, name, callsite]
sql = args
else
raise ArgumentError, "args must be an array, hash or sql string."
end
db.exec(<<~SQL, params)
WITH data AS (
#{sql}
)
INSERT INTO jobs (
queue,
name,
callsite,
args
)
SELECT
$#{param_offset + 1},
$#{param_offset + 2},
$#{param_offset + 3},
data.args::jsonb
FROM
data
ON CONFLICT DO NOTHING
RETURNING
id
SQL
end
end
end
Enqueue jobs:
require_relative "lib/db"
require_relative "lib/jobs/insert"
i = Jobs::Insert.new(DB.pool)
# Single job
i.call(
queue: "github",
name: "JobOne",
args: {
company_id: 42
}
)
# Multiple jobs from SQL
i.call(
queue: "github",
name: "JobOne",
args: <<~SQL
SELECT
jsonb_build_object('company_id', id) AS args
FROM
companies
WHERE
status = 'active'
SQL
)
Scheduling jobs
A Clock process inserts recurring jobs into the queue. I did not use cron, which keeps the schedule apart from the application code.
Run one Clock process:
# schedule/clock.rb
require_relative "../lib/db"
require_relative "../lib/jobs/insert"
require_relative "../lib/calendar"
module Schedule
JOBS = [
# Daily at midnight
{
queue: "discord",
name: "Discord::Ingest",
at?: proc { |t| t.min == 30 && t.hour == 0 }
},
# Daily cleanup at 8:30am
{
queue: "no_throttle",
name: "Queues::CleanUp",
at?: proc { |t| t.min == 30 && t.hour == 8 }
},
# Weekly on Tuesdays at 3am
{
queue: "github",
name: "Github::IngestStars",
at?: proc { |t| t.min == 0 && t.hour == 3 && t.tuesday? }
}
]
class Clock
attr_reader :db
def initialize(db)
@db = db
end
def tick(seconds:)
i = Jobs::Insert.new(db)
loop do
Schedule::JOBS.each do |job|
if job.fetch(:at?).call(Time.now.utc)
puts "insert #{job.fetch(:name)}"
i.call(
queue: job.fetch(:queue),
name: job.fetch(:name)
)
end
end
sleep(seconds)
end
end
end
end
if $0 == __FILE__
$stdout.sync = true
Schedule::Clock.new(DB.new).tick(seconds: 60)
end
Each job defines when it runs with a proc that receives a UTC
Time:
proc { |t| t.min % 15 == 0 } # every 15 minutes
proc { |t| t.min == 0 && t.hour == 3 } # daily at 3am
proc { |t| t.min == 0 && t.hour == 3 && t.tuesday? } # weekly
proc { |t| t.to_date == Date.new(t.year, t.month, -1) } # last of month
proc { |t| t.min == 0 && t.hour == 12 && !Calendar.holiday?(t) }
If the Clock crashes, the process supervisor restarts it. A job missed during downtime runs at the next matching time.
Maintenance
A monthly job deletes old rows. The rows that remain let a job check for recent work before it calls a paid API:
# lib/queues/clean_up.rb
module Queues
class CleanUp
attr_reader :db
def initialize(db)
@db = db
end
def call
db.exec <<~SQL
DELETE FROM jobs
WHERE created_at < now() - '1 month'::interval
SQL
"ok"
end
end
end