Skip to content

A job queue that keeps working when workers fall behind

Design a Distributed Job Queue

Built from Slack's account of the outage that made it rebuild its job queue: make enqueues safe when workers fall behind, keep one slow job type from starving the rest, and drain a backlog without causing the next outage.

Intermediate, about 45 minutes, 9 stages

The situation

A team messaging app does a lot of work outside web requests: push notifications, link previews, search indexing, billing events, exports. Web servers enqueue a job and return; workers pick jobs up and run them.

Today, jobs are pushed onto lists in a set of Redis clusters, and a fleet of workers pops them off. At peak the system handles about 33,000 jobs a second, 1.4 billion a day. It has worked for years, and engineers use it for everything.

Then a database slowed down. Workers that wrote to it slowed down too, jobs piled up, and Redis ran out of memory. Dequeuing a job also needed a little free memory, so when Redis filled up the queue could not drain at all, and every feature built on jobs stopped. Slack wrote about this outage, and the system they built afterwards, in "Scaling Slack's Job Queue". This investigation follows their reasoning.

What it has to do

Functional

  • Enqueue a job from any web request.
  • Run every job at least once, retrying failures.
  • Keep job types (and priorities) separate, with per-type throttling.
  • Show how far behind each job type is.

Non-functional

  • Enqueueing keeps working when workers are slow; jobs wait instead of failing.
  • Once enqueue returns, the job is not lost.
  • A backlog in one job type does not delay the others.
  • Interactive jobs (notifications) still run within seconds while batch work is backlogged.

Constraints and assumptions

  • About 1.4 billion jobs a day, peaking at about 33,000 a second.
  • Thousands of job handlers and the worker fleet talk to Redis; they cannot all be rewritten at once.
  • A Kafka cluster can be provisioned.
  • Jobs are small: a few kilobytes of JSON.
  • Most jobs finish in milliseconds; a few take seconds or minutes.
  • Jobs call downstream databases and services that have their own limits.

Interview questions it prepares you for

  • “Design a distributed job queue or task scheduler.”
  • “Design a background job system like Sidekiq, Celery or SQS.”
  • “Your queue is backing up. What do you do?”
  • “How do you guarantee a background job runs exactly once?”

Read and practise next

How Slack built it · Job queues under backlog, in their engineers' own words

Concepts to know first: Message queues, Backpressure and capacity.

Similar systems: Design a Video Processing Pipeline, Design a Notification System.