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

.NET 6中MQTT Broker启用WebSocket及消息推送问题咨询

解决方案:.NET 6中自定义MQTT Broker的WebSocket支持及消息推送触发

问题1:启用WebSocket支持

你的配置存在重复创建MQTT Server实例的问题:既通过AddHostedMqttServerWithServices注册了托管的MQTT Server,又在MqttBrokerService中自行new MqttFactory().CreateMqttServer(),导致WebSocket适配器无法关联到正确的服务实例,最终连接验证失败。

修复步骤:

  1. 重构MqttBrokerService,移除自行创建Server的逻辑,改为依赖注入IMqttServer:
public class MqttBrokerService : IDisposable, IHostedService, IMqttBrokerService
{
    private readonly IMqttServer _mqttServer;

    // 依赖注入托管的MQTT Server实例
    public MqttBrokerService(IMqttServer mqttServer)
    {
        _mqttServer = mqttServer;
        // 绑定自定义消息/连接处理逻辑
        _mqttServer.ClientConnectedAsync += OnClientConnectedAsync;
        _mqttServer.ApplicationMessageReceivedAsync += OnApplicationMessageReceivedAsync;
    }

    public Task StartAsync(CancellationToken cancellationToken)
    {
        return _mqttServer.StartAsync(cancellationToken);
    }

    public Task StopAsync(CancellationToken cancellationToken)
    {
        return _mqttServer.StopAsync(cancellationToken);
    }

    public void Dispose()
    {
        _mqttServer.ClientConnectedAsync -= OnClientConnectedAsync;
        _mqttServer.ApplicationMessageReceivedAsync -= OnApplicationMessageReceivedAsync;
        (_mqttServer as IDisposable)?.Dispose();
    }

    // 自定义连接验证逻辑示例
    private Task OnClientConnectedAsync(MqttClientConnectedEventArgs args)
    {
        // 你的连接校验或初始化逻辑
        return Task.CompletedTask;
    }

    // 自定义消息接收逻辑示例
    private Task OnApplicationMessageReceivedAsync(MqttApplicationMessageReceivedEventArgs args)
    {
        // 你的消息处理逻辑
        return Task.CompletedTask;
    }
}
  1. 正确配置服务与中间件,确保WebSocket适配器与托管Server关联:
var builder = WebApplication.CreateBuilder(args);
var services = builder.Services;

// 1. 注册托管MQTT Server及基础配置
services.AddHostedMqttServerWithServices(optionsBuilder =>
{
    optionsBuilder
        .WithDefaultEndpoint() // 保留TCP端点(可选)
        .WithDefaultEndpointPort(1883);
});

// 2. 添加WebSocket适配器和连接处理
services.AddMqttConnectionHandler();
services.AddMqttWebSocketServerAdapter();
services.AddConnections();

// 3. 注册自定义Broker服务
services.AddHostedService<MqttBrokerService>();
// 注册接口以便其他类注入
services.AddSingleton<IMqttBrokerService>(sp => sp.GetRequiredService<MqttBrokerService>());

var app = builder.Build();

// 4. 配置WebSocket端点
app.UseRouting();
app.UseEndpoints(endpoints =>
{
    // 映射MQTT WebSocket路径,筛选合法的MQTT子协议
    endpoints.MapConnectionHandler<MqttConnectionHandler>(
        "/ws",
        options =>
        {
            options.WebSockets.SubProtocolSelector = protocols => 
                protocols.FirstOrDefault(p => p.StartsWith("mqtt")) ?? string.Empty;
        });

    endpoints.MapGet("/", async context =>
    {
        await context.Response.WriteAsync("MQTT Broker Running");
    });
});

// 5. 启用MQTT TCP端点(如需保留TCP连接)
app.UseMqttEndpoint();

app.Run();

关键说明:

  • 移除了app.UseMqttServer()的手动配置,AddHostedMqttServerWithServices已完成Server的初始化与托管。
  • WebSocket子协议必须筛选mqtt开头的(如mqtt、mqttv3.1),否则客户端会因协议不匹配断开连接。

问题2:通过其他类触发Broker推送消息

核心是让业务类能访问到MQTT Server实例,调用其InjectApplicationMessageAsync方法发布消息。

实现方式:

  1. 在MqttBrokerService中封装推送方法(推荐,保持逻辑内聚):
public async Task PublishMessageAsync(string topic, string payload, 
    MqttQualityOfServiceLevel qos = MqttQualityOfServiceLevel.AtLeastOnce, bool retain = false)
{
    var message = new MqttApplicationMessageBuilder()
        .WithTopic(topic)
        .WithPayload(payload)
        .WithQualityOfServiceLevel(qos)
        .WithRetainFlag(retain)
        .Build();

    await _mqttServer.InjectApplicationMessageAsync(
        new InjectedMqttApplicationMessage(message)
        {
            SenderClientId = "InternalSystem" // 标记为内部推送消息
        });
}
  1. 在业务类中注入IMqttBrokerService并触发推送:
public class XClass
{
    private readonly IMqttBrokerService _mqttBroker;

    public XClass(IMqttBrokerService mqttBroker)
    {
        _mqttBroker = mqttBroker;
    }

    public async Task TriggerBusinessEvent()
    {
        // 业务逻辑执行完成后,调用推送方法
        await _mqttBroker.PublishMessageAsync(
            "device/alert", 
            "{\"type\": \"temperature\", \"value\": 85}", 
            MqttQualityOfServiceLevel.AtLeastOnce);
    }
}

替代方案:直接注入IMqttServer

若无需封装,可直接在业务类中注入IMqttServer:

public class XClass
{
    private readonly IMqttServer _mqttServer;

    public XClass(IMqttServer mqttServer)
    {
        _mqttServer = mqttServer;
    }

    public async Task TriggerBusinessEvent()
    {
        var message = new MqttApplicationMessageBuilder()
            .WithTopic("device/status")
            .WithPayload("Online")
            .Build();

        await _mqttServer.InjectApplicationMessageAsync(new InjectedMqttApplicationMessage(message));
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 01:25:02