@queuebert/bullmq
v0.0.1
Published
BullMQ queue, worker, processor, and stats helpers for Queuebert-aware applications.
Maintainers
Readme
@queuebert/bullmq
BullMQ helpers for Queuebert-aware NestJS applications.
This package provides wrappers and base classes for collecting Queuebert stats
from BullMQ queues and workers without putting the integration inside
@queuebert/nest.
Installation
npm install @queuebert/bullmq @queuebert/nest @nestjs/bullmq @nestjs/common @nestjs/core bullmqExports
QueuebertBullMQModuleQueuebertBullMQServiceQueuebertQueueQueuebertWorkerBaseQueueProcessorDurationStatsCollectorGlobalStatsCollectorQueuebertProcessorAdapter- Queue, worker, lifecycle, dispatch, and stats types
NestJS Module
Use QueuebertBullMQModule when you want to create Queuebert-wrapped queues and
workers from a service.
import { Module } from '@nestjs/common';
import { QueuebertBullMQModule } from '@queuebert/bullmq';
@Module({
imports: [
QueuebertBullMQModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
defaultQueueOptions: {
defaultJobOptions: {
attempts: 3,
removeOnComplete: true,
},
},
statsConfig: {
windowMs: 60000,
maxSamples: 1000,
},
}),
],
})
export class AppModule {}Async configuration is also supported:
QueuebertBullMQModule.forRootAsync({
imports: [ConfigModule],
useFactory: (config: ConfigService) => ({
connection: {
host: config.getOrThrow('REDIS_HOST'),
port: config.getOrThrow('REDIS_PORT'),
},
}),
inject: [ConfigService],
});BaseQueueProcessor
BaseQueueProcessor extends Nest's WorkerHost and implements the
QueuebertProcessor interface from @queuebert/nest. It tracks job duration,
success/failure counts, throughput, Redis connection status, and optional cache
stats.
import { Processor } from '@nestjs/bullmq';
import { Injectable } from '@nestjs/common';
import { BaseQueueProcessor } from '@queuebert/bullmq';
import type { Job } from 'bullmq';
@Injectable()
@Processor('emails', { concurrency: 10 })
export class EmailProcessor extends BaseQueueProcessor {
constructor(private readonly emailService: EmailService) {
super({ queueName: 'emails' });
}
async processJob(job: Job<{ to: string; template: string }>) {
await this.emailService.send(job.data);
return true;
}
protected getCustomStats() {
return {
provider: this.emailService.providerName,
};
}
}Register that processor with @queuebert/nest:
import { QueuebertModule } from '@queuebert/nest';
QueuebertModule.forRoot({
queues: [
{
name: 'emails',
processor: EmailProcessor,
},
],
});QueuebertBullMQService
Create Queuebert-wrapped queues and workers programmatically:
import { Injectable, OnModuleInit } from '@nestjs/common';
import { QueuebertBullMQService } from '@queuebert/bullmq';
@Injectable()
export class WorkerBootstrap implements OnModuleInit {
constructor(private readonly queuebertBullMQ: QueuebertBullMQService) {}
async onModuleInit() {
const queue = this.queuebertBullMQ.createQueue<{ email: string }>('emails');
this.queuebertBullMQ.createWorker('emails', async (job) => {
await sendEmail(job.data.email);
});
await queue.add('welcome', {
email: '[email protected]',
});
}
}Direct Queue and Worker Wrappers
Use QueuebertQueue and QueuebertWorker directly outside NestJS when you want
typed dispatch results and local stats collection.
import { QueuebertQueue, QueuebertWorker } from '@queuebert/bullmq';
const queue = new QueuebertQueue<{ email: string }>('emails', {
connection: {
host: 'localhost',
port: 6379,
},
});
const worker = new QueuebertWorker(
'emails',
async (job) => {
await sendEmail(job.data.email);
},
{
connection: {
host: 'localhost',
port: 6379,
},
},
);
await queue.add('welcome', {
email: '[email protected]',
});
console.log(worker.getStats());Import Path
Import this integration directly from @queuebert/bullmq.
import { BaseQueueProcessor } from '@queuebert/bullmq';The old nested import shape is not part of the package contract.
