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

如何在.NET Framework中通过MQTT实现Broker向所有客户端发送消息?

问题描述

我正在开发一个应用,需要实现Broker向所有客户端发送消息,同时客户端也能向Broker发送消息。目前客户端向Broker发送消息的功能正常,但我无法实现Broker向所有客户端发送消息的功能。我知道需要用订阅机制,也试过订阅主题,但没效果,也试过MqttFactory,还是没解决。

Broker代码
private static int MessageCounter = 0;

static void Main(string[] args)
{
    Log.Logger = new LoggerConfiguration()
        .MinimumLevel.Debug()
        .Enrich.FromLogContext()
        .WriteTo.Console()
        .CreateLogger();

    MqttServerOptionsBuilder options = new MqttServerOptionsBuilder()
        .WithDefaultEndpoint()
        .WithDefaultEndpointPort(707)
        .WithConnectionValidator(OnNewConnection)
        .WithApplicationMessageInterceptor(OnNewMessage);


    IMqttServer mqttServer = new MqttFactory().CreateMqttServer();

    mqttServer.StartAsync(options.Build()).GetAwaiter().GetResult();

    Console.ReadLine();
}

public static void OnNewConnection(MqttConnectionValidatorContext context)
{
    Log.Logger.Information(
            "New connection: ClientId = {clientId}, Endpoint = {endpoint}, CleanSession = {cleanSession}",
            context.ClientId,
            context.Endpoint,
            context.CleanSession);
}

public static void OnNewMessage(MqttApplicationMessageInterceptorContext context)
{
    var payload = context.ApplicationMessage?.Payload == null ? null : Encoding.UTF8.GetString(context.ApplicationMessage?.Payload);

  
    Log.Logger.Information(
        "MessageId: {MessageCounter} - TimeStamp: {TimeStamp} -- Message: ClientId = {clientId}, Topic = {topic}, Payload = {payload}, QoS = {qos}, Retain-Flag = {retainFlag}",
        MessageCounter,
        DateTime.Now,
        context.ClientId,
        context.ApplicationMessage?.Topic,
        payload,
        context.ApplicationMessage?.QualityOfServiceLevel,
        context.ApplicationMessage?.Retain);
}
客户端代码
static void Main(string[] args)
{
    Log.Logger = new LoggerConfiguration()
        .MinimumLevel.Debug()
        .Enrich.FromLogContext()
        .WriteTo.Console()
        .CreateLogger();

    MqttClientOptionsBuilder builder = new MqttClientOptionsBuilder()
                                .WithClientId("Dev.To")
                                .WithTcpServer("localhost", 707);

    ManagedMqttClientOptions options = new ManagedMqttClientOptionsBuilder()
                                .WithAutoReconnectDelay(TimeSpan.FromSeconds(60))
                                .WithClientOptions(builder.Build())
                                .Build();

    IManagedMqttClient _mqttClient = new MqttFactory().CreateManagedMqttClient();

    _mqttClient.ConnectedHandler = new MqttClientConnectedHandlerDelegate(OnConnected);
    _mqttClient.DisconnectedHandler = new MqttClientDisconnectedHandlerDelegate(OnDisconnected);
    _mqttClient.ConnectingFailedHandler = new ConnectingFailedHandlerDelegate(OnConnectingFailed);

    _mqttClient.ApplicationMessageReceivedHandler = new MqttApplicationMessageReceivedHandlerDelegate(a =>
    {
        Log.Logger.Information("Message recieved: {payload}", a.ApplicationMessage);
    });

    _mqttClient.StartAsync(options).GetAwaiter().GetResult();

    //Try this maybeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee
    MqttClientSubscribeOptions _subscribeOptions = new MqttClientSubscribeOptions();
    
    while (true)
    {
        string json = JsonConvert.SerializeObject(new { message = Console.ReadLine(), sent= DateTimeOffset.UtcNow });
        _mqttClient.PublishAsync("topic", json);
        Task.Delay(100).GetAwaiter().GetResult();
    }
}

public static void OnConnected(MqttClientConnectedEventArgs obj)
{
    Log.Logger.Information("Successfully connected.");
}

public static void OnConnectingFailed(ManagedProcessFailedEventArgs obj)
{
    Log.Logger.Warning("Couldn't connect to broker.");
}

public static void OnDisconnected(MqttClientDisconnectedEventArgs obj)
{
    Log.Logger.Information("Successfully disconnected.");
}
解决方案

1. 客户端补全订阅逻辑

你的客户端代码仅创建了订阅选项但未执行订阅操作,这是核心问题。需要在客户端连接成功后,主动订阅目标主题:

首先将_mqttClient改为类级静态变量,确保OnConnected方法能访问它:

private static IManagedMqttClient _mqttClient;

static void Main(string[] args)
{
    // ... 原有代码 ...
    _mqttClient = new MqttFactory().CreateManagedMqttClient();
    // ... 原有代码 ...
}

然后修改OnConnected方法添加订阅:

public static void OnConnected(MqttClientConnectedEventArgs obj)
{
    Log.Logger.Information("Successfully connected.");
    // 订阅广播主题,这里和Broker发送的主题保持一致
    _mqttClient.SubscribeAsync("broadcast/topic", MqttQualityOfServiceLevel.AtLeastOnce).GetAwaiter().GetResult();
}

2. Broker添加主动广播逻辑

当前Broker仅接收消息,没有主动推送的逻辑。可以在Broker启动后添加控制台输入广播的功能:

修改Broker的Main方法:

static void Main(string[] args)
{
    // ... 原有代码 ...
    IMqttServer mqttServer = new MqttFactory().CreateMqttServer();
    mqttServer.StartAsync(options.Build()).GetAwaiter().GetResult();

    // 新增广播逻辑
    Task.Run(async () =>
    {
        while (true)
        {
            Console.WriteLine("输入要广播的消息:");
            string content = Console.ReadLine();
            if (!string.IsNullOrWhiteSpace(content))
            {
                var broadcastMsg = new MqttApplicationMessageBuilder()
                    .WithTopic("broadcast/topic")
                    .WithPayload(Encoding.UTF8.GetBytes(content))
                    .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce)
                    .WithRetainFlag(false)
                    .Build();

                await mqttServer.PublishAsync(broadcastMsg);
                Log.Logger.Information("已向所有客户端广播消息:{content}", content);
            }
        }
    });

    Console.ReadLine();
}

3. 关键注意事项

  • 客户端订阅的主题必须和Broker发送的主题完全一致,否则无法收到消息
  • 订阅操作必须在客户端连接成功后执行,避免因连接未建立导致订阅失败
  • 若使用ManagedMqttClient,在Connected事件中执行订阅是最安全的时机

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 16:39:21