NestJS微服务:如何捕获消息代理断开连接并处理重连?
NestJS微服务中消息代理重连事件的处理方式
在NestJS微服务中,针对消息代理(以MQTT为例)的重连、断开等事件,可通过以下几种方式实现自定义处理逻辑:
1. 自定义客户端实例并绑定事件监听
手动创建消息代理客户端实例,直接监听其原生事件,再将该客户端绑定到NestJS微服务:
import { NestFactory } from '@nestjs/core'; import { MicroserviceOptions, Transport } from '@nestjs/microservices'; import { MqttClient } from 'mqtt'; async function bootstrap() { // 初始化MQTT客户端 const mqttClient = require('mqtt').connect('mqtt://localhost:1883'); // 监听重连事件 mqttClient.on('reconnect', () => { console.log('MQTT客户端正在重连...'); // 这里可添加Sentry通知、日志上报等逻辑 }); // 监听连接断开事件 mqttClient.on('close', () => { console.log('MQTT客户端已断开连接'); }); // 创建微服务并注入自定义客户端 const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, { transport: Transport.MQTT, options: { client: mqttClient, }, }); await app.listen(); } bootstrap();
2. 通过ClientProxy获取底层客户端实例监听事件
如果在服务内部使用ClientProxy进行消息交互,可在模块初始化时获取底层客户端实例,再绑定事件:
import { Injectable, OnModuleInit } from '@nestjs/common'; import { ClientProxy, ClientProxyFactory, Transport } from '@nestjs/microservices'; import { MqttClient } from 'mqtt'; @Injectable() export class MessageService implements OnModuleInit { private client: ClientProxy; private mqttClient: MqttClient; constructor() { this.client = ClientProxyFactory.create({ transport: Transport.MQTT, options: { url: 'mqtt://localhost:1883', }, }); } async onModuleInit() { // 获取底层MQTT客户端实例 this.mqttClient = await this.client.connect() as unknown as MqttClient; // 绑定重连事件处理 this.mqttClient.on('reconnect', () => { console.log('MQTT客户端触发重连'); // 执行业务恢复、告警通知等操作 }); // 可选:监听连接成功事件 this.mqttClient.on('connect', () => { console.log('MQTT客户端连接成功'); }); } }
其他消息代理的适配逻辑
对于RabbitMQ、Kafka等其他消息代理,核心逻辑一致:
- 获取对应代理的底层客户端实例
- 监听其原生的断开、重连类事件(如RabbitMQ的
close事件,Kafka的disconnect事件) - 在事件回调中实现自定义处理逻辑
内容的提问来源于stack exchange,提问作者Vahid Najafi
相关产品推荐
相关产品推荐

