如何在独立进程运行的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

