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

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,提问作者呂學霖

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:35:28