You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 03:45:08