Integrations
BullMQ
Run typed background jobs on full BullMQ: Redis-durable work, stalled-job recovery, and replica-safe schedules behind the same queueService.
On this page
@rhythmjs/bullmq gives a typed job map and a kernel module over BullMQ. It exposes the same queueService as @rhythmjs/queue, so producers and handlers do not change if you switch between the two. Start from the template; the full app, with tests, is examples/bullmq.
Install#
BullMQ needs Redis and brings ioredis with it.
bun add @rhythmjs/bullmq @rhythmjs/rhythm
docker run --rm -p 6379:6379 redisWhy BullMQ#
Use it for BullMQ's production machinery: Redis-durable jobs, stalled-job recovery with lock renewal (a job whose worker dies is picked up again), atomic Lua-scripted state transitions, replica-coordinated Job Schedulers, per-job history and retention, and the Bull Board and Taskforce ecosystem. If you want zero queue dependencies on Bun natives instead, use @rhythmjs/queue.
The honest differences: job ids are string | undefined (BullMQ's own typing), and job options are BullMQ's JobsOptions passed through untouched.
Register the module#
The module owns the BullMQ lifecycle: on the kernel's teardown() every worker closes first, then the queue. A job map types producers and processors.
// src/jobs.ts
export type AppJobs = {
"email.send": { to: string; subject: string };
"order.process": { orderId: string };
};// src/app.module.ts
import { Rhythm } from "@rhythmjs/rhythm";
import { bullmqModule } from "@rhythmjs/bullmq";
import type { AppJobs } from "./jobs";
const app = new Rhythm().register(
bullmqModule.forRoot<AppJobs>({
name: "app",
connection: { host: "127.0.0.1", port: 6379 },
defaultJobOptions: { attempts: 3, backoff: { type: "exponential", delay: 1000 }, removeOnComplete: 100 },
}),
({ queueService }) => ({ queueService }),
);Options: name (the queue name, default rhythm), connection (ioredis ConnectionOptions: an options object or an existing instance, default 127.0.0.1:6379), prefix (the Redis key prefix), and defaultJobOptions (BullMQ defaults for every job).
Start a worker with the app#
To process jobs in the same process, build the queue service in a provider, start a worker on it, and close it with the app. close() closes every worker first, then the queue, so one dispose function releases Redis:
// src/emails/queue.ts
import { createQueueService, type QueueService } from "@rhythmjs/bullmq";
export function createQueue({ mailService }: { mailService: typeof import("./mail.service").mailService }) {
const queueService = createQueueService<AppJobs>({ name: "emails", connection: { host: "127.0.0.1", port: 6379 } });
queueService.process({ "email.send": (payload) => mailService.send(payload) }, { concurrency: 4 });
return { queueService };
}
export function closeQueue({ queueService }: { queueService: QueueService<AppJobs> }) {
return queueService.close();
}// src/emails/emails.module.ts
export const emailsModule = new Rhythm<RhythmHttpContext>({ name: "emails", type: "module" })
.provide(() => ({ mailService }))
.provide(createQueue, closeQueue)
.use(emailsController.middleware());Controllers then read ctx.queueService, for example await ctx.queueService.add("email.send", { to, subject }), and tests stub it. bullmqModule.forRoot above is the same service as a ready-made module, for apps that only produce jobs.
Produce and process#
await queueService.add("email.send", { to: "a@b.c", subject: "hi" }, { delay: 5000 });
await queueService.addBulk([
{ name: "email.send", payload: { to: "x@y.z", subject: "s" } },
{ name: "order.process", payload: { orderId: "o1" }, options: { priority: 1 } },
]);
const worker = queueService.process(
{
"email.send": async (payload, job) => sendMail(payload),
"order.process": async (payload, job) => fulfil(payload.orderId),
},
{
concurrency: 4,
onCompleted: (name, id) => metrics.increment(`jobs.${name}.ok`),
onFailed: (name, id, error) => log.error(`job ${name}#${id} failed`, error),
},
);
await worker.close();Job options are BullMQ's own JobsOptions: delay, attempts, backoff, priority, removeOnComplete and removeOnFail, jobId, lifo, deduplication, and everything else BullMQ accepts. Each process() call creates one BullMQ Worker on the queue, dispatching to processors by name with a typed payload and JobInfo ({ id, name, attemptsMade }).
A job with no registered processor throws no processor registered for job "…" and lands in BullMQ's retry and failed flow like any other failure. counts() forwards getJobCounts(): waiting, active, delayed, completed, failed, and the rest of BullMQ's states.
Repeats#
schedule(name, repeat, payload, options?) upserts a BullMQ Job Scheduler keyed by the job name. Schedulers are Redis-coordinated: register the same schedule from every replica and it still fires once, so declare schedules at boot on all instances. repeat is BullMQ's own RepeatOptions (pattern, every, tz, limit, startDate, endDate). unschedule(name) stops future firings; jobs already enqueued are unaffected.
await queueService.schedule("order.process", { pattern: "0 3 * * *", tz: "UTC" }, { orderId: "nightly" });
await queueService.schedule("email.send", { every: 60_000, limit: 10 }, { to: "digest@x", subject: "s" });
await queueService.unschedule("order.process");The whole Queue#
queueService.queue is the underlying BullMQ Queue, so everything the typed service does not wrap is one property away, on the same connection:
await queueService.queue.pause(); // and resume()
await queueService.queue.getJobs(["failed"]);
await queueService.queue.getJobSchedulers();
await queueService.queue.drain();It is also what dashboards want: hand it to Bull Board's BullMQAdapter or Taskforce and they see the queue exactly as your app does. For parent-child flows, construct BullMQ's FlowProducer with the same connection options.
Things to check#
- Tests can substitute BullMQ with
mock.module("bullmq", …)frombun:test(recordingQueueandWorkerfakes), so the wiring is tested without Redis. Run live behavior (retries, schedulers) in a separate suite skipped whenREDIS_URLis unset. tzis the timezone field on BullMQ repeats, where@rhythmjs/queuespells ittimezone.- Redis has to be reachable when the module starts and when jobs are added.