如何通过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. 验证步骤
- 启动本地RabbitMQ服务(确保5672端口正常运行)
- 启动主应用(端口3000)
- 启动微服务(无需指定额外端口,通过RabbitMQ与主应用通信)
- 客户端通过WebSocket连接主应用,携带唯一
clientId参数 - 在微服务业务逻辑中调用
pushToClient方法,传入目标clientId和消息内容,客户端即可收到WebSocket推送
内容的提问来源于stack exchange,提问作者Dart Vinnie
相关产品推荐
相关产品推荐

