bullmq-outbox
v0.1.2
Published
Outbox fallback for BullMQ: when Redis is out of memory or unreachable, jobs land in a store you own instead of disappearing. Works with BullMQ v5, v6 and Pro.
Maintainers
Readme
bullmq-outbox
When Redis runs out of memory or goes away, queue.add() rejects and the job
is gone. Not delayed — gone. There is nothing to retry, because nothing was
ever written down.
This package writes it down. A failed enqueue lands in a store you own, and a scheduler you control puts it back into Redis when Redis comes back.
import { createOutbox } from 'bullmq-outbox';
import { Queue } from 'bullmq';
const outbox = createOutbox({ store: myStore });
const emails = outbox.wrapQueue(new Queue('emails'));
// Behaves exactly like the queue you passed in, until Redis says no.
await emails.add('welcome', { userId: 1 });- No dependencies. Not even
bullmq. - Works with BullMQ v5, v6 and BullMQ Pro. The wrapper is structural, so it does not care which one you have.
- Any store. Four functions against whatever database you already run — Postgres, Redis, DynamoDB, Mongo, SQLite.
- Not a broker. It is a fallback. Redis stays the queue.
Install
npm install bullmq-outboxThe store
The package ships no adapters, because your outbox belongs in a database you already operate. You write four functions:
import type { OutboxStore } from 'bullmq-outbox';
const myStore: OutboxStore = {
save(entry) { /* persist it */ },
loadPending(limit) { /* oldest first, at most `limit` */ },
markProcessed(id) { /* it made it back into Redis */ },
markFailed(failure) { /* it did not — {id, error, attempts, expired} */ },
};Working implementations to copy, not install. Each ships with its schema:
| Store | File | Notes |
|---|---|---|
| Postgres | examples/postgres-store.ts · schema | The simplest road. A table and one partial index. |
| Redis | examples/redis-store.ts | Must be a different Redis from your queues. See the caveat in the file. |
| MongoDB | examples/mongodb-store.ts · setup | Partial index plus a TTL that cannot reap a pending job. |
| DynamoDB | examples/dynamodb-store.ts | Table shape taken from the production system this came from. |
There is also a MemoryOutboxStore for tests. It is not durable and it is not
for production.
Setting up the database
Run the schema once, before the app starts. Both files are idempotent, so they are safe in a migration step or at boot:
psql "$DATABASE_URL" -f node_modules/bullmq-outbox/examples/schema/postgres.sql
# or
node node_modules/bullmq-outbox/examples/schema/mongodb.js "$MONGO_URL" mydbThen copy the matching store file into your codebase and point it at your existing connection pool.
Using an AI assistant?
AGENTS.mdis a dense integration guide written for coding agents — contracts, the decisions that matter, and the mistakes that cost jobs. Point your agent at it and it should be able to wire this up without reading the source.
Draining
flush() re-enqueues what it finds. Call it however you schedule things:
const result = await outbox.flush(50);
// { requeued: 12, failed: 0, skipped: 3, expired: 0, details: [...] }Run the scheduler somewhere that does not share the failure you are recovering from. This is the part people get wrong. A drain loop that lives in the Redis that just died never fires. Use a cron, a Lambda, a separate small Redis — anything with an independent failure domain.
A job that keeps failing to re-enqueue is marked expired after maxAttempts
(default 10) and your store decides what that means: a dead letter table, an
alert, a row someone looks at on Monday.
NestJS
With @nestjs/bullmq you never call new Queue() — the module builds it and
you get it through @InjectQueue(). There is no instance to wrap, so swap the
class instead:
// main.ts, before NestFactory.create()
BullModule.queueClass = outbox.queueClass(Queue);Every queue in the app is now covered, including ones added later by someone who never read this. Coverage stops depending on anyone remembering.
Both @nestjs/bullmq and @taskforcesh/nestjs-bullmq-pro expose that setter —
it is the documented way to substitute QueuePro, and it takes any subclass.
It must run before the app bootstraps: queue providers read it at construction,
so a module's onModuleInit is too late.
Full working example, including the drain processor:
examples/nestjs/
Full example
import { createOutbox } from 'bullmq-outbox';
import { Queue } from 'bullmq';
import { createPostgresOutboxStore } from './outbox-store'; // copied from examples/
const outbox = createOutbox({
store: createPostgresOutboxStore(pool),
maxAttempts: 10,
// Real-time queues are better off dropping a job than replaying it later.
// A chat turn delivered fifteen minutes late is worse than one never sent.
shouldCapture: ({ queueName }) => !queueName.startsWith('live-'),
onJobSaved: (e) => metrics.increment('outbox.saved', { queue: e.queueName }),
onJobRequeued: (e) => metrics.timing('outbox.recovery_ms', e.ageMs),
onJobExpired: (e) => alerts.page('outbox job expired', e),
});
export const emails = outbox.wrapQueue(new Queue('emails', { connection }));
// somewhere with its own failure domain
setInterval(() => outbox.flush(50), 60_000);If the process that drains is not the process that enqueues, register the queues there instead of wrapping them:
outbox.registerQueue(new Queue('emails', { connection }));
await outbox.flush();Entries for queues this process does not know about are reported as skipped
and left alone, so several services can share one outbox table and each drains
only what it owns.
API
createOutbox(options)
| Option | Default | |
|---|---|---|
| store | required | Your four functions. |
| maxAttempts | 10 | Re-enqueue attempts before an entry is expired. |
| shouldCapture | capture everything | Return false to let a failure through unstored. |
| generateId | randomUUID() | Entry ids. |
| onJobSaved · onJobRequeued · onJobExpired · onSaveFailed · onFlush | — | Optional. Exceptions inside a hook are swallowed. |
outbox.wrapQueue(queue)
Returns a Proxy over your queue. add and addBulk gain the fallback;
everything else — getJob, pause, upsertJobScheduler, Pro's group and
batch APIs — passes straight through. Chainable methods (on, once) return
the wrapped queue, not the bare one, so wrapQueue(q).on('error', log) keeps
the fallback.
outbox.queueClass(BaseQueue)
Returns a subclass of BaseQueue with the fallback built in, for frameworks
that construct queues for you. Instances self-register, so flush() finds them
without a registerQueue call. See NestJS.
outbox.flush(limit?)
Re-enqueues pending entries. Returns counts plus a per-entry breakdown.
outbox.capture(queueName, jobName, data, opts, error)
Store a failure by hand, if you would rather not wrap the queue.
Things worth knowing
The original error is always re-thrown. The outbox buys you a replay, not a lie. Your caller still finds out the enqueue failed and still decides what to tell the user.
A failed addBulk stores every job in the batch. At this layer a partial
failure is indistinguishable from a total one, and a job replayed twice is
cheaper than a job lost. Set a jobId if you need the replay to dedupe.
parent is stripped from stored options. A flow parent may not exist by
the time the outbox drains, and BullMQ would reject the whole re-enqueue.
Flow children come back as standalone jobs.
The drain never captures its own failures. Replay re-enqueues through the raw queue, so a drain that cannot reach Redis counts an attempt against the existing entry instead of storing a second copy of it. Without that, every failed drain would double the backlog.
Jobs that fail while a drain is running are still captured. The bypass is scoped to the replay call, not to a window of time — live traffic failing mid-drain is exactly the traffic worth keeping.
Payloads are snapshotted at capture time. data goes through JSON when it
is stored, so mutating your object afterwards cannot change what gets replayed,
and an unserializable payload is reported through onSaveFailed rather than
blowing up inside your store.
A store that fails to record a successful replay does not lose the job.
The job is already in Redis at that point, so it counts as requeued and the
store error surfaces through onSaveFailed. Treating it as a failed replay
would re-enqueue the job and eventually expire one that had succeeded.
Replay is at-least-once. If the store write succeeds and the process dies
before the error propagates, you may get the job twice. Idempotent handlers,
or a jobId.
A broken store never breaks a job. If save throws, onSaveFailed fires
and the Redis error propagates unchanged — you do not want the fallback's
failure hiding the real one.
Tested against a real Redis
Unit tests with a fake queue prove the logic; they cannot prove the package
survives a Redis that is actually out of memory. integration/
runs against real Redis, Postgres and MongoDB containers and breaks the Redis
for real — maxmemory set just above current usage, so writes are rejected
with the same OOM command not allowed error ElastiCache throws.
cd integration && pnpm install && pnpm run verifyThe Postgres and MongoDB examples are tested verbatim, so what you copy is what was exercised.
Does the drain stay cheap as the store fills up?
The worry: after months in production the table holds hundreds of thousands of already-processed jobs. Does finding the handful still pending get slower?
pnpm run scale measures it. Every row below has exactly 50 pending entries
to fetch — what grows is the pile of already-processed rows sitting around
them:
| already-processed rows in the table | time to fetch the 50 pending | |---|---| | 0 (fresh table) | Postgres 1.01 ms · Mongo 0.87 ms | | 50,000 | Postgres 0.88 ms · Mongo 1.29 ms | | 200,000 | Postgres 0.66 ms · Mongo 0.69 ms | | 500,000 | Postgres 0.94 ms · Mongo 0.71 ms |
Flat. The differences between rows are measurement noise — they do not even
move in one direction — and that is the result: history does not cost you
anything. Always on an index, never a Seq Scan or COLLSCAN. Mongo reads
exactly 50 documents to return 50, with half a million in the collection.
The control case, same 500k table with the partial index dropped:
13.70 ms | SEQ SCAN ⚠15× slower, and it degrades with every job you ever process. So the flat numbers above are the index earning its keep, not the database being fast.
Index size tells the same story from another angle: a 74 MB table with a 16 kB index, because a partial index only covers what is pending. It tracks your backlog, not your history — the same property that makes the DynamoDB GSI this came from scale.
Measured on Docker containers on a laptop, so treat the absolute milliseconds as indicative. What the numbers establish is the shape — constant rather than growing — which is what survives a move to real infrastructure.
Where this came from
Built at Monest after losing jobs to an ElastiCache OOM. Three services, three Redis instances, one outbox table, running since April 2026 — where this package is the same design, rewritten to be framework-free and storage-agnostic.
Two things we learned that are not in the code:
- Set
reserved-memory-percenton your Redis parameter group. A Redis at 100% ofmaxmemorydoes not fail cleanly — it hangs, and BullMQ hangs with it (maxRetriesPerRequest: nullis required by BullMQ, so ioredis retries forever). Reserving a slice makes the OOM arrive as an error you can catch. - Watch the age of what you replay, not the count.
onJobRequeuedgives youageMs; that number is your real recovery time, and it is the one that tells you whether the fallback is working or quietly filling up.
Seeing that second point is easier with a dashboard. Bullpane is a self-hosted BullMQ dashboard — free edition, no login, does everything bull-board does.
See also
bullmq-fanout — publish one domain event to many queues, each with its own retries and failure domain. Wrap your queues with this package first and the fan-out inherits the durability.
License
MIT © Matheus Morett
