Running BullMQ on Redis Cluster
BullMQ relies on Lua scripts that touch many keys of one queue atomically, and Redis Cluster only allows that when all those keys live on the same hash slot. This guide configures BullMQ for Redis Cluster so it works on the first try and actually spreads load, as part of BullMQ for Node.js Ecosystems in Backend Frameworks & Worker Scaling.
Problem Statement
A platform team moved BullMQ from a single Redis instance to a 3-shard managed Redis Cluster to get past the single-core throughput ceiling. The first deploy failed immediately with CROSSSLOT Keys in request don't hash to the same slot. After adding a hash tag to the prefix, everything worked — but all 25 queues were now on one shard, whose CPU sat at 90% while the other two idled. You want BullMQ's scripts to work on Cluster, queues distributed across shards by expected load, flows (parent and child jobs in different queues) that still work, and predictable behaviour during failover.
Prerequisites
- BullMQ 5.x with
ioredis5.x. - A Redis Cluster (self-managed 7.x or a managed service such as ElastiCache/MemoryDB cluster mode, Azure Cache Enterprise, or Redis Cloud).
- An estimate of load per queue (jobs per second, payload size) to plan shard placement.
maxmemory-policy noevictionon all shards.
Step 1 — Understand Why Hash Tags Are Required
Redis Cluster assigns each key to one of 16,384 slots by hashing the key. A Lua script or multi-key command must only touch keys in one slot. BullMQ keys look like bull:emails:wait, bull:emails:active, bull:emails:1234 — different keys that hash to different slots, so scripts fail. A hash tag — a substring in braces — makes Redis hash only that substring, forcing all of a queue's keys into one slot.
without tag: bull:emails:wait -> slot 5601
bull:emails:active -> slot 12870 -> CROSSSLOT in any BullMQ script
with tag: {emails}:emails:wait -> slot of "emails"
{emails}:emails:active -> same slot -> scripts work
The consequence is that one queue lives entirely on one shard. Cluster scales BullMQ by spreading different queues across shards, never by splitting one queue.
Step 2 — Give Each Queue Its Own Hash Tag
The common mistake — and the cause of the one-hot-shard problem in the scenario — is using the same tag for every queue, such as prefix: "{bull}". That puts every queue in the same slot. Use a per-queue tag so each queue can land on a different slot.
import { Queue, Worker } from "bullmq";
import { Cluster } from "ioredis";
const connection = new Cluster(
[{ host: "clustercfg.jobs.abc123.use1.cache.amazonaws.com", port: 6379 }],
{
dnsLookup: (address, callback) => callback(null, address), // TLS + AWS config endpoint
redisOptions: { tls: {}, maxRetriesPerRequest: null }, // BullMQ requirement
slotsRefreshTimeout: 2000,
},
);
// Per-queue tag: each queue's keys share a slot, different queues can differ
const queueFor = (name: string) => new Queue(name, { connection, prefix: `{${name}}` });
const workerFor = (name: string, fn: Processor) => new Worker(name, fn, { connection, prefix: `{${name}}` });
Producers and workers must use the same prefix for a queue, or they read and write different keys. Put the prefix rule in one helper, as above.
Step 3 — Place Busy Queues on Different Shards
Which shard a tag lands on is decided by the hash of the tag, not by you. With a handful of queues, check where each landed and choose tags that spread the heavy ones.
import { calculateSlot } from "cluster-key-slot";
async function shardFor(tag: string): Promise<string> {
const slot = calculateSlot(tag);
const slots = await connection.cluster("SLOTS") as [number, number, [string, number]][];
const owner = slots.find(([from, to]) => slot >= from && slot <= to)!;
return `${owner[2][0]}:${owner[2][1]}`;
}
for (const q of ["emails", "webhooks", "thumbnails", "reports"]) {
console.log(q, calculateSlot(q), await shardFor(q));
}
// emails 11232 10.0.3.14:6379
// webhooks 5840 10.0.2.21:6379
// thumbnails 14101 10.0.3.14:6379 <- same shard as emails; both are heavy
// reports 3316 10.0.1.9:6379
If two heavy queues collide, change one queue's tag (for example {thumbnails-2}) — a new tag means new keys, so migrate by draining the old queue while producing to the new one. For one very heavy workload, split it into several queues with different tags ({events-0} to {events-5}) and route jobs by a hash of their key, as in consistent hashing for queue shards.
Step 4 — Keep Flows Working Across Queues
BullMQ flows link parent and child jobs, which may live in different queues. Flow operations touch both parent and child keys, so on Cluster all queues in one flow must share a hash tag. Give flow-related queues a common tag, accepting that they share a shard.
// All queues participating in the "render" flow share the {render} tag
const FLOW_TAG = "{render}";
const flow = new FlowProducer({ connection, prefix: FLOW_TAG });
await flow.add({
name: "assemble", queueName: "render-assemble", data: { docId },
children: pages.map((p) => ({ name: "page", queueName: "render-page", data: { docId, p } })),
});
new Worker("render-assemble", assemble, { connection, prefix: FLOW_TAG });
new Worker("render-page", renderPage, { connection, prefix: FLOW_TAG });
If a flow's queues are too busy to share one shard, reconsider using a flow for that workload — a workflow record in your database, as in Job Chaining & Workflow Orchestration, can coordinate steps across queues on different shards.
Step 5 — Plan for Failover
When a primary shard fails, the cluster promotes a replica. For the few seconds of failover, commands to that shard fail; BullMQ workers on queues in that slot see connection errors and retry. Because replication is asynchronous, the last writes before the failure (a job just added, a completion just recorded) can be lost.
connection.on("+node", (node) => log.info({ node: node.options.host }, "cluster node added"));
connection.on("-node", (node) => log.warn({ node: node.options.host }, "cluster node removed"));
connection.on("error", (err) => log.error({ err }, "cluster connection error"));
// Workers: tolerate failover pauses without stalling jobs
new Worker("emails", sendEmail, { connection, prefix: "{emails}", lockDuration: 60_000 });
A lock duration longer than typical failover time (usually 10–30 seconds) avoids a wave of stalls when a shard fails over. Lost completions cause re-processing, which idempotent jobs absorb; lost additions are why enqueue should happen through an outbox when losing a job is unacceptable — see Redis Sentinel for queue high availability for the same trade-off with Sentinel.
Step 6 — Monitor Per Shard, Not Per Cluster
Cluster-wide averages hide a hot shard. Watch CPU, memory, and command latency per shard, and map queues to shards in your dashboards.
# Per-shard main-thread CPU (redis_exporter scraping each node)
rate(redis_cpu_user_seconds_total[1m]) + rate(redis_cpu_sys_seconds_total[1m])
# Memory headroom per shard: a busy queue's backlog fills its shard only
redis_memory_used_bytes / redis_memory_max_bytes
A shard at high CPU tells you which queues to move or split; a shard near its memory limit tells you which queue's backlog to worry about. Throughput limits per shard follow benchmarking Redis broker throughput.
Verification
it("runs every queue's scripts on the cluster", async () => {
for (const name of ["emails", "webhooks", "thumbnails", "reports"]) {
const q = queueFor(name);
const job = await q.add("probe", {});
const w = workerFor(name, async () => "ok");
const events = new QueueEvents(name, { connection, prefix: `{${name}}` });
await expect(job.waitUntilFinished(events, 10_000)).resolves.toBe("ok");
await w.close(); await events.close();
}
});
Run it in staging against the real cluster, then trigger a manual failover of one shard (CLUSTER FAILOVER on a replica, or the managed service's test failover) under load and confirm workers recover without manual action.
Gotchas & Edge Cases
Same tag everywhere. prefix: "{bull}" makes Cluster pointless for BullMQ. Tag per queue.
Changing a prefix orphans jobs. New prefix means new keys; jobs under the old prefix are invisible to new workers. Drain before switching.
Local development without Cluster. Hash-tagged prefixes work on a standalone Redis too, so use the same prefix helper everywhere and run a small cluster in CI to catch CROSSSLOT errors before production.
Resharding. Moving slots between shards while BullMQ runs can cause MOVED/ASK redirects and brief errors. Schedule resharding in quiet periods.
QueueEvents and streams. Event streams live under the queue's tag too; QueueEvents must use the same prefix.
FAQ
Is Redis Cluster the right way to scale BullMQ? When the bottleneck is one Redis instance's CPU and you have several busy queues, yes. For one very hot queue, you must split it into several queues; Cluster cannot split a queue.
Can I use ElastiCache Serverless? Serverless caches present a cluster-compatible endpoint; the same hash-tag rules apply. Check that the service supports the Lua and blocking commands BullMQ uses.
How many queues do I need to make Cluster worthwhile? Enough busy ones to spread: with three shards and one queue carrying 90% of the load, Cluster buys almost nothing until that queue is split. Count load per queue first; if one dominates, split it into several tagged queues before adding shards.
Does Sentinel avoid these restrictions? Sentinel provides high availability for a single primary, with no slot restrictions but no horizontal scaling. Many deployments start with Sentinel and move to Cluster only when CPU demands it.
Related
- BullMQ for Node.js Ecosystems — the overall BullMQ setup.
- Redis Sentinel for Queue High Availability — the single-primary HA alternative.
- Benchmarking Redis Broker Throughput — deciding when a single instance is not enough.
- BullMQ Flows for Parent-Child Jobs — flows and their key layout.