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

如何通过RabbitMQ微服务实现主应用与WebSocket消息推送联动

实现RabbitMQ微服务与主应用WebSocket的消息互通

核心思路

让微服务通过RabbitMQ发送消息到主应用,主应用监听RabbitMQ队列,收到消息后通过WebSocket网关推送给客户端。


1. 主应用配置调整

1.1 集成RabbitMQ客户端

在主应用中添加RabbitMQ客户端依赖,同时启动WebSocket服务与RabbitMQ监听:

// 主应用 main.ts
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';

async function bootstrap() {
  const app = await NestFactory.create(AppModule);
  
  // 保留现有WebSocket相关配置(如跨域)
  app.enableCors();
  
  // 配置RabbitMQ客户端,监听微服务消息队列
  app.connectMicroservice<MicroserviceOptions>({
    transport: Transport.RMQ,
    options: {
      urls: ['amqp://localhost:5672'],
      queue: 'ws_notification_queue',
      queueOptions: { durable: false },
    },
  });

  await app.startAllMicroservices();
  await app.listen(3000);
}
bootstrap();

1.2 创建RabbitMQ消息处理服务

在主应用中新增服务,负责接收微服务消息并调用WebSocket网关推送:

// src/rabbitmq/rabbitmq.service.ts
import { Injectable } from '@nestjs/common';
import { MessagePattern } from '@nestjs/microservices';
import { WsGateway } from '../ws/ws.gateway';

@Injectable()
export class RabbitmqService {
  constructor(private readonly wsGateway: WsGateway) {}

  // 监听微服务发送的消息指令
  @MessagePattern('send_ws_message')
  handleWsMessage(data: { clientId: string; message: string }) {
    this.wsGateway.sendMessageToClient(data.clientId, data.message);
  }
}

1.3 增强WebSocket网关功能

给网关添加定向发送消息的方法,同时维护客户端连接映射:

// src/ws/ws.gateway.ts
import { WebSocketGateway, WebSocketServer, OnGatewayConnection, OnGatewayDisconnect } from '@nestjs/websockets';
import { Server, Socket } from 'socket.io';

@WebSocketGateway({ cors: true })
export class WsGateway implements OnGatewayConnection, OnGatewayDisconnect {
  @WebSocketServer()
  server: Server;

  // 存储客户端ID与socket实例的映射
  private clientConnections = new Map<string, Socket>();

  handleConnection(client: Socket) {
    const clientId = client.handshake.query.clientId as string;
    if (clientId) {
      this.clientConnections.set(clientId, client);
    }
  }

  handleDisconnect(client: Socket) {
    const clientId = client.handshake.query.clientId as string;
    if (clientId) {
      this.clientConnections.delete(clientId);
    }
  }

  // 新增:向指定客户端发送消息
  sendMessageToClient(clientId: string, message: string) {
    const targetClient = this.clientConnections.get(clientId);
    if (targetClient) {
      targetClient.emit('notification', message);
    }
  }
}

2. 微服务配置调整

2.1 微服务RabbitMQ客户端配置

修改微服务入口文件,配置RabbitMQ通信:

// 微服务 main.ts
import { NestFactory } from '@nestjs/core';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';
import { AppModule } from './app.module';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
    transport: Transport.RMQ,
    options: {
      urls: ['amqp://localhost:5672'],
      queue: 'ws_notification_queue',
      queueOptions: { durable: false },
    },
  });
  await app.listen();
}
bootstrap();

2.2 微服务消息发送服务

在微服务中创建服务,用于触发WebSocket消息推送请求:

// src/notification/notification.service.ts
import { Injectable } from '@nestjs/common';
import { Client, ClientProxy, Transport } from '@nestjs/microservices';

@Injectable()
export class NotificationService {
  @Client({
    transport: Transport.RMQ,
    options: {
      urls: ['amqp://localhost:5672'],
      queue: 'ws_notification_queue',
      queueOptions: { durable: false },
    },
  })
  private clientProxy: ClientProxy;

  // 调用该方法即可向指定客户端发送WebSocket消息
  async pushToClient(clientId: string, message: string) {
    await this.clientProxy.send('send_ws_message', { clientId, message }).toPromise();
  }
}

3. 验证步骤

  1. 启动本地RabbitMQ服务(确保5672端口正常运行)
  2. 启动主应用(端口3000)
  3. 启动微服务(无需指定额外端口,通过RabbitMQ与主应用通信)
  4. 客户端通过WebSocket连接主应用,携带唯一clientId参数
  5. 在微服务业务逻辑中调用pushToClient方法,传入目标clientId和消息内容,客户端即可收到WebSocket推送

内容的提问来源于stack exchange,提问作者Dart Vinnie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:05:11