bullmq-fanout
v0.1.0
Published
Turn BullMQ into a message bus: publish one domain event to many queues, each with its own worker, retries and failure domain. Works with BullMQ v5, v6 and Pro.
Downloads
4,449
Maintainers
Readme
bullmq-fanout
You already run BullMQ. Then a second team needs to know when an order is paid, and a third one after that.
The usual next step is a queue called order-paid with three workers on it —
and now analytics being slow delays invoices, one consumer's retries are
everyone's retries, and nobody can say who is listening without grepping four
repositories.
This package fans out at publish time instead. One event, one add per
consumer, each into its own queue with its own worker, concurrency,
retries and dead letters.
export const ORDER_PAID = defineEvent({
name: 'order.paid',
payload: OrderPaidSchema,
subscribers: [
{ queue: 'invoices', job: 'issue-invoice', attempts: 10 },
{ queue: 'emails', job: 'send-receipt' },
{ queue: 'analytics', job: 'track-revenue', attempts: 1 },
],
});
await publisher.publish(ORDER_PAID, order, { eventId: String(order.id) });- No dependencies. Not even
bullmq. - Works with BullMQ v5, v6 and BullMQ Pro.
- No broker, no runtime. Consumers are ordinary BullMQ workers.
- Consumers are isolated. One that is down does not cost you the others.
Install
npm install bullmq-fanoutWhy one queue per consumer
This is the whole design, so it is worth being explicit.
| Shared queue, N workers | A queue per consumer |
|---|---|
| One backlog. A slow consumer delays everyone. | Backlogs are independent. |
| One retry policy for all. | Analytics retries once, invoicing ten times. |
| A poison job blocks every consumer. | It blocks one. |
| Pausing means pausing everything. | Pause one consumer. |
| "Who consumes this?" — grep the org. | Read subscribers. |
The cost is N enqueues per event instead of one. That is a few hundred microseconds against a Redis you are already talking to, and it buys you failure domains.
Setup
import { createPublisher } from 'bullmq-fanout';
import { Queue } from 'bullmq';
const publisher = createPublisher({
queues: {
invoices: new Queue('invoices', { connection }),
emails: new Queue('emails', { connection }),
analytics: new Queue('analytics', { connection }),
},
});Fail at boot rather than on the first publish:
import { requiredQueues } from 'bullmq-fanout';
const missing = requiredQueues([ORDER_PAID, ORDER_REFUNDED])
.filter((name) => !(name in queues));
if (missing.length) throw new Error(`Unregistered queues: ${missing}`);NestJS
Queues come from DI, so the wiring resolves them by token. One module you copy once:
@Module({
imports: [FanoutModule.forFeature([ORDER_PAID])],
providers: [OrdersService],
})
export class OrdersModule {}It reads subscribers off the events, registers those queues and provides a
Publisher wired to them — so adding a consumer stays a one-line change to
subscribers, with no edit to the publishing module.
Full working example, publisher and consumer:
examples/nestjs/
Using BullMQ Pro? Import BullModule and getQueueToken from
@taskforcesh/nestjs-bullmq-pro. Nothing else changes.
Declaring events
Keep them in one directory. The value is that subscribers is the answer to
"who consumes this?" — reviewable in a diff, not scattered across services.
import { z } from 'zod';
import { defineEvent } from 'bullmq-fanout';
export const ORDER_PAID = defineEvent({
name: 'order.paid',
payload: z.object({ orderId: z.number(), amountCents: z.number() }),
subscribers: [{ queue: 'invoices', job: 'issue-invoice' }],
});payload is any object with a parse method — a zod schema, a valibot
schema, or a function you wrote. The package depends on none of them, and the
field is optional if you do not want validation.
It runs once per publish, at the producer. A drifted contract throws where someone can see it, instead of in four workers at 3am.
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.
Full examples: examples/events.ts ·
examples/publish.ts ·
examples/consumer.ts
Publishing
const result = await publisher.publish(ORDER_PAID, order, {
eventId: String(order.id),
});
// { event: 'order.paid', published: ['invoices', 'emails'], failures: [...] }eventId is what makes retries safe. It derives the job id per
subscriber — order.paid:1234:invoices — so publishing the same event twice
does not double-invoice anyone while the first job is still in the queue. Pass
the business id: an order id, a payment id.
A failed subscriber does not stop the fan-out. Everyone else still gets
their job; the failure comes back in failures and through onPublishFailed.
One dead consumer must not cost you the other four.
Only an invalid payload throws. A broken contract should be loud.
Batches use one addBulk per subscriber:
await publisher.publishBulk(ORDER_PAID, orders);Consuming
Ordinary BullMQ workers. This package is not in the path at consume time:
new Worker('invoices', async (job) => issueInvoice(job.data), {
connection,
concurrency: 2,
});Making the fan-out durable
publish is as reliable as queue.add is. If Redis is down, the job is gone
— the same as any direct add.
Wrap the queues with bullmq-outbox and the fan-out inherits the fallback without knowing it exists:
const outbox = createOutbox({ store: myStore });
const publisher = createPublisher({
queues: {
invoices: outbox.wrapQueue(new Queue('invoices', { connection })),
emails: outbox.wrapQueue(new Queue('emails', { connection })),
},
});A subscriber whose enqueue fails is persisted and replayed on the next drain.
The others never noticed. This works because the publisher only ever calls
add — anything with that shape is a valid queue here.
API
createPublisher(options)
| Option | Default | |
|---|---|---|
| queues | required | Name → queue. A Record or a Map. |
| defaultAttempts | 3 | For subscribers without their own. |
| defaultBackoff | { type: 'exponential', delay: 5000 } | |
| onPublished · onPublishFailed | — | Optional. Exceptions inside a hook are swallowed. |
publisher.publish(event, payload, options?)
options.eventId derives a deterministic job id per subscriber.
options.opts merges extra job options into every subscriber of this publish.
publisher.publishBulk(event, payloads, options?)
One addBulk per subscriber. Every payload is validated before anything is
enqueued, so a bad item cannot leave the batch half delivered. No eventId —
a deterministic id needs a per-payload business key, and there is no honest
way to guess one from an array.
defineEvent(definition) · requiredQueues(events) · publisher.register(name, queue)
Things worth knowing
This is not a broker. There is no delivery guarantee beyond what BullMQ gives you per queue, no ordering across subscribers, no replay of past events for a consumer that joins later. If you need those, you need Kafka. If you need three teams to react to an order being paid, you need this.
Subscribers are wired at the publisher. A new consumer means editing
subscribers and deploying the producer. That is a deliberate trade: you give
up runtime subscription and you get a list you can read.
Delivery is at-least-once. A subscriber can get the same event twice if
the producer retries. Use eventId, and write idempotent handlers.
Tested against a real Redis
integration/ runs the fan-out against real Redis containers
with real BullMQ workers — including the composition above, where an OOM
during a publish puts every subscriber in Postgres and the drain restores all
of them.
cd integration && npm install && npm run verifyWhere this came from
Built at Monest — a single tag.applied event now
feeds four independent consumers, one of which is a client-facing integration
with its own retry policy and its own failure budget.
Once a single event lands in four queues, the question stops being "did it work?" and becomes "which of the four is behind?". Bullpane is a self-hosted BullMQ dashboard for exactly that — free edition, no login, does everything bull-board does.
See also
bullmq-outbox — when Redis is out of memory or down, jobs land in a store you own and are replayed later.
License
MIT © Matheus Morett
