@nathapp/nestjs-queue
v4.1.0
Published
NestJS Queue Module with multi-provider support (BullMQ, Kafka, RabbitMQ, Redis Pub/Sub)
Downloads
579
Readme
@nathapp/nestjs-queue
Multi-provider queue abstraction for NestJS — BullMQ, Kafka, and RabbitMQ
behind one IQueueProvider interface, with decorator-based processors,
idempotency, circuit breaking, and delayed-job scheduling.
Installation
npm install @nathapp/nestjs-queueInstall the client library for whichever provider you use (e.g. bullmq +
ioredis, kafkajs, or amqplib) as a peer dependency. This package is part
of the @nathapp/nestjs-* peer-layer stack and builds on
@nathapp/nestjs-common.
What it provides
QueueModule—register()/registerAsync()for full processing mode (wiresQueueService+ProcessorExplorerfor@Processor/@Processdecorators);forProducer()/forProducerAsync()for producer-only mode (justQueueService, no decorator scanning).QueueService—add,addBulk,process,getJob,getJobs,removeJob,pause,resume, and cleanup methods, delegating to the configuredIQueueProvider.QueueProviderTypeenum —BULLMQ,KAFKA,RABBITMQ,CUSTOM.- Decorators —
@Processor(queueName | options)(class) and@Process(name? | options)(method) to declare job handlers;@OnQueueEvent,@BatchProcessfor event hooks and batch processing. - Interfaces —
IJob<T>,JobOptions,IQueueProvider,QueueModuleOptions,QueueModuleAsyncOptions,HealthCheckResult,IdempotencyStoreOptions. - Utilities —
DelayedJobScheduler,QueueError,ShutdownOptions, sanitization/validation helpers, anInMemoryIdempotencyStore(understore), and a Redis-backed idempotency store (underutils). - Circuit breaking —
QueueModuleOptions.circuitBreaker(global) andJobOptions.circuitBreaker(per-processor, via@Processor) accept a circuit-breaker config object to guard against cascading failures. The breaker implementation and its option/state/metrics types are internal and not part of the package's exported surface.
Provider implementations (BullMQ/Kafka/RabbitMQ) are lazy-loaded internally
and are not part of the package's static export surface — select one via
provider: QueueProviderType.<X> in the module options.
Usage
Register a queue and declare a processor:
import { Module } from '@nestjs/common';
import { QueueModule, Processor, Process, IJob } from '@nathapp/nestjs-queue';
import { QueueProviderType } from '@nathapp/nestjs-queue';
@Module({
imports: [
QueueModule.register({
provider: QueueProviderType.BULLMQ,
options: { connection: { host: 'redis', port: 6379 } },
}),
],
providers: [EmailProcessor],
})
export class AppModule {}
@Processor('email-queue')
export class EmailProcessor {
@Process('send')
async handleSend(job: IJob<{ to: string }>) {
// ... send email
}
}Batch processing and failure settlement
@BatchProcess() handlers return one result for each input job. Results must
contain unique matching jobId values and a boolean success; the queue checks
the complete result set before it settles any delivery. A failed result without
an Error receives a QueueError with code PROCESS_FAILED. Missing,
duplicate, unknown, or malformed results fail the batch instead of acknowledging
partial success.
continueOnError defaults to true. When enabled, the provider continues to
settle independent jobs after a job-level failure. For RabbitMQ and Kafka,
false stops after the first failed result and leaves later jobs for redelivery;
successful jobs already durably settled remain settled. BullMQ receives the
whole result set and uses each job's result to settle its worker callback.
Processor exceptions and invalid result sets are treated as batch failures.
retryFailedIndividually defaults to false. RabbitMQ and Kafka use it to
publish a failed job as a durable retry before acknowledging/resolving the
original delivery, while respecting the job's attempt limit and backoff. If
retry publication fails, the original delivery remains eligible for redelivery.
At attempt exhaustion RabbitMQ rejects to its configured dead-letter queue;
Kafka uses its DLQ when enableDLQ is enabled, otherwise it leaves the offset
unresolved for redelivery. BullMQ uses Bull's per-job attempts and rejection
semantics for failed batch results.
Kafka batch consumers disable automatic batch offset resolution and advance
only through a contiguous settled prefix. Delayed jobs use
delayedOffsetCommit: 'onExecute' by default: the offset remains a barrier until
the delayed job has been durably re-enqueued. The legacy 'immediate' mode
commits when the local timer is scheduled; that mode has weaker durability
because a process crash can lose the in-memory scheduled job.
RabbitMQ also retains delayed deliveries until durable handoff. When a message
carries a future x-process-after, the provider schedules it locally but does
not acknowledge the original delivery: the broker's copy is the only durable
copy until the delayed job is re-published and the publish is confirmed. Only
then is the original acked; if re-publication fails it is nacked with requeue. On
a crash before the timer fires, or on reconnect, the unacknowledged original is
redelivered by the broker. Re-publication preserves the job/correlation IDs,
source headers, and bounded attempt counters, and clears the delay headers so the
copy is not delayed again. This is at-least-once: an ack that fails after a
confirmed publication can duplicate work, deduplicated by the idempotency store.
Because retained deliveries stay unacknowledged, each one holds a unit of the
consumer's prefetch capacity. At a low prefetchCount — the default is 1 — this
is head-of-line blocking: a single long delay occupies the only slot and the
queue will not deliver another message until that delay elapses (or the broker's
consumer_timeout redelivers/closes the channel). For queues that mix long
delays with regular traffic, use a dedicated queue for delayed jobs, raise
prefetchCount to at least the number of concurrently delayed jobs you expect,
or use the broker-native delayed message exchange
(rabbitmq_delayed_message_exchange plugin) so the broker holds the delay instead
of an in-flight delivery. A retained delivery with an effective prefetchCount
of 1 logs a one-time warning at startup of the first delayed delivery.
noAck is deprecated and ignored: RabbitMQ consumers always use manual
acknowledgement, because automatic acknowledgement double-acks the manual
retry/DLQ paths — an ack issued after the broker already auto-acknowledged
targets a stale delivery tag and fails the whole channel. Setting
noAck: true logs a one-time warning per consumer registration and has no
effect. Delayed messages rely on the retained-delivery behaviour described
above.
To only enqueue jobs (no local processing), use producer-only mode:
QueueModule.forProducer({
provider: QueueProviderType.BULLMQ,
options: { connection: { host: 'redis', port: 6379 } },
});import { Injectable } from '@nestjs/common';
import { QueueService } from '@nathapp/nestjs-queue';
@Injectable()
export class EmailQueueClient {
constructor(private readonly queue: QueueService) {}
enqueue(to: string) {
return this.queue.add('email-queue', 'send', { to });
}
}Job payload validation
add and addBulk validate every job before it reaches the provider SDK.
The following defaults apply and are breaking for payloads over 1 MB:
- Payload size: serialized job data must not exceed 1 MB (1,048,576 bytes of JSON); larger payloads are rejected.
- Job ID length: an explicit
jobIdmust be at most 256 characters. - Job ID characters: an explicit
jobIdmay only contain[a-zA-Z0-9_-]. - Prototype-pollution guard: job data containing
__proto__,constructor, orprototypekeys is rejected.
Keep large blobs out of job data — pass a reference (e.g. an object-store key or URL) instead.
