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

调用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应用的客户端:

  1. 安装依赖:
npm install @socket.io/redis-adapter redis
  1. 配置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方法
}
  1. 确保微服务和HTTP应用都连接同一个Redis实例,这样两边的Socket.io状态就能共享。

方案2:通过消息队列解耦通信

让微服务处理完逻辑后发送消息到消息队列(如Redis MQ、NATS),HTTP应用监听队列消息,收到后触发WebSocket事件:

  1. 微服务在@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发送逻辑
  }
}
  1. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 05:25:33