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

如何将BullMQ队列事件仅推送给对应任务的所属用户

核心问题点

你的定向推送失效、消息全量发送是三个逻辑错误共同导致的:

  • 全局socket实例赋值逻辑完全错误:this.socket = this.socket == null ? socket : this.socket只会把第一个连接到服务的用户socket存在全局变量里,后续所有新用户连接都不会更新这个值,所有队列事件最终都只会通过第一个用户的连接发送,必然出现消息错发。
  • BullMQ事件监听绑定位置错误:你把integrationQueueEvents.on("progress")写在了connection连接回调内部,每有一个用户连入服务就会新增一个重复的事件监听器,队列触发进度更新时会同时执行N个监听回调,不仅会重复推送消息,闭包缓存的旧socket实例还会进一步导致推送目标错乱。
  • 推送逻辑和任务归属没有绑定:你现在的代码里,推送目标的userID是从全局/闭包里的socket拿的,根本和当前触发进度的BullMQ任务没有关联,自然没法把任务进度推给对应的发起用户。另外你用socket.to(socket.id).emit()本身用法就错了:to()方法是广播接口,会自动排除调用它的当前socket,你拿当前socket调用to给自己发,当然收不到消息。
修复方案

你已经实现了用户连接时自动加入以自身userID命名的房间,配合Redis适配器,直接用房间做定向推送是多实例部署下最稳定的方案,不需要依赖临时的socketID做推送,按以下步骤改即可:

  1. 清理connection回调里的错误代码
    直接删掉这两行错误逻辑:
    // 删掉全局socket赋值
    // this.socket = this.socket == null ? socket : this.socket;
    
    // 删掉connection回调里所有绑定integrationQueueEvents的代码,监听器不能重复绑定
    
    补充断开连接时的会话状态更新逻辑,避免在线用户列表不准:
    socket.on("disconnect", async () => {
      await this.sessionStore.saveSession(socket.sessionID, {
        userID: socket.userID,
        username: socket.username,
        connected: false
      });
      socket.broadcast.emit("user disconnected", {
        userID: socket.userID,
        username: socket.username,
        connected: false
      });
    });
    
  2. 任务入队时绑定归属用户ID
    你必须在往BullMQ队列添加任务的时候,就把发起任务的用户ID存在任务数据里,不然后续队列触发事件时,根本没法知道这个任务属于哪个用户:
    // 任务入队逻辑示例
    await integrationQueue.add("yourTaskName", {
      // 你原有的业务字段
      ...originalTaskData,
      // 新增字段,标记任务所属用户
      belongUserID: 发起任务的当前登录用户ID
    });
    
  3. 全局仅绑定一次BullMQ事件监听
    把队列事件监听的逻辑移到Socket.IO实例初始化完成的位置,绝对不要放在connection回调里,触发事件时直接向对应用户的房间推送消息即可,Redis适配器会自动跨实例路由消息到用户的连接上:
    // 服务初始化阶段执行一次即可
    integrationQueueEvents.on("progress", async (job: any) => {
      try {
        console.log("Job Progressing", job);
        const targetUserID = job.data.belongUserID;
        if (!targetUserID) return;
    
        const payload = {
          status: true,
          data: job.data,
          jobId: job.jobId
        };
    
        console.log("推送进度到用户", targetUserID, payload);
        // 向对应用户的房间发消息,自动适配多实例部署
        this.instance.to(targetUserID).emit("integrationProgress", payload);
      } catch (error) {
        console.log("队列进度推送报错", error);
      }
    });
    
    任务完成、失败的事件监听也按同样逻辑写即可,全部全局绑定一次,通过job里存的belongUserID定向推到对应房间。
补充说明

socket.to(socketId).emit()的设计是用来给其他连接发消息的,会自动排除调用该方法的socket本身,所以你尝试用这个方法给当前socket对应的用户发消息必然失效。如果是单连接内的即时消息,直接用socket.emit()即可,但涉及到异步队列这种脱离了连接请求上下文的场景,用用户ID作为房间名做定向推送是最可靠的方案,不需要关心用户当前的连接状态、连接到了哪个服务节点、socketID是否变化,只要用户在线就能收到消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 16:18:21