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

如何在独立进程运行的Bull队列Worker中发送Socket.IO事件?

如何在独立进程运行的Bull队列Worker中发送Socket.IO事件?

你遇到的这个问题其实非常常见——因为主进程和Worker进程是完全隔离的,内存不共享,所以Worker里根本拿不到主进程初始化的Socket.IO实例,直接调用getIO()自然会报错。下面给你介绍几种业界常用的解决方案,按推荐度由高到低排序:

1. 使用Socket.IO Redis适配器(最推荐)

既然你的Bull队列本身就依赖Redis,那直接用Socket.IO官方的Redis适配器是最省心的方案。它能让所有进程(主进程、Worker)通过Redis共享Socket.IO的房间、连接状态等信息,Worker可以直接发送Socket事件,事件会通过Redis广播到主进程,再推送给客户端。

步骤:

首先安装适配依赖(如果是Socket.IO v4+版本,用@socket.io/redis-adapter):

npm install @socket.io/redis-adapter redis

然后修改主进程的socket.ts,配置Redis适配器:

import { Server } from "socket.io";
import { createAdapter } from "@socket.io/redis-adapter";
import { createClient } from "redis";
import type { HttpServer } from "http";

let io: Server | null = null;

export function initSocket(server: HttpServer) {
  io = new Server(server, {
    cors: {
      origin: process.env.FRONTEND_URL,
      methods: ["GET", "POST"]
    },
  });

  // 创建Redis客户端
  const pubClient = createClient({ url: process.env.REDIS_URL });
  const subClient = pubClient.duplicate();
  
  // 连接Redis并设置适配器
  Promise.all([pubClient.connect(), subClient.connect()]).then(() => {
    io!.adapter(createAdapter(pubClient, subClient));
    console.log("Socket.IO Redis adapter initialized");
  });

  io.on("connection", (socket) => {
    console.log("Client connected:", socket.id);
    registerEvents(io!, socket);
  });

  return io;
}

export function getIO() {
  if (!io) {
    throw new Error("Socket.io not initialized. Call initSocket first.");
  }
  return io;
}

接下来修改Worker进程的worker.ts,自己初始化一个连接到Redis适配器的Socket.IO实例:

import { Server } from "socket.io";
import { createAdapter } from "@socket.io/redis-adapter";
import { createClient } from "redis";

// 初始化Worker的Socket.IO实例
let io: Server;

async function initWorkerIO() {
  const pubClient = createClient({ url: process.env.REDIS_URL });
  const subClient = pubClient.duplicate();
  
  await Promise.all([pubClient.connect(), subClient.connect()]);

  io = new Server({
    adapter: createAdapter(pubClient, subClient)
  });
}

// 先完成IO初始化
await initWorkerIO();

aiResponseQueue.process(5, async (job) => {
  const { conversationId, content } = job.data;
  try {
    // ... 你的任务处理逻辑 ...
    
    // 现在可以直接发送Socket事件,Redis会帮你广播到主进程
    io.to(conversationId).emit("ai-status", { status: "finished" });
  } catch (err) {
    console.error("Job processing failed:", err);
  }
});

这样一来,Worker和主进程通过Redis共享Socket.IO状态,发送的事件会自动同步到主进程的Socket服务,推送给对应的客户端。

2. Worker通过HTTP请求通知主进程

如果你的项目比较小,不想引入Redis适配器的复杂度,可以在主进程里加一个专用的API接口,Worker完成任务后调用这个接口,由主进程负责发送Socket事件。

步骤:

在主进程的app.ts中添加一个POST接口:

import { getIO } from "./socket";

// 新增通知客户端的接口
app.post("/api/notify-client", (req, res) => {
  const { conversationId, event, data } = req.body;
  try {
    const io = getIO();
    io.to(conversationId).emit(event, data);
    res.status(200).json({ success: true });
  } catch (err) {
    console.error("Failed to emit socket event:", err);
    res.status(500).json({ success: false, error: "Socket.IO not initialized" });
  }
});

然后在Worker里用axios或fetch调用这个接口:

import axios from "axios";

aiResponseQueue.process(5, async (job) => {
  const { conversationId, content } = job.data;
  try {
    // ... 你的任务处理逻辑 ...
    
    // 调用主进程的API发送事件
    await axios.post("http://localhost:8000/api/notify-client", {
      conversationId,
      event: "ai-status",
      data: { status: "finished" }
    });
  } catch (err) {
    console.error("Job processing failed:", err);
  }
});

这个方案实现简单,但如果Worker数量多、任务频率高,会给主进程带来额外的HTTP请求压力,适合小型项目使用。

3. 使用Redis Pub/Sub直接通信

因为Bull本身就依赖Redis,你也可以直接用Redis的发布订阅功能:Worker发布事件消息,主进程订阅这个频道,收到消息后再用自己的Socket.IO实例发送事件。

步骤:

修改主进程的socket.ts,添加Redis订阅逻辑:

import { createClient } from "redis";

// ... 原有的initSocket函数 ...
export function initSocket(server: HttpServer) {
  io = new Server(server, {
    cors: {
      origin: process.env.FRONTEND_URL,
      methods: ["GET", "POST"]
    },
  });

  // 订阅Redis频道
  const redisClient = createClient({ url: process.env.REDIS_URL });
  redisClient.connect().then(() => {
    redisClient.subscribe("socket-worker-events", (message) => {
      const { conversationId, event, data } = JSON.parse(message);
      io?.to(conversationId).emit(event, data);
    });
    console.log("Subscribed to Redis socket events channel");
  });

  io.on("connection", (socket) => {
    console.log("Client connected:", socket.id);
    registerEvents(io!, socket);
  });

  return io;
}

然后在Worker里发布消息到Redis频道:

import { createClient } from "redis";

// 创建Redis客户端
const redisClient = createClient({ url: process.env.REDIS_URL });
await redisClient.connect();

aiResponseQueue.process(5, async (job) => {
  const { conversationId, content } = job.data;
  try {
    // ... 你的任务处理逻辑 ...
    
    // 发布消息到Redis频道
    await redisClient.publish("socket-worker-events", JSON.stringify({
      conversationId,
      event: "ai-status",
      data: { status: "finished" }
    }));
  } catch (err) {
    console.error("Job processing failed:", err);
  } finally {
    // 记得断开Redis连接
    await redisClient.disconnect();
  }
});

这个方案复用了现有Redis依赖,比HTTP方案更轻量,但需要自己处理消息的序列化和反序列化,不如Socket.IO适配器省心。


内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 09:22:57