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发送成功'); } } ); } }
关键说明
- 复用客户端:通过
MqttContext.getClient()拿到的是微服务监听用的同一个MQTT客户端,确保clientID一致,符合Broker的单连接限制,同时ACK的会话上下文与接收消息的上下文匹配,避免重复推送。 - 直接调用publish:MQTT是单向通信协议,用底层客户端的
publish方法比NestJS的client.send()更直接(send()是针对RPC模式的封装)。 - 会话保持:设置
clean: false让Broker保存会话状态,客户端重连后能正确处理离线消息和ACK关联。
内容的提问来源于stack exchange,提问作者Jorge
相关产品推荐
相关产品推荐

