A worker fleet at scale is not “the same worker, times a hundred.” It is a set of isolated pools, each sized against the bottleneck it actually contends for, scaled on how long work has been waiting rather than how much work is waiting, with explicit lease and drain semantics so that adding or removing a worker is never a data-loss event. The three decisions that determine whether a fleet holds up — how you partition it, what signal you scale on, and what happens to in-flight work when a worker goes away — all get made before you configure an autoscaler, and none of them are the autoscaler’s job.
I get called into fleets that are already “scaled.” They have fifty workers, an autoscaling policy, a dashboard with a queue-depth graph. And they still page someone at 3am because a single slow downstream API turned every worker in the fleet into a thread blocked on a socket, and the password-reset emails that share the queue stopped going out. The workers weren’t the problem. The architecture was one pool.
The default architecture, and exactly where it breaks
Nearly every fleet starts in the same place, and for good reason — it’s the right starting point:
One queue. One worker deployment. Every job type goes in the queue, any worker takes the next job. Scale the deployment on queue depth.
This works, genuinely, for a long time. It breaks in a specific and predictable order:
First, head-of-line starvation across job classes. A batch of 200,000 nightly re-index jobs lands in the queue. Your password reset email is job number 200,001. Every worker is busy with re-index work for the next forty minutes. Nothing is broken, no error fires, and your users are locked out. The queue was FIFO and FIFO was never what you wanted — you wanted two independent latency guarantees, and a single queue can only offer one.
Second, coupled failure domains. One job type calls a third-party enrichment API. That API starts taking 30 seconds instead of 300ms. Your workers are IO-bound, so they sit there holding a slot. Within a minute every worker in the fleet is blocked on that one dependency, and every unrelated job type is now down. You have built a system where any one dependency’s latency is every job’s latency.
Third, the scaling signal stops meaning anything. Queue depth is a count of heterogeneous items. Ten thousand jobs might be four seconds of work or nine hours of work. When one number covers both, the autoscaler is guessing, and it will guess wrong in the direction that costs you money or your SLO.
Each of these has the same fix underneath it: stop treating the fleet as one resource.
Partition by failure domain, not by job name
The instinct when you split a queue is to split it by what the job is — emails, reports, webhooks. That’s a reasonable first cut but it’s the wrong axis. The axis that matters is: what does this job contend for, and what takes it down?
Two jobs belong in the same pool when they share a bottleneck and a failure mode. They belong in different pools when a failure in one should not be able to consume the capacity of the other.
In practice that gives me a small number of pools, usually along these lines:
- Interactive / latency-bound. Work a human or an API caller is waiting on. Small, fast, tightly capped. This pool exists so that nothing bulk can ever occupy it. It is often over-provisioned on purpose, because idle capacity here is cheap and a queue here is a user-visible outage.
- Bulk / throughput-bound. Re-indexing, backfills, nightly rollups. Deep queue, high latency tolerance, scales aggressively and scales to near-zero. Nobody is waiting; the SLO is “finishes by 6am,” not “finishes in 2 seconds.”
- One pool per hostile dependency. Any job whose runtime is dominated by a third party I do not control gets its own pool. This is the single highest-value split in most fleets. It converts “the enrichment API is slow” from a fleet-wide outage into a bounded, visible backlog on one pool, which is a Tuesday, not an incident.
- Poison-prone / unbounded-input. Work that parses user-supplied files, renders untrusted documents, or otherwise has a real chance of hanging or OOMing on a single input. Isolate it so that a memory bomb kills one pool’s workers and not the fleet.
The rule I actually apply: if a job class can saturate the fleet for longer than another job class’s latency budget, they don’t share a pool.
The cost of this is real and worth stating. Every pool is separate capacity, separate idle time, separate config, and a separate thing to monitor. Ten pools is usually a smell — you’ve partitioned by job name and you now have an operations surface nobody can hold in their head. Three to five is typical for a fleet that’s genuinely at scale. Start with two: interactive and bulk. Add a pool when you can name the incident it would have prevented.
A weaker version worth knowing about: if your queue broker supports per-consumer-group concurrency limits (or you enforce them yourself with a semaphore keyed on job class), you can get isolation within one pool without separate deployments. It’s less robust — a process-level failure like OOM still takes shared workers down — but it’s a genuinely good intermediate step, and much cheaper to run.
Scale on backlog age, not queue depth
This is the change that surprises people most, and it’s the one I’d make first.
Queue depth is unitless. “There are 4,000 jobs waiting” tells you nothing about whether you’re in trouble, because it doesn’t contain the service rate. Four thousand 50ms jobs on twenty workers is a two-minute blip. Four thousand 20-second jobs is over five hours of backlog. The same number, the same alert threshold, wildly different situations — and any threshold you pick will be wrong for one of them.
Backlog age has units, and they’re the units of your SLO. “The oldest job in this queue has been waiting 4 minutes” is directly comparable to “no job in this pool should wait more than 5 minutes.” It automatically accounts for service time, worker count, worker health, and arrival rate, because all of those are already baked into how long the front of the queue has been sitting there. Most brokers expose it: ApproximateAgeOfOldestMessage on SQS, consumer lag translated to time on Kafka, or a simple now() - min(enqueued_at) on a database-backed queue.
So the autoscaling policy I want is: scale to keep the oldest-message age under the pool’s latency target. That’s a policy you can explain to a product owner, and it’s per-pool, which is exactly why the partitioning came first.
For sizing the actual target — how many workers to ask for — Little’s Law gives you the honest arithmetic. To keep up with arrivals and burn down an existing backlog B within T seconds:
import math
def desired_concurrency(
arrival_rate: float, # jobs/sec arriving, measured over a recent window
service_time: float, # mean seconds of work per job, measured
backlog: int, # jobs currently waiting
drain_target: float, # seconds allowed to clear that backlog
utilization: float = 0.7, # never plan for 100% — queues explode near saturation
) -> int:
"""Concurrency slots needed to keep up and clear the backlog in drain_target."""
required_throughput = arrival_rate + (backlog / drain_target) # jobs/sec
slots = (required_throughput * service_time) / utilization
return max(1, math.ceil(slots))
def desired_workers(slots: int, slots_per_worker: int) -> int:
return max(1, math.ceil(slots / slots_per_worker))
# 12 jobs/sec arriving, 800ms each, 9,000 backlogged, clear it in 5 minutes:
# required_throughput = 12 + 30 = 42 jobs/sec
# slots = 42 * 0.8 / 0.7 = 48
# with 4 slots per worker -> 12 workers
Two things in there are not decoration. The utilization divisor is the one people drop, and dropping it is why fleets sized “exactly right” fall over: as utilization approaches 1.0, queueing delay goes to infinity, so planning for 100% busy workers means planning for unbounded latency. Seventy percent is a reasonable default. And service_time must be measured, per pool, from real job durations — the moment you hardcode it, the model silently drifts away from reality and you’re back to guessing.
Two guardrails that matter more than the formula:
- Asymmetric scaling rates. Scale out fast, scale in slowly. Scaling out late costs you SLO; scaling in early costs you a re-scale-out three minutes later, plus cold starts, plus flapping. I typically allow scale-out every evaluation period and scale-in on a much longer cooldown.
- A hard ceiling per pool, set by the downstream, not by your budget. Which is the next section.
Size each pool against its actual bottleneck
The most expensive scaling mistake I see is a fleet that scales up correctly and takes down the thing it depends on.
Your worker count is not a free variable. It is bounded by whichever of these binds first:
- CPU. For genuinely CPU-bound work, concurrency per worker above the core count buys you nothing but context switching and memory. One slot per core, maybe minus one for the runtime.
- Memory. Concurrency × peak-RSS-per-job must fit with headroom. This is the constraint that kills workers rather than slowing them, and peak matters, not mean — one 900MB document in a stream of 20MB ones is what OOMs the box.
- Connection pools. Every concurrent job that touches Postgres holds a connection. Fleet-wide concurrency above your database’s
max_connectionsdoesn’t produce more throughput; it produces connection errors and a thundering herd of retries. If you have a pooler in front, the bound moves, but it doesn’t disappear. - The downstream’s rate limit or capacity. This is the one that belongs to the dependency, not to you. If an API allows 50 requests/second, then the pool that calls it has a fixed concurrency ceiling regardless of how deep its queue gets — and the correct behaviour when the backlog grows is to let it grow, visibly, not to add workers.
That last point is worth stating plainly because it inverts the usual instinct: for a downstream-bound pool, the right response to a backlog is to alert, not to scale. An autoscaler that adds workers against a saturated dependency is a load generator pointed at something already failing. Cap the pool at the concurrency the downstream can actually absorb, enforce it as a fleet-wide limit rather than a per-worker one, and let the backlog age be the signal that a human needs to make a decision. The retry logic underneath this matters as much as the cap — a fleet with no shared ceiling on retries will manufacture the same overload from the other direction, which I’ve written about in the retry with exponential backoff and jitter guide.
Acquire work with a lease, not a delete
Here is the failure that quietly loses data in fleets that otherwise look well-run.
A worker pops a job off the queue and the job is now gone from the queue. The worker starts processing. The worker’s instance is reclaimed mid-job — spot interruption, OOM kill, node drain, a deploy. The job existed only in that process’s memory. It is gone, silently, and the only evidence is a customer asking where their report went.
The fix is that dequeuing must be a lease, not a removal: the job becomes invisible to other workers for a bounded time, and is only deleted after the work is confirmed complete. If the worker dies, the lease expires and the job becomes visible again. SQS calls this a visibility timeout; a Postgres-backed queue does it with a locked_until column and a FOR UPDATE SKIP LOCKED claim.
The trap is that the lease duration must exceed the job’s actual runtime, and for any job with a variable runtime, it eventually won’t. When the lease expires mid-execution, a second worker picks up the same job and now you have two workers doing the same work concurrently — which is either a duplicate charge, a duplicate email, or a corrupted counter, depending on what the job does. So the lease gets extended by a heartbeat while the work is genuinely in progress:
import threading
LEASE_SECONDS = 60
HEARTBEAT_SECONDS = 20 # comfortably under LEASE_SECONDS
MAX_LEASE_EXTENSIONS = 30 # ~10 min ceiling: past this, the job is stuck, not slow
def run_with_lease(queue, job) -> None:
"""Process a job, extending its lease while it's genuinely still running."""
done = threading.Event()
extensions = 0
def heartbeat() -> None:
nonlocal extensions
while not done.wait(HEARTBEAT_SECONDS):
if extensions >= MAX_LEASE_EXTENSIONS:
# Refuse to hold the lease forever. Let it expire so the job
# is redelivered and eventually dead-lettered, rather than
# pinning a worker on something that will never finish.
return
queue.extend_lease(job.receipt, LEASE_SECONDS)
extensions += 1
beat = threading.Thread(target=heartbeat, daemon=True)
beat.start()
try:
handle(job) # must be idempotent — see below
queue.delete(job.receipt) # delete only after the work is durable
finally:
done.set()
beat.join(timeout=1)
Note the ordering: the job is deleted after the handler returns, not before it starts. That ordering is what makes the queue at-least-once, and at-least-once is what makes crash recovery possible. It also means duplicate execution is a normal event rather than an exception — a worker that finished the work and died before the delete will see that job again. Every handler in the fleet has to be safe to run twice, which is a design constraint on the job itself, not something the fleet can paper over. That’s the whole subject of how to design an idempotent job queue, and it’s a prerequisite for everything in this guide, not an optional extra.
The MAX_LEASE_EXTENSIONS ceiling is there because an unbounded heartbeat turns a hung job into a permanently occupied worker slot. A job that has held a lease for ten minutes when the p99 is forty seconds isn’t slow, it’s stuck, and the fleet is better off letting it fail into the retry-and-dead-letter path than dedicating a worker to it indefinitely.
Draining: the deploy path most fleets get wrong
You deploy ten times a day. Each deploy terminates every worker in the fleet. That means the shutdown path executes far more often than any of the failure paths you spent time on, and it’s usually the least-tested code in the system.
What should happen when a worker receives SIGTERM:
- Stop accepting new work immediately. Break the poll loop first, before anything else.
- Let in-flight jobs finish, up to a bounded grace period. Keep the heartbeat running during the drain — an in-flight job still needs its lease extended.
- For anything still running at the end of the grace period, release the lease explicitly rather than just exiting. An explicit release makes the job immediately visible to another worker; letting the lease expire means waiting out the full timeout with the job invisible and nobody working on it. That difference is the gap between a seamless deploy and a several-minute latency spike on every deploy.
import signal
draining = False
def _on_sigterm(signum, frame) -> None:
global draining
draining = True # only set a flag; never do teardown work in a handler
signal.signal(signal.SIGTERM, _on_sigterm)
def poll_loop(queue) -> None:
while not draining:
job = queue.receive(wait_seconds=20) # long-poll, so we notice SIGTERM promptly
if job is None:
continue
run_with_lease(queue, job)
# Loop exited: in-flight work has returned or been released by run_with_lease.
Two operational details make or break this. The grace period your orchestrator gives you must be longer than your longest expected job — in Kubernetes that’s terminationGracePeriodSeconds, and the default of 30 seconds is shorter than a great many jobs, so it silently SIGKILLs work in the middle every single deploy. And the long-poll wait needs to be short enough that a draining worker notices the flag quickly; a 20-second receive means up to 20 seconds of a worker sitting idle before it even starts draining.
If you run on spot or preemptible capacity, the same drain path handles the interruption notice — you get roughly two minutes of warning, which is only useful if SIGTERM already does the right thing.
The failure modes that only appear at fleet scale
These don’t show up with three workers. They show up at eighty.
Retry storms. Every worker retrying a failing dependency simultaneously generates more load during the outage than during normal operation, which prevents recovery. Backoff bounds one worker; only a fleet-wide retry budget bounds the fleet.
Scale-out thundering herd. The autoscaler adds forty workers at once. All forty cold-start, all forty open database connections, all forty hit the cache cold, and the burst of connection setup is itself the load spike that trips the dependency. Ramp scale-out in steps rather than jumping to target, and stagger startup with a small random delay.
The autoscaler feedback loop. Backlog grows because a downstream is degraded. The autoscaler adds workers. More workers means more pressure on the degraded downstream. Service time rises. Backlog grows further. The autoscaler adds more workers. The scaling policy is now actively making the incident worse, and it will keep doing so until it hits the ceiling — which is precisely why every pool needs a ceiling derived from its bottleneck.
Poison messages taking a pool down. One job that reliably OOMs the worker gets redelivered after each crash, killing worker after worker. Without a delivery-count cap routing it out of the queue, a single bad input becomes a rolling fleet outage.
Silent partial capacity. Twenty of your eighty workers are alive to the orchestrator’s health check but wedged — a deadlocked thread pool, an exhausted connection pool. Throughput is at 75% and nothing is red. Health checks that only prove the process is running will not catch this; the check has to assert that the worker has completed work recently.
What to actually watch
Fleet dashboards are usually full of the wrong graph. Per pool, the ones that have earned their place:
- Age of the oldest queued job — the SLO-shaped number, and the one to alert on.
- Throughput vs. arrival rate — whether you’re gaining or losing ground, which depth alone can’t tell you.
- Lease expiries / redeliveries — a rising count means jobs are dying mid-flight or leases are too short. This is the metric that catches silent job loss, and almost nobody has it.
- Effective concurrency — slots actually doing work vs. slots provisioned. The gap is your wedged workers.
- Dead-letter arrival rate, with an alert on any sustained non-zero rate.
Notice what isn’t there: CPU utilization as a primary signal. It’s useful for sizing a single worker and nearly useless for judging a fleet whose work is mostly IO-bound. The distinction between the numbers worth waking someone up for and the numbers worth looking at during an investigation is the whole argument in monitoring vs. alerting, and it applies with particular force here, where it’s easy to build twenty graphs that never tell you the fleet is behind. Getting this instrumentation right is most of what I do on a monitoring and operations engagement.
When you don’t need any of this
Most systems don’t. If your queue drains to empty within seconds most of the day, every job class has effectively the same latency requirement, and you can restart the fleet without anyone noticing — one pool, a fixed worker count, and a good alert on queue age is a complete and correct architecture. Adding pools and an autoscaler to that system buys you nothing but configuration surface and more ways to be paged.
The signals that you’ve actually outgrown it: one job class’s backlog is visibly delaying another; a single dependency’s slowness becomes a fleet-wide event; your worker count is either wasting money at night or short during peak; or deploys cause a latency spike. Until at least one of those is true, the simple version is genuinely the better engineering. Fleet architecture is a response to specific, observed failures, and if you can’t name yours yet, the honest move is to wait for one.
The broader failure-classification model behind all of this — how I think about which failures are transient, environmental, semantic, or terminal — is in worker fleets in practice, and the rest of the reliability patterns these pools depend on live in the reliable systems guide.
FAQ
How many workers should I run?
Measure rather than guess: (arrival_rate + backlog/drain_target) × service_time / 0.7 gives you concurrency slots, and slots divided by per-worker concurrency gives you workers. Then check that number against your connection pool, your memory ceiling, and your downstream’s rate limit, and take the smallest. If those three checks don’t produce a lower number than your formula does, you probably haven’t found your real bottleneck yet.
Should each worker run one job at a time or many?
Many, if the work is IO-bound — a worker blocked on a network call is wasting a whole process. One at a time, or close to it, if the work is CPU-bound or memory-heavy, where extra concurrency only adds contention and OOM risk. The mistake is applying one answer fleet-wide when different pools have different profiles; per-pool concurrency is the point of having pools.
Is queue depth ever a useful metric?
For capacity planning and cost forecasting, yes. For alerting and autoscaling, no — use age. Depth is a fine secondary graph; it just can’t be the number that decides anything, because it has no units you can compare to an SLO.
Do I need Kubernetes for this?
No. Everything here — pool isolation, age-based scaling, leases, drain-on-SIGTERM — works on plain VMs with an autoscaling group, or on any managed container platform. Kubernetes gives you convenient primitives (and KEDA makes queue-age scaling straightforward), but the architecture is what matters and none of it is Kubernetes-specific.
What about exactly-once processing so I don’t need idempotency?
You can’t buy your way out of this one. Exactly-once delivery across a network is not achievable in the general case; what real systems provide is at-least-once delivery plus idempotent handlers, or effectively-once semantics within one system’s boundary that stop applying the moment you make an external call. Build the handlers to be safe on redelivery.
How do I stop one tenant from monopolizing the fleet?
Pools split by failure domain, not by tenant, so this needs a separate mechanism: a per-tenant concurrency cap enforced fleet-wide, or per-tenant queues with weighted round-robin consumption. Without one, a single customer’s bulk import is indistinguishable from a fleet-wide outage for everyone else.