NestJS中基于BullMQ的定时任务无法运行求助
问题:NestJS中BullMQ生产者与@nestjs/bull消费者无法协同工作
问题场景
在NestJS后端使用BullMQ实现定时任务时遇到以下问题:
- 生产者代码位于独立library中,启动后端时能打印
Event Cron Job Called日志,说明生产者已执行 - 消费者使用
@nestjs/bull注册,但始终无法触发任务执行,没有输出日志 - 生产者依赖
bullmq库,消费者依赖@nestjs/bull(底层基于bull库)
生产者代码
import { Queue } from 'bullmq'; import { BullQueues, BullWorkers } from '../../../../../apps/api-server/src/shared/utils/constants'; import { CronJobDataSource } from '../data-source'; import { PlatformIntegration } from '../../../../../apps/api-server/src/modules/platform-integrations/entities/platform-integration.entity'; import { CalendarPlatforms } from '../../../../../apps/api-server/src/modules/platform-integrations/domain/calendar-platforms.enum'; /* eslint-disable @typescript-eslint/no-var-requires */ require('dotenv').config(); async function getUsersToSyncWithPlatform(platform: CalendarPlatforms) { ... } export const eventCronJob = async () => { console.log('Event Cron Job Called'); try { await CronJobDataSource.initialize(); // Initialize the BullMQ queue const syncQueue = new Queue(BullQueues.SYNC_EVENTS, { connection: { host: process.env.REDIS_HOSTNAME, port: Number(process.env.REDIS_PORT) }, }); // Function to enqueue jobs const enqueueSyncJobs = async ( platform: CalendarPlatforms, users: { id: string; account: string; }[], period: number, ) => { for await (const user of users) { await syncQueue.add( BullWorkers.SYNC_EVENTS_FOR_PLATFORM, { platform, userId: user.id, account: user.account, }, { repeat: { every: period } }, ); } }; // Enqueue jobs for Google Calendar const usersToSyncGoogle = await getUsersToSyncWithPlatform(CalendarPlatforms.GOOGLE); await enqueueSyncJobs(CalendarPlatforms.GOOGLE, usersToSyncGoogle, 20000); process.exit(); } catch (error) { console.error(error); } }
消费者代码(原版本)
import { Process, Processor } from '@nestjs/bull'; import { Job } from 'bull'; @Processor(BullQueues.SYNC_EVENTS) export class SyncEventsConsumer { constructor() {} @Process(BullWorkers.SYNC_EVENTS_FOR_PLATFORM) async readOperationJob(job: Job<{ platform: CalendarPlatforms; userId: string; account: string }>) { const { data: { platform, userId, account }, } = job; console.log('Cron Job: >>>>>>>>>>>>>', platform, userId, account); }
解决方案
1. 核心原因:Bull与BullMQ不兼容
@nestjs/bull底层依赖的是bull库(v3/v4版本),而你使用的生产者依赖的是bullmq库——这是两个完全独立的项目,API相似但Redis存储的数据格式不兼容,因此消费者无法识别生产者创建的任务。
2. 统一使用BullMQ生态
将消费者替换为NestJS官方适配BullMQ的@nestjs/bullmq模块,步骤如下:
步骤1:安装正确的依赖
npm install @nestjs/bullmq bullmq # 或使用yarn yarn add @nestjs/bullmq bullmq
步骤2:修改AppModule的队列配置
用BullMQModule替代原有的BullModule:
import { BullMQModule } from '@nestjs/bullmq'; import { Module } from '@nestjs/common'; @Module({ imports: [ BullMQModule.forRoot({ connection: { host: process.env.REDIS_HOSTNAME, port: Number(process.env.REDIS_PORT), }, }), ], }) export class AppModule {}
步骤3:更新消费者代码
使用@nestjs/bullmq的装饰器,并导入bullmq的Job类型:
import { Process, Processor } from '@nestjs/bullmq'; import { Job } from 'bullmq'; // 注意此处导入bullmq的Job import { BullQueues, BullWorkers } from '../shared/utils/constants'; import { CalendarPlatforms } from '../modules/platform-integrations/domain/calendar-platforms.enum'; @Processor(BullQueues.SYNC_EVENTS) export class SyncEventsConsumer { constructor() {} @Process(BullWorkers.SYNC_EVENTS_FOR_PLATFORM) async readOperationJob(job: Job<{ platform: CalendarPlatforms; userId: string; account: string }>) { const { platform, userId, account } = job.data; console.log('Cron Job: >>>>>>>>>>>>>', platform, userId, account); } }
步骤4:在子模块中注册队列与消费者
如果之前用BullModule注册消费者,替换为BullMQModule的注册方式:
import { BullMQModule } from '@nestjs/bullmq'; import { Module } from '@nestjs/common'; import { SyncEventsConsumer } from './sync-events.consumer'; import { BullQueues } from '../shared/utils/constants'; @Module({ imports: [ BullMQModule.registerQueue({ name: BullQueues.SYNC_EVENTS, }), ], providers: [SyncEventsConsumer], }) export class YourSubModule {}
3. 额外优化建议
- 移除生产者中的
process.exit():这行代码会强制终止Node进程,可能导致任务未完全入队。建议在关闭队列连接后再退出:// 生产者代码末尾修改 await syncQueue.close(); // 关闭队列连接 process.exit(); - 清理Redis旧数据:如果之前用bull或bullmq创建过队列,执行
FLUSHDB(生产环境需谨慎)清除残留数据,避免格式冲突。 - 验证Redis连接:确保生产者和消费者的Redis主机、端口配置完全一致。
内容的提问来源于stack exchange,提问作者Eric Ryan
相关产品推荐
相关产品推荐

