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

如何避免PM2集群模式下BullMQ重复创建任务

问题:PM2集群模式下BullMQ任务完成事件被多实例重复触发

我用多个队列构建任务流,比如把包含两个操作的任务拆分为Queue A和Queue B。收到HTTP请求后在Queue A中创建任务;Queue A完成后,再创建Queue B的任务。该流程在开发环境运行正常,但使用PM2集群模式启动多线程时出现问题:第一步仅一个Node实例响应HTTP请求并创建Queue A任务(符合预期),但Queue A的completed事件被多个BullMQ Worker同时监听,导致Queue B中被重复创建任务。

相关代码示例:

@QueueEventsListener('archiveR0Queue')
export class ArchiveR0QueueEvents extends BaseQueueEvents {
  constructor(processService: ProcessService) {
    super(processService, ArchiveR0QueueEvents.name)
  }
  @OnQueueEvent('completed')
  async onCompleted(job: { jobId: string; returnvalue: any; prev?: string }) {
    const { orderId } = job.returnvalue
    if (orderId) {
      await this.processService.queueManagementService.handlePendingJob(
        orderId,
        job.returnvalue,
      )
    }
  }
}

解决方案

方案1:指定单实例监听队列事件

在PM2集群模式下,仅让一个实例负责监听队列的completed事件,其他实例只处理任务,不监听事件。

  1. 在PM2配置文件(如ecosystem.config.js)中,给指定实例设置专属环境变量:
module.exports = {
  apps: [
    {
      name: 'your-app',
      script: 'dist/main.js',
      instances: 'max',
      exec_mode: 'cluster',
      // 默认实例不监听事件
      env: {},
      // 监听实例的环境变量
      env_event_listener: {
        IS_EVENT_LISTENER: 'true'
      }
    }
  ]
};
  1. 在代码中判断环境变量,仅在指定实例中注册QueueEvents监听类:
// 初始化应用时判断
if (process.env.IS_EVENT_LISTENER === 'true') {
  app.get(ArchiveR0QueueEvents);
}
  1. 启动时指定监听实例:
pm2 start ecosystem.config.js --env event_listener

方案2:基于Redis实现分布式锁

在处理completed事件时,先通过Redis获取分布式锁,只有拿到锁的实例才能执行创建Queue B任务的逻辑,从根源避免重复执行。

import { createClient } from 'redis';

// 初始化Redis客户端
const redisClient = createClient({
  url: 'redis://your-redis-host:6379'
});
await redisClient.connect();

@QueueEventsListener('archiveR0Queue')
export class ArchiveR0QueueEvents extends BaseQueueEvents {
  constructor(processService: ProcessService) {
    super(processService, ArchiveR0QueueEvents.name)
  }
  @OnQueueEvent('completed')
  async onCompleted(job: { jobId: string; returnvalue: any; prev?: string }) {
    const { orderId } = job.returnvalue
    if (!orderId) return;

    // 生成唯一锁键,确保同一orderId的任务只处理一次
    const lockKey = `job:archive:${orderId}:lock`;
    // 设置锁,过期时间10秒(根据实际任务处理时长调整)
    const lockAcquired = await redisClient.set(lockKey, 'locked', {
      NX: true, // 仅当键不存在时设置
      EX: 10    // 自动过期时间
    });

    // 未拿到锁,直接返回
    if (!lockAcquired) return;

    try {
      // 执行创建Queue B任务的逻辑
      await this.processService.queueManagementService.handlePendingJob(
        orderId,
        job.returnvalue,
      )
    } finally {
      // 无论执行成功与否,释放锁
      await redisClient.del(lockKey);
    }
  }
}

方案3:校验Queue B中是否已有重复任务

在创建Queue B任务前,先查询Queue B中是否存在同一orderId的待处理/活跃任务,若存在则跳过创建。

@QueueEventsListener('archiveR0Queue')
export class ArchiveR0QueueEvents extends BaseQueueEvents {
  constructor(processService: ProcessService) {
    super(processService, ArchiveR0QueueEvents.name)
  }
  @OnQueueEvent('completed')
  async onCompleted(job: { jobId: string; returnvalue: any; prev?: string }) {
    const { orderId } = job.returnvalue
    if (!orderId) return;

    // 获取Queue B的实例
    const queueB = this.processService.queueManagementService.getQueue('queueB');
    // 查询Queue B中处于等待/活跃状态的任务
    const existingJobs = await queueB.getJobs(['waiting', 'active'], 0, 1, true);
    // 检查是否已有同一orderId的任务
    const hasDuplicateJob = existingJobs.some(j => j.data.orderId === orderId);

    if (hasDuplicateJob) return;

    // 执行创建Queue B任务的逻辑
    await this.processService.queueManagementService.handlePendingJob(
      orderId,
      job.returnvalue,
    )
  }
}

内容的提问来源于stack exchange,提问作者Junglebook

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 00:49:58