Integration Testing BullMQ Workers with Testcontainers
BullMQ has no in-memory mode — its behaviour lives in Lua scripts running inside Redis — so meaningful tests need a real Redis, and this guide shows how to run them quickly and reliably with Testcontainers, as part of Testing Background Jobs in Backend Frameworks & Worker Scaling. The result is a suite that exercises real retries, backoff, flows, rate limiting, and stalled-job recovery in a few seconds per test.
Problem Statement
A Node.js service processes image moderation jobs with BullMQ. Its tests mock the Queue and Worker classes entirely, so they verify that queue.add was called and nothing else. Production incidents have included a flow whose parent never ran because a child's removeOnComplete deleted the result the parent needed, a backoff configuration that retried instantly because delay was specified in seconds instead of milliseconds, and a worker that stalled jobs because its CPU-bound handler blocked the event loop past the lock duration. You want tests that run real workers against real Redis, assert on actual job state transitions, and stay fast and isolated enough to run on every pull request.
Prerequisites
- Node 20+, BullMQ 5.x, and Jest or Vitest.
@testcontainers/redis(ortestcontainers) and Docker available locally and in CI.- Processor functions exported separately from the
Workerconstruction, so tests can build workers with test-specific options. - An understanding of BullMQ job states: waiting, delayed, active, completed, failed.
Step 1 — Start One Redis Container per Test Run
Starting a container takes a second or two; do it once in a global setup and share the connection details through environment variables.
// test/global-setup.ts (Vitest globalSetup or Jest globalSetup)
import { RedisContainer, StartedRedisContainer } from "@testcontainers/redis";
let container: StartedRedisContainer;
export async function setup() {
container = await new RedisContainer("redis:7.2-alpine")
.withCommand(["redis-server", "--maxmemory-policy", "noeviction"]) // BullMQ requirement
.start();
process.env.TEST_REDIS_URL = container.getConnectionUrl();
}
export async function teardown() {
await container?.stop();
}
Configure noeviction exactly as in production; BullMQ warns about other policies because evicted keys mean lost jobs, as explained in Redis maxmemory policy for queues.
Step 2 — Build a Per-Test Harness
Each test gets a unique queue name, its own connections, and guaranteed cleanup. Unique names let test files run in parallel against one Redis without interfering.
// test/harness.ts
import { Queue, Worker, QueueEvents, Processor, WorkerOptions } from "bullmq";
import IORedis from "ioredis";
import { randomUUID } from "node:crypto";
export async function harness<T>(processor: Processor<T>, opts: Partial<WorkerOptions> = {}) {
const name = `moderation-${randomUUID().slice(0, 8)}`;
const connection = new IORedis(process.env.TEST_REDIS_URL!, { maxRetriesPerRequest: null });
const queue = new Queue<T>(name, { connection });
const events = new QueueEvents(name, { connection: connection.duplicate() });
const worker = new Worker<T>(name, processor, { connection: connection.duplicate(), ...opts });
await Promise.all([events.waitUntilReady(), worker.waitUntilReady()]);
return {
queue, events, worker,
async close() {
await worker.close();
await events.close();
await queue.obliterate({ force: true }); // remove every key for this queue
await queue.close();
await connection.quit();
},
};
}
maxRetriesPerRequest: null is required for BullMQ's blocking commands. QueueEvents is the reliable way to wait for a specific job's outcome without polling.
Step 3 — Wait on Job Outcomes, Not Timers
job.waitUntilFinished(queueEvents, ttl) resolves with the return value on completion and rejects with the failure reason. It replaces setTimeout-based waiting, which is both slow and flaky.
// moderation.test.ts
import { describe, it, expect, afterEach } from "vitest";
import { harness } from "./harness";
import { moderateImage } from "../src/processors/moderate";
let h: Awaited<ReturnType<typeof harness>>;
afterEach(() => h?.close());
describe("moderation worker", () => {
it("completes and returns the verdict", async () => {
h = await harness(moderateImage);
const job = await h.queue.add("moderate", { imageId: "img-1" });
const result = await job.waitUntilFinished(h.events, 10_000);
expect(result).toEqual({ imageId: "img-1", verdict: "approved" });
});
it("fails permanently on unsupported formats without retrying", async () => {
h = await harness(moderateImage);
const job = await h.queue.add("moderate", { imageId: "img-heic" }, { attempts: 5 });
await expect(job.waitUntilFinished(h.events, 10_000)).rejects.toThrow(/unsupported format/);
const fresh = await h.queue.getJob(job.id!);
expect(fresh!.attemptsMade).toBe(1); // UnrecoverableError: no retries
});
});
The second test asserts that the processor throws BullMQ's UnrecoverableError for a permanent failure, so the job fails on the first attempt despite attempts: 5. That classification is the same idea as in BullMQ retries and backoff strategies.
Step 4 — Test Retries and Backoff with Real Timing, Scaled Down
Real backoff delays are seconds to minutes. Keep the shape of the production configuration and scale the base delay through configuration, so tests run the real retry path quickly.
// src/config.ts
export const BACKOFF_BASE_MS = Number(process.env.BACKOFF_BASE_MS ?? 2000);
export const defaultJobOptions = {
attempts: 5,
backoff: { type: "exponential", delay: BACKOFF_BASE_MS }, // milliseconds!
};
// retry.test.ts — BACKOFF_BASE_MS=20 in the test environment
it("retries transient failures with increasing delays", async () => {
let calls = 0;
const stamps: number[] = [];
h = await harness(async () => {
stamps.push(Date.now());
if (++calls < 3) throw new Error("timeout from vision API");
return "ok";
});
const job = await h.queue.add("moderate", { imageId: "img-2" }, defaultJobOptions);
await expect(job.waitUntilFinished(h.events, 10_000)).resolves.toBe("ok");
expect(calls).toBe(3);
const gaps = stamps.slice(1).map((t, i) => t - stamps[i]);
expect(gaps[1]).toBeGreaterThan(gaps[0]); // exponential, not constant or instant
});
This catches the seconds-vs-milliseconds bug directly: with delay: 2 intended as seconds, the gaps are two milliseconds and the assertion on growth is fragile but the absolute check gaps[0] >= BACKOFF_BASE_MS fails outright. Add that check for the configuration you ship.
Step 5 — Test Flows and Stalled-Job Recovery
Flows (parent jobs waiting on children) and stalled-job detection are two behaviours that exist only inside Redis.
import { FlowProducer } from "bullmq";
it("runs the parent after all children and sees their results", async () => {
h = await harness(async (job) => {
if (job.name === "scan") return { part: job.data.part, clean: true };
const children = await job.getChildrenValues(); // needs child results retained
return Object.values(children).every((c: any) => c.clean) ? "approved" : "rejected";
});
const flow = new FlowProducer({ connection: h.queue.opts.connection as any });
const tree = await flow.add({
name: "verdict", queueName: h.queue.name, data: {},
children: [1, 2, 3].map((part) => ({ name: "scan", queueName: h.queue.name, data: { part } })),
});
await expect(tree.job.waitUntilFinished(h.events, 10_000)).resolves.toBe("approved");
await flow.close();
});
it("recovers a job whose worker stops renewing its lock", async () => {
h = await harness(async () => { const end = Date.now() + 1500; while (Date.now() < end) {} return "done"; },
{ lockDuration: 500, stalledInterval: 300, maxStalledCount: 1 });
const job = await h.queue.add("moderate", { imageId: "img-3" });
const stalled = new Promise((res) => h.events.on("stalled", ({ jobId }) => jobId === job.id && res(true)));
await expect(stalled).resolves.toBe(true); // busy loop blocked lock renewal
});
The first test fails if a child is added with removeOnComplete: true, because getChildrenValues then finds nothing — the flow bug from the problem statement. The second reproduces the event-loop-blocking stall and is the reason CPU-heavy processors belong in BullMQ sandboxed processors. Lock and stall mechanics are covered in BullMQ lock duration and stalled jobs.
Step 6 — Run It in CI
Testcontainers works in GitHub Actions and most CI systems with Docker. Run integration tests as their own job so unit tests keep instant feedback.
# .github/workflows/test.yml
jobs:
integration:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with: { node-version: 20, cache: npm }
- run: npm ci
- run: npx vitest run --project integration --pool=forks --poolOptions.forks.maxForks=4
env:
BACKOFF_BASE_MS: "20"
TESTCONTAINERS_RYUK_DISABLED: "false" # keep container cleanup on
Four parallel forks sharing one container is typically faster than one container per fork; unique queue names from Step 2 keep them independent.
Verification
Reintroduce each production bug and confirm a test fails:
# 1. removeOnComplete on flow children -> flow test fails (parent sees no child values)
# 2. backoff delay in seconds (delay: 2) -> backoff gap assertion fails
# 3. remove the sandbox for CPU processor -> stall test observes "stalled" (expected in that test)
npx vitest run --project integration --reporter=verbose
Typical timings: container start about 1.5 s once, each test 50–500 ms.
Gotchas & Edge Cases
Unclosed workers hang the run. A worker left open keeps a blocking connection alive and Jest/Vitest never exits. Always close in afterEach, and set a test-runner teardown timeout.
Shared connections and blocking commands. Reusing one ioredis connection for a Worker and a QueueEvents can deadlock because both issue blocking reads; use connection.duplicate() for each.
Timing assertions on busy runners. Assert on ordering and lower bounds, not exact durations, and keep scaled-down delays well above scheduler jitter (tens of milliseconds, not one).
Global state in processors. Module-level caches or clients leak between tests in the same fork. Construct dependencies inside a factory and inject them into the processor.
FAQ
Can I use ioredis-mock instead of a container? Not for BullMQ: it relies on Lua scripts and blocking commands that mocks do not implement faithfully. A container is the only reliable option.
Should every worker test be an integration test? No. Keep processor logic in plain functions with unit tests, and use this harness for queue behaviour: retries, flows, rate limits, stalls, and events.
How do I test that a job is added with the right options from application code?
Call the application function against the real test queue, then read the job back with queue.getJobs(["waiting", "delayed"]) and assert on job.opts — attempts, backoff, delay, priority, and jobId for deduplication. That checks the exact options that will reach Redis, including defaults merged from defaultJobOptions, which a mocked queue.add cannot show you.
Do these tests work with Redis Cluster? Run a cluster container only if production uses Redis Cluster; the main difference is that queue names must use a hash tag prefix so all of a queue's keys land on one slot. A single dedicated test for that configuration is usually enough.
How do I test the rate limiter?
Configure the worker's limiter in the harness, add more jobs than the limit, and assert on the completion timestamps — see configuring the BullMQ rate limiter.
Related
- Testing Background Jobs — where this layer fits in the overall strategy.
- BullMQ for Node.js Ecosystems — the production setup these tests protect.
- BullMQ Flows for Parent-Child Jobs — the flow semantics under test.
- Chaos Testing Worker Crashes and Redelivery — killing workers mid-job.