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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 05:54:17