如何避免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事件,其他实例只处理任务,不监听事件。
- 在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' } } ] };
- 在代码中判断环境变量,仅在指定实例中注册QueueEvents监听类:
// 初始化应用时判断 if (process.env.IS_EVENT_LISTENER === 'true') { app.get(ArchiveR0QueueEvents); }
- 启动时指定监听实例:
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
相关产品推荐
相关产品推荐

