MQTTNet 4.x版本:如何向在线客户端发送消息修改主题值?
在MQTTNet 4.x中实现服务器端客户端状态查询与消息发布
MQTTNet 4.x对服务器端API做了重构,移除了旧的IMqttServer接口,改用IMqttServerInstance作为核心操作接口,同时客户端状态、会话信息的获取方式也有调整。以下是适配4.x版本的实现方案:
一、调整服务注册(ASP.NET Core .NET 6)
简化原注册代码,确保正确注入服务器实例:
var optionBuilder = new MqttServerOptionsBuilder() .WithDefaultEndpoint() .WithDefaultCommunicationTimeout(TimeSpan.FromMilliseconds(5000)) .Build(); // 注册MQTT服务器及相关适配器 builder.Services .AddHostedMqttServer(optionBuilder) .AddMqttConnectionHandler() .AddConnections() .AddMqttTcpServerAdapter() .AddMqttWebSocketServerAdapter();
注意:无需重复调用AddMqttConnectionHandler(),一次即可完成注册。
二、实现服务器端操作服务
通过注入IMqttServerInstance替代旧的IMqttServer,实现客户端状态查询、消息发布等核心功能:
using MQTTnet; using MQTTnet.Server; public class MqttBrokerService { private readonly IMqttServerInstance _mqttServer; public MqttBrokerService(IMqttServerInstance mqttServer) { _mqttServer = mqttServer ?? throw new ArgumentNullException(nameof(mqttServer)); } // 获取当前连接的客户端状态 public async Task<IEnumerable<MqttClientConnection>> GetClientStatusAsync() { return await _mqttServer.GetClientsAsync(); } // 获取会话状态(包含订阅、保留消息等信息) public async Task<IEnumerable<MqttSession>> GetSessionStatusAsync() { return await _mqttServer.GetSessionsAsync(); } // 清除所有保留消息 public Task ClearRetainedApplicationMessagesAsync() { return _mqttServer.ClearRetainedMessagesAsync(); } // 获取所有保留消息 public Task<IEnumerable<MqttApplicationMessage>> GetRetainedApplicationMessagesAsync() { return _mqttServer.GetRetainedMessagesAsync(); } // 发布消息到指定主题(客户端订阅后会收到) public async Task<MqttServerPublishResult> PublishAsync(MqttApplicationMessage applicationMessage) { if (applicationMessage == null) { throw new ArgumentNullException(nameof(applicationMessage)); } return await _mqttServer.PublishAsync(applicationMessage, CancellationToken.None); } }
在Program.cs中注册该服务:
builder.Services.AddScoped<MqttBrokerService>();
三、在API控制器中使用示例
创建API控制器暴露这些功能,方便调用:
using Microsoft.AspNetCore.Mvc; using MQTTnet; [ApiController] [Route("api/mqtt")] public class MqttController : ControllerBase { private readonly MqttBrokerService _mqttBrokerService; public MqttController(MqttBrokerService mqttBrokerService) { _mqttBrokerService = mqttBrokerService; } // 获取所有连接的客户端信息 [HttpGet("clients")] public async Task<IActionResult> GetClients() { var clients = await _mqttBrokerService.GetClientStatusAsync(); var clientInfos = clients.Select(c => new { c.ClientId, c.IsConnected, c.ConnectionStartTimestamp, SessionStartTime = c.Session?.SessionStartTimestamp }); return Ok(clientInfos); } // 发布消息到指定主题 [HttpPost("publish")] public async Task<IActionResult> Publish([FromBody] PublishRequest request) { var message = new MqttApplicationMessageBuilder() .WithTopic(request.Topic) .WithPayload(request.Payload) .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .WithRetainFlag(request.Retain) .Build(); var result = await _mqttBrokerService.PublishAsync(message); return Ok(new { result.ReasonCode, result.IsSuccess }); } } public class PublishRequest { public string Topic { get; set; } public string Payload { get; set; } public bool Retain { get; set; } }
四、关于ManagedMqttClient的说明
你之前尝试的MQTTNet.Extensions.ManagedClient是客户端侧的工具库,用于实现自动重连的客户端逻辑,并不适用于服务器端操作自身的客户端会话。服务器端直接通过IMqttServerInstance即可完成所有所需操作。
内容的提问来源于stack exchange,提问作者itprodavets
相关产品推荐
相关产品推荐

