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

NestJS接收MQTT消息后向Broker发送ACK的实现方案

问题描述

我用NestJS开发了一个与MQTT Broker通信的应用,用于接收设备发布的消息。需求是成功接收消息后,向Broker的特定主题发送ACK响应,避免消息重复推送。

目前遇到的核心问题:

  • 通过主文件创建微服务时,能正常监听接收消息,但找不到发送响应的方式。
  • 按文档创建代理客户端时,无法使用同一个clientID(Broker同一时间仅允许一个客户端用该ID连接);换用不同ID的话,发送的ACK无法对应接收消息的客户端,导致消息重复推送。
  • 尝试不在主文件配置连接,仅在模块内使用客户端并在onApplicationBootstrap中连接,此时控制器完全无法接收消息。

希望找到能同时实现监听和发送消息的配置方案。


现有代码

main.ts

import { NestFactory } from '@nestjs/core';
import { MqttOptions, Transport } from '@nestjs/microservices';
import { AppModule } from './app.module';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<MqttOptions>(AppModule, {
    transport: Transport.MQTT,
    options: {
      url: 'mqtt://XX.XXX.XXX.XXXX:1883',
      clientId: 'my-client-id-test-001',
    },
  });
  await app.listen();
}
bootstrap();

app.module.ts

import { Module } from '@nestjs/common';
import {
  ClientsModule,
  OutgoingResponse,
  Transport,
} from '@nestjs/microservices';
import { AppController } from './app.controller';

@Module({
  imports: [
    ClientsModule.register([
      {
        name: 'MQTT_CLIENT',
        transport: Transport.MQTT,
        options: {
          url: 'mqtt://XX.XXX.XXX.XXX:1883',
          clientId: 'my-client-id-test-001',
          serializer: {
            serialize: (value: any): OutgoingResponse => value.data,
          },
          clean: false,
        },
      },
    ]),
  ],
  controllers: [AppController],
})
export class AppModule {}

app.controller.ts

import { Controller, Inject, OnApplicationBootstrap } from '@nestjs/common';
import {
  ClientProxy,
  Ctx,
  MessagePattern,
  MqttContext,
  Payload,
} from '@nestjs/microservices';

import { Message } from 'src/Message';

@Controller()
export class AppController implements OnApplicationBootstrap {
  constructor(@Inject('MQTT_CLIENT') private client: ClientProxy) {}

  async onApplicationBootstrap() {
    await this.client.connect();
  }

  @MessagePattern('GW/GPUB/682719248464')
  getPublishMessages(@Ctx() context: MqttContext, @Payload() payload: string) {
    console.log('recived data...');
    const message = new Message(payload);
    this.sendAck(
      'GW/SACK/682719248464',
      `ti=0F:${message.packetId}&id=${message.gatewayId}`,
    );
  }

  private sendAck(pattern: string, payload: string) {
    return this.client.send(pattern, payload);
  }
}

解决方案

核心思路是复用微服务的MQTT客户端发送ACK,不需要额外创建代理客户端,确保用同一个clientID连接Broker,解决ACK上下文不匹配的问题。

步骤1:调整main.ts启动逻辑

保留微服务连接,同时可选启动HTTP应用(若不需要HTTP可省略):

import { NestFactory } from '@nestjs/core';
import { MqttOptions, Transport } from '@nestjs/microservices';
import { AppModule } from './app.module';

async function bootstrap() {
  // 创建HTTP应用(可选,不需要可删除)
  const app = await NestFactory.create(AppModule);
  
  // 连接MQTT微服务
  const mqttMicroservice = app.connectMicroservice<MqttOptions>({
    transport: Transport.MQTT,
    options: {
      url: 'mqtt://XX.XXX.XXX.XXXX:1883',
      clientId: 'my-client-id-test-001',
      clean: false, // 保持会话,避免离线消息重复推送
    },
  });

  await mqttMicroservice.listen();
  await app.listen(3000); // 可选HTTP端口
}
bootstrap();

步骤2:简化app.module.ts

删除多余的ClientsModule注册,不需要额外客户端:

import { Module } from '@nestjs/common';
import { AppController } from './app.controller';

@Module({
  controllers: [AppController],
})
export class AppModule {}

步骤3:修改控制器,复用监听客户端发送ACK

通过MqttContext获取底层MQTT客户端实例,直接调用publish发送ACK:

import { Controller } from '@nestjs/common';
import { Ctx, MessagePattern, MqttContext, Payload } from '@nestjs/microservices';
import { Message } from 'src/Message';

@Controller()
export class AppController {
  @MessagePattern('GW/GPUB/682719248464')
  getPublishMessages(@Ctx() context: MqttContext, @Payload() payload: string) {
    console.log('recived data...');
    const message = new Message(payload);
    
    // 获取监听用的MQTT客户端实例
    const mqttClient = context.getClient();
    
    // 发布ACK到指定主题
    mqttClient.publish(
      'GW/SACK/682719248464',
      `ti=0F:${message.packetId}&id=${message.gatewayId}`,
      (err) => {
        if (err) {
          console.error('ACK发送失败:', err);
        } else {
          console.log('ACK发送成功');
        }
      }
    );
  }
}

关键说明
  1. 复用客户端:通过MqttContext.getClient()拿到的是微服务监听用的同一个MQTT客户端,确保clientID一致,符合Broker的单连接限制,同时ACK的会话上下文与接收消息的上下文匹配,避免重复推送。
  2. 直接调用publish:MQTT是单向通信协议,用底层客户端的publish方法比NestJS的client.send()更直接(send()是针对RPC模式的封装)。
  3. 会话保持:设置clean: false让Broker保存会话状态,客户端重连后能正确处理离线消息和ACK关联。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:05:39