Autoscaling BullMQ Workers on ECS

ECS scales services on CloudWatch metrics, and CPU is the metric most teams start with — which is the wrong signal for queue workers whose jobs mostly wait on I/O. This guide scales BullMQ workers on ECS from queue backlog instead, as part of Horizontal Worker Scaling in Backend Frameworks & Worker Scaling.

Problem Statement

A Node.js service runs BullMQ workers as an ECS service on Fargate, scaled by average CPU with a 60% target. The jobs call external APIs and spend most of their time waiting, so CPU stays around 20% even when 200,000 jobs are waiting; the service never scales out, and nightly batch loads take six hours to drain. When someone lowered the CPU target to 15%, the service oscillated between 2 and 30 tasks every few minutes. Scale-in also killed tasks mid-job, causing stalled jobs and duplicates. You want scaling driven by backlog per task, calm behaviour without oscillation, and scale-in that never interrupts running jobs.

Prerequisites

  • An ECS service running BullMQ workers (Fargate or EC2 launch type) with Application Auto Scaling available.
  • A component that can read BullMQ queue counts and publish to CloudWatch (a small Lambda on a schedule, or a sidecar).
  • Graceful shutdown in the worker (worker.close() on SIGTERM) and an ECS stopTimeout longer than typical jobs.
  • Known concurrency per task (e.g. 20 jobs in parallel per worker task).

Step 1 — Publish Backlog per Task as a CloudWatch Metric

Target tracking works best with a metric that should stay constant as the service scales: backlog per task — waiting jobs divided by running tasks. Publish it every minute.

// backlog-metric.ts — Lambda on a 1-minute EventBridge schedule
import { Queue } from "bullmq";
import { CloudWatchClient, PutMetricDataCommand } from "@aws-sdk/client-cloudwatch";
import { ECSClient, DescribeServicesCommand } from "@aws-sdk/client-ecs";

const queue = new Queue("sync", { connection: { host: process.env.REDIS_HOST!, port: 6379 } });
const cw = new CloudWatchClient({});
const ecs = new ECSClient({});

export async function handler() {
  const counts = await queue.getJobCounts("waiting", "prioritized", "delayed");
  const backlog = counts.waiting + counts.prioritized;           // due now; delayed excluded
  const svc = await ecs.send(new DescribeServicesCommand({ cluster: "jobs", services: ["sync-worker"] }));
  const running = Math.max(svc.services?.[0]?.runningCount ?? 1, 1);
  await cw.send(new PutMetricDataCommand({
    Namespace: "Jobs",
    MetricData: [
      { MetricName: "Backlog", Dimensions: [{ Name: "Queue", Value: "sync" }], Value: backlog, Unit: "Count" },
      { MetricName: "BacklogPerTask", Dimensions: [{ Name: "Queue", Value: "sync" }],
        Value: backlog / running, Unit: "Count" },
    ],
  }));
}

Delayed jobs are excluded because they are not work that can be done now; counting them scales out for jobs scheduled hours ahead.

CPU is the wrong signal for I/O-bound jobs During a nightly batch, the waiting backlog rises to two hundred thousand jobs. Average CPU of the worker tasks stays around twenty percent because jobs wait on external APIs, so a sixty-percent CPU target never triggers scale-out. Backlog per task rises sharply and is the signal that should drive scaling. Nightly batch: two signals CPU target 60% CPU ~20%: never scales backlog per task 22:00 04:00

Step 2 — Choose the Target from Job Duration and Wait Objective

The target backlog per task is how many waiting jobs one task should have queued in front of it. It follows from how fast a task drains jobs and how long jobs may wait.

task throughput = concurrency / avg job duration = 20 / 2 s = 10 jobs/s per task
wait objective  = 60 s
target backlog per task = throughput x objective = 10 x 60 = 600

