调用Nest微服务@EventPattern时Socket.io连接未更新问题
问题:Nest微服务中调用@EventPattern时无法获取Socket.io已连接客户端
在调用NestMicroservice的@EventPattern时,尝试通过this.websocketService.emitToRoom(...)触发Socket.io事件,但此时Socket没有任何已连接客户端。通过GraphQL等其他模块测试时,Socket能正常获取到已连接客户端,功能符合预期。已确认客户端已成功连接Socket并加入指定房间。
相关代码
FCMTokenMicroserviceController
import { Controller } from '@nestjs/common'; import { EventPattern } from '@nestjs/microservices'; import WebsocketEvent from 'src/common/interfaces/websocket/WebsocketEvent'; import { WebsocketRoom } from 'src/common/interfaces/websocket/WebsocketRoom'; import { FCMTokenService } from 'src/modules/fcm-token/fcm.token.service'; import { NotificationService } from 'src/modules/notification/notification.service'; import { WebsocketService } from 'src/modules/websocket/websocket.service'; @Controller() export class FCMTokenMicroserviceController { constructor( private readonly fcmTokenService: FCMTokenService, private readonly notificationService: NotificationService, private readonly websocketService: WebsocketService ) {} @EventPattern({ cmd: 'sendNotificationByUserIds', service: 'fcm-token' }) async sendNotificationByUserId( messages: { userId: number, notification: INotificationPayload, data?: Record<string, any> }[] ): Promise<void> { for (const message of messages) { const notification = await this.notificationService.save({ payload: message.data, notification: message.notification, userId: message.userId }); // 此处Socket无已连接客户端 this.websocketService.emitToRoom(WebsocketRoom.notification, WebsocketEvent.notification(message.userId), notification); await this.fcmTokenService.sendNotificationByUserIds({ userIds: [message.userId], notification: message.notification, data: message.data }); } } }
WebsocketService
import { Injectable } from '@nestjs/common'; import { WebsocketGateway } from './websocket.gateway'; import { WebsocketRoom } from 'src/common/interfaces/websocket/WebsocketRoom'; @Injectable() export class WebsocketService { constructor( private readonly websocketGateway: WebsocketGateway ) {} async emitToRoom(room: WebsocketRoom, event: string, data: any) { // 微服务调用时此处无客户端 console.log('\n', await this.websocketGateway.server.in(room).fetchSockets()); console.log('\n', this.websocketGateway.server.listenerCount(event), this.websocketGateway.server.listeners(event)); return this.websocketGateway.server.to(room).emit(event, data); } }
WebsocketGateway
import { UseGuards } from "@nestjs/common"; import { ConnectedSocket, MessageBody, SubscribeMessage, WebSocketGateway, WebSocketServer } from "@nestjs/websockets"; import { Server, Socket } from "socket.io"; import { WSGuard } from "src/common/guard/WSGuard"; @WebSocketGateway({ cors: true, path: '/ws', }) export class WebsocketGateway { @WebSocketServer() server: Server; @UseGuards(WSGuard) @SubscribeMessage('joinRooms') joinRoom( @ConnectedSocket() client: Socket, @MessageBody() body: JoinRooms ) { client.join(body.rooms); } }
问题原因与解决方案
核心原因
微服务实例和WebSocket网关所在的HTTP应用实例是两个独立进程,它们的Socket.io服务器完全隔离。微服务进程中的网关实例没有任何客户端连接,自然无法触发事件到已连接的客户端。
解决方案
方案1:使用Socket.io Redis适配器共享连接状态
通过Redis让多个进程共享Socket.io的连接、房间等状态,这样微服务进程的Socket.io操作能同步到HTTP应用的客户端:
- 安装依赖:
npm install @socket.io/redis-adapter redis
- 配置Redis适配器到WebSocketGateway:
import { Injectable, OnModuleInit } from '@nestjs/common'; import { WebSocketGateway, WebSocketServer } from '@nestjs/websockets'; import { Server } from 'socket.io'; import { createAdapter } from '@socket.io/redis-adapter'; import { createClient } from 'redis'; @WebSocketGateway({ cors: true, path: '/ws', }) export class WebsocketGateway implements OnModuleInit { @WebSocketServer() server: Server; async onModuleInit() { // 连接Redis const pubClient = createClient({ url: 'redis://localhost:6379' }); const subClient = pubClient.duplicate(); await Promise.all([pubClient.connect(), subClient.connect()]); // 设置Redis适配器 this.server.adapter(createAdapter(pubClient, subClient)); } // ...原有joinRoom方法 }
- 确保微服务和HTTP应用都连接同一个Redis实例,这样两边的Socket.io状态就能共享。
方案2:通过消息队列解耦通信
让微服务处理完逻辑后发送消息到消息队列(如Redis MQ、NATS),HTTP应用监听队列消息,收到后触发WebSocket事件:
- 微服务在
@EventPattern方法中,处理完数据后发送消息到队列:
// 引入Redis客户端或其他MQ客户端 async sendNotificationByUserId(messages: ...) { for (const message of messages) { // ...原有保存通知逻辑 // 发送消息到队列 await redisClient.publish('notification_events', JSON.stringify({ room: WebsocketRoom.notification, event: WebsocketEvent.notification(message.userId), data: notification })); // ...原有FCM发送逻辑 } }
- 在HTTP应用的WebsocketService中监听队列消息:
@Injectable() export class WebsocketService implements OnModuleInit { constructor( private readonly websocketGateway: WebsocketGateway, private readonly redisClient: RedisClient ) {} async onModuleInit() { this.redisClient.subscribe('notification_events', (message) => { const { room, event, data } = JSON.parse(message); this.websocketGateway.server.to(room).emit(event, data); }); } // ...原有emitToRoom方法 }
内容的提问来源于stack exchange,提问作者Melvin Jovano
相关产品推荐
相关产品推荐

