Nest.js ClientProxy连接RabbitMQ后无法重连问题求助
问题描述
使用Nest.js实现微服务时遇到以下问题:
- 需求:向RabbitMQ发送消息,同时将消息日志保存到MongoDB;ClientProxy连接正常时,MongoDB中
sendMessageStatus设为ok;连接异常时设为fail。 - 问题现象:关闭RabbitMQ服务器后发送消息,MongoDB能正确记录
sendMessageStatus: fail;但重启RabbitMQ后再次发送消息,日志状态仍为fail,说明ClientProxy未自动重连。已知Nest.js官方文档提到ClientProxy是懒加载的,需要实现ClientProxy的自动重连逻辑。
相关代码如下:
jobQueue.module.ts
return ClientProxyFactory.create({ transport: Transport.RMQ, options: { urls: [`amqp://${account}:${password}@${IP}:${port}`], queue: outputQueueName, serializer: { serialize: value => value.data, }, noAck: false, persistent: true, queueOptions: { durable: true, } } });
jobQueue.service.ts
constructor( @Inject(CONNECTION_NAME) private readonly client: ClientRMQ, ) {}; async sendMessage(data: SendMessageDto) { try { this.logger.serviceDebug(SENDMESSAGE_METHOD); data.id = this.messageID++; return await this.client.connect() .then(() => { return this.client.emit('', data) }).catch(err => { return this.client.emit('', data) .pipe( catchError(connectionError => { throw connectionError; }) ); }); } catch (err) { console.log('catch in job', err); throw err; }; };
调用
client.connect()并无效果。
myService.service.ts
const messageObserver = await this.jobQueueService.sendMessage(MQCLI); const createdLog: CreateScheduleExecutionLogDto = { ...data, scheduleID: scheduleID, schedule: item, processDatetime: new Date(), }; messageObserver.subscribe({ next: x => { console.log(x); createdLog.processStatus = OK; this.scheduleExecutionLogModel.create(createdLog); }, error: e => { console.log(e); createdLog.processStatus = ERROR; this.scheduleExecutionLogModel.create(createdLog); }, })
解决方案
1. 添加RabbitMQ自动重连配置
在RMQ的options中添加重连相关参数,让ClientProxy自动尝试重连:
return ClientProxyFactory.create({ transport: Transport.RMQ, options: { urls: [`amqp://${account}:${password}@${IP}:${port}`], queue: outputQueueName, serializer: { serialize: value => value.data, }, noAck: false, persistent: true, queueOptions: { durable: true, }, // 新增重连配置 retryDelay: 5000, // 每次重连间隔5秒 retryAttempts: Number.MAX_SAFE_INTEGER, // 无限重试 connectionOptions: { heartbeat: 60, // 心跳检测,保持连接活性 reconnect: true, // 开启自动重连 } } });
2. 修正sendMessage方法逻辑
ClientProxy本身是懒加载的,无需手动调用connect(),它会在首次发送消息时自动建立连接;连接失败后,会根据配置的重连参数自动重试。修改后的方法如下:
async sendMessage(data: SendMessageDto) { this.logger.serviceDebug(SENDMESSAGE_METHOD); data.id = this.messageID++; // 直接返回emit的Observable,无需手动处理connect return this.client.emit('', data).pipe( catchError(err => { console.log('send message error:', err); throw err; }) ); };
3. 确保错误订阅逻辑生效
原订阅逻辑无需修改,当ClientProxy重连成功后,后续的消息发送会自动触发next回调,此时MongoDB会记录ok状态。
核心原理说明
ClientProxy的懒加载特性意味着它不会提前建立连接,只有当调用emit()/send()等方法时才会尝试连接。通过配置retryDelay、retryAttempts和connectionOptions.reconnect,可以让ClientProxy在连接断开后自动重试,直到连接恢复。手动调用connect()会提前建立连接,但不会处理后续的重连逻辑,因此不需要手动调用。
内容的提问来源于stack exchange,提问作者呂學霖
相关产品推荐
相关产品推荐