With a target of 600, a backlog of 200,000 asks for about 333 tasks — so the maximum task count, not the target, will be the binding limit during big batches. Set the maximum from what downstream dependencies can take, as discussed in preventing retry storms after an outage. The throughput arithmetic comes from calculating worker count with Little's Law.

Step 3 — Configure Target Tracking with Sensible Bounds

Register the service as a scalable target and attach a target-tracking policy on the custom metric.

resource "aws_appautoscaling_target" "sync" {
  service_namespace  = "ecs"
  resource_id        = "service/jobs/sync-worker"
  scalable_dimension = "ecs:service:DesiredCount"
  min_capacity       = 2
  max_capacity       = 60            # bounded by partner API capacity, not by backlog
}

resource "aws_appautoscaling_policy" "sync_backlog" {
  name               = "sync-backlog-per-task"
  policy_type        = "TargetTrackingScaling"
  service_namespace  = aws_appautoscaling_target.sync.service_namespace
  resource_id        = aws_appautoscaling_target.sync.resource_id
  scalable_dimension = aws_appautoscaling_target.sync.scalable_dimension

  target_tracking_scaling_policy_configuration {
    target_value       = 600
    scale_out_cooldown = 60          # react quickly to growth
    scale_in_cooldown  = 600         # shrink slowly to avoid flapping
    customized_metric_specification {
      metric_name = "BacklogPerTask"
      namespace   = "Jobs"
      statistic   = "Average"
      dimensions { name = "Queue"  value = "sync" }
    }
  }
}

The asymmetric cool-downs are what stop the oscillation from the problem statement: scale out within a minute, scale in only after ten stable minutes. A metric that divides by running tasks is self-correcting as tasks are added, which also damps oscillation compared with raw backlog.

Calm scaling with asymmetric cool-downs With a low CPU target, the service oscillated between 2 and 30 tasks every few minutes. With backlog-per-task target tracking, a one-minute scale-out cool-down, and a ten-minute scale-in cool-down, the task count rises to the maximum of 60 during the batch, holds while the backlog drains, and steps down gradually afterwards. Running tasks over the batch window CPU 15%: flapping backlog/task: max 60, then slow scale-in before after

Step 4 — Protect Busy Tasks During Scale-In

When ECS scales in, it stops tasks with SIGTERM followed by SIGKILL after stopTimeout (default 30 s, maximum 120 s on Fargate). Two measures prevent interrupted jobs.

// task definition container settings
{ "name": "worker", "stopTimeout": 120 }
// worker.ts — close gracefully, and protect the task while it has active jobs
import { ECSClient, UpdateTaskProtectionCommand } from "@aws-sdk/client-ecs";
const ecs = new ECSClient({});
const taskArn = await getOwnTaskArn();                  // from the ECS task metadata endpoint

async function setProtection(enabled: boolean) {
  await ecs.send(new UpdateTaskProtectionCommand({
    cluster: "jobs", tasks: [taskArn], protectionEnabled: enabled, expiresInMinutes: 15 }));
}

let active = 0;
worker.on("active", async () => { if (active++ === 0) await setProtection(true); });
worker.on("completed", async () => { if (--active === 0) await setProtection(false); });
worker.on("failed", async () => { if (--active === 0) await setProtection(false); });

process.on("SIGTERM", async () => { await worker.close(); process.exit(0); });   // waits for active jobs
Scale in idle tasks first Six worker tasks are running and the service scales in to three. Tasks 1, 2 and 5 have active jobs and have set scale-in protection. Tasks 3, 4 and 6 are idle and unprotected, so ECS stops those three. The busy tasks drop protection when their jobs complete and become eligible for the next scale-in. Desired count 6 to 3 task 1 busy, protected task 2 busy, protected task 3 idle: stopped task 4 idle: stopped task 5 busy, protected task 6 idle: stopped No job is interrupted; protected tasks become eligible once their jobs complete.

ECS task scale-in protection makes the scheduler choose unprotected (idle) tasks to stop. Jobs longer than two minutes cannot be protected by stopTimeout alone, so for those, scale-in protection is the main safeguard; the protection expiry bounds how long a stuck task can block scale-in. Graceful close mechanics are covered in draining BullMQ workers before shutdown.

Step 5 — Add a Step Policy for Sudden Bursts

Target tracking adjusts gradually. For sudden large batches, a step-scaling policy on raw backlog can jump capacity faster.

resource "aws_appautoscaling_policy" "sync_burst" {
  name        = "sync-burst"
  policy_type = "StepScaling"
  # ...target references as above...
  step_scaling_policy_configuration {
    adjustment_type         = "ExactCapacity"
    cooldown                = 120
    metric_aggregation_type = "Maximum"
    step_adjustment { metric_interval_lower_bound = 0      metric_interval_upper_bound = 50000  scaling_adjustment = 20 }
    step_adjustment { metric_interval_lower_bound = 50000                                        scaling_adjustment = 60 }
  }
}
# attached to a CloudWatch alarm on Backlog > 10000

When both policies are active, Application Auto Scaling takes the larger desired count, so the step policy handles bursts and target tracking handles steady state.

Step 6 — Choose Fargate or EC2 Capacity

Fargate tasks start in about 30–90 seconds with no instance management; EC2-backed services with a capacity provider can pack workers densely and use Spot, but add instance scaling to the chain. For bursty batch workloads, Fargate Spot combined with a small on-demand base is often the simplest cost-efficient mix, as long as jobs tolerate interruption — see cutting worker costs with spot instances.

capacity_provider_strategy { capacity_provider = "FARGATE"       base = 2  weight = 1 }
capacity_provider_strategy { capacity_provider = "FARGATE_SPOT"  base = 0  weight = 4 }

Verification

Replay a nightly batch in staging and watch three graphs together: Backlog, DesiredTaskCount/RunningTaskCount, and queue wait p95. The service should reach its maximum within a few minutes of the batch starting, hold while backlog is high, step down over the following half hour, and produce no stalled-job events during scale-in.

aws application-autoscaling describe-scaling-activities --service-namespace ecs \
  --resource-id service/jobs/sync-worker --max-items 20 --query 'ScalingActivities[].[StartTime,Description]'

Gotchas & Edge Cases

Metric gaps. If the metric Lambda fails, target tracking sees missing data and may scale in. Alarm on the metric's absence and set treat_missing_data appropriately on step alarms.

Counting delayed jobs. Including delayed jobs scales out for work that cannot run yet. Exclude them.

Minimum of zero. Target tracking cannot scale from zero on a per-task metric (division by zero). Keep a minimum of one or two, or use the approach in scaling workers to zero.

Task startup time. A Fargate task pulling a large image can take over a minute to become ready. Keep images small, and remember that scale-out lag adds to backlog during bursts.

Downstream limits. Scaling to 60 tasks × 20 concurrency is 1,200 concurrent API calls; make sure the partner allows it, or add a limiter.

FAQ

Why not use the SQS-based ECS scaling examples? The same pattern applies; SQS publishes backlog metrics natively, while BullMQ needs the small publisher in Step 1.

Should one ECS service handle several queues? Only if they share a scaling signal and job profile. A service serving a fast interactive queue and a slow batch queue will scale on the combined backlog and give the interactive queue no guaranteed capacity. Run one service per queue class, each with its own metric and bounds.

How does protection interact with deployments? Deployments stop old tasks regardless of scale-in protection only after the protection expires or is cleared; set expiresInMinutes to the longest job duration so deploys are not blocked indefinitely by a stuck task.

How often should the metric be published? Every minute is enough for most workloads; target tracking evaluates over several data points anyway.

Can I scale on queue wait time instead? Yes, if you emit it; wait time is a better user-facing signal, but it lags backlog. Many teams scale on backlog per task and alert on wait time.

Related