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

如何为MQTTnet Broker实现冗余以构建高可用双节点集群

双节点MQTTnet Broker消息复制与高可用实现方案

核心实现思路

要完成双节点Broker的消息复制,核心是让两个节点互相同步收到的消息。可以通过让每个节点作为MQTT客户端连接到对方节点,订阅全量主题,再将收到的同步消息转发到本地Broker,同时通过消息标识过滤避免循环转发。

代码修改方案(以你的Server2为例,Server1仅需调整端口和节点标识)

1. 添加同步客户端相关字段

在Server类中新增用于跨节点同步的客户端实例和本地节点标识:

// 跨节点同步用的MQTT客户端
private static IMqttClient _syncClient;
// 本地节点唯一标识,用于过滤循环消息
private static readonly string _localNodeId = "Server2"; // Server1需设为"Server1"

2. 初始化同步客户端并建立跨节点连接

在本地Broker启动后,创建同步客户端并连接到另一个节点,订阅所有主题:

// 初始化同步客户端配置,连接到Server1的端口(Server1则连接1884)
var syncClientOptions = new MqttClientOptionsBuilder()
    .WithTcpServer("localhost", 1883)
    .WithClientId($"SyncClient_{_localNodeId}")
    .Build();

_syncClient = new MqttFactory().CreateMqttClient();

// 同步客户端收到消息时,转发到本地Broker
_syncClient.ApplicationMessageReceivedAsync += async e =>
{
    // 跳过本地节点发送的消息,避免循环转发
    if (e.ApplicationMessage.UserProperties.Any(p => p.Name == "NodeId" && p.Value == _localNodeId))
    {
        return;
    }

    // 构建带本地标识的转发消息
    var forwardedMessage = new MqttApplicationMessageBuilder()
        .WithTopic(e.ApplicationMessage.Topic)
        .WithPayload(e.ApplicationMessage.Payload)
        .WithQualityOfServiceLevel(e.ApplicationMessage.QualityOfServiceLevel)
        .WithRetainFlag(e.ApplicationMessage.Retain)
        .WithUserProperty("NodeId", _localNodeId)
        .Build();

    // 注入到本地Broker
    await server.InjectApplicationMessageAsync(new InjectApplicationMessageEventArgs
    {
        ApplicationMessage = forwardedMessage,
        ClientId = $"SyncForward_{_localNodeId}"
    });

    WriteLine($"[{_localNodeId}] 同步消息: Topic={e.ApplicationMessage.Topic}, Payload={Encoding.UTF8.GetString(e.ApplicationMessage.Payload)}");
};

// 连接到对方节点并订阅全量主题
await _syncClient.ConnectAsync(syncClientOptions);
await _syncClient.SubscribeAsync(new MqttTopicFilterBuilder()
    .WithTopic("#")
    .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.ExactlyOnce)
    .Build());

3. 修改消息拦截器添加本地节点标识

在本地消息处理逻辑中,给每条消息添加节点标识,方便同步时过滤:

Task ServerInterceptingPublishAsync(InterceptingPublishEventArgs arg)
{
    var payloadSegment = arg.ApplicationMessage.PayloadSegment;
    var payload = Encoding.UTF8.GetString(payloadSegment.Array!, payloadSegment.Offset, payloadSegment.Count);

    WriteLine(
        " TimeStamp: {0} -- Message: ClientId = {1}, Topic = {2}, Payload = {3}, QoS = {4}, Retain-Flag = {5}",
        DateTime.Now,
        arg.ClientId,
        arg.ApplicationMessage?.Topic,
        payload,
        arg.ApplicationMessage?.QualityOfServiceLevel,
        arg.ApplicationMessage?.Retain);

    // 添加本地节点标识,用于同步时过滤
    arg.ApplicationMessage.UserProperties.Add(new MqttUserProperty("NodeId", _localNodeId));

    return Task.CompletedTask;
}

4. 优化保留消息的同步逻辑

确保本地加载的保留消息也带有节点标识,避免同步时重复处理:

private static async Task ServerOnLoadingRetainedMessageAsync(LoadingRetainedMessagesEventArgs arg)
{
    // 给加载的保留消息添加本地节点标识
    var modifiedMessages = arg.LoadedRetainedMessages.Select(msg => 
    {
        var newMsg = MqttApplicationMessageBuilder.FromApplicationMessage(msg.ApplicationMessage)
            .WithUserProperty("NodeId", _localNodeId)
            .Build();
        return new MqttRetainedMessage(msg.ClientId, newMsg);
    });

    var models = modifiedMessages.Select(MqttRetainedMessageModel.Create);
    var buffer = JsonSerializer.SerializeToUtf8Bytes(models);
    await File.WriteAllBytesAsync(storePath, buffer);
    WriteLine("Retained messages saved with node identifier.");
}

关键注意事项

  • 端口区分:两个节点必须使用不同端口(如Server1用1883,Server2用1884),同步客户端互相连接对方端口。
  • QoS选择:同步消息建议用ExactlyOnce(QoS2),确保消息不丢失、不重复。
  • 资源清理:程序退出时需断开同步客户端连接,停止本地Broker。
  • 节点存活检测:可额外添加心跳机制,当对方节点故障时触发客户端重连或切换逻辑(基础版本可先忽略)。

完整修改后的Server2代码

using System.Text;
using System.Text.Json;
using MQTTnet;
using MQTTnet.Client;
using MQTTnet.Server;
using static System.Console;

namespace Server2;

internal class Server2
{
    static string storePath = "Server2/RetainedMessages.json";
    private static IMqttClient _syncClient;
    private static readonly string _localNodeId = "Server2";

    private static async Task Main(string[] args)
    {
        // 配置本地MQTT Broker
        var options = new MqttServerOptionsBuilder()
            .WithDefaultEndpoint()
            .WithDefaultEndpointPort(1884)
            .WithPersistentSessions();

        var server = new MqttFactory().CreateMqttServer(options.Build());

        server.LoadingRetainedMessageAsync += ServerOnLoadingRetainedMessageAsync;

        // 启动本地Broker
        await server.StartAsync();

        // 初始化跨节点同步客户端
        var syncClientOptions = new MqttClientOptionsBuilder()
            .WithTcpServer("localhost", 1883)
            .WithClientId($"SyncClient_{_localNodeId}")
            .Build();

        _syncClient = new MqttFactory().CreateMqttClient();

        // 同步消息转发逻辑
        _syncClient.ApplicationMessageReceivedAsync += async e =>
        {
            if (e.ApplicationMessage.UserProperties.Any(p => p.Name == "NodeId" && p.Value == _localNodeId))
            {
                return;
            }

            var forwardedMessage = new MqttApplicationMessageBuilder()
                .WithTopic(e.ApplicationMessage.Topic)
                .WithPayload(e.ApplicationMessage.Payload)
                .WithQualityOfServiceLevel(e.ApplicationMessage.QualityOfServiceLevel)
                .WithRetainFlag(e.ApplicationMessage.Retain)
                .WithUserProperty("NodeId", _localNodeId)
                .Build();

            await server.InjectApplicationMessageAsync(new InjectApplicationMessageEventArgs
            {
                ApplicationMessage = forwardedMessage,
                ClientId = $"SyncForward_{_localNodeId}"
            });

            WriteLine($"[{_localNodeId}] 同步消息: Topic={e.ApplicationMessage.Topic}, Payload={Encoding.UTF8.GetString(e.ApplicationMessage.Payload)}");
        };

        // 连接到Server1并订阅全量主题
        await _syncClient.ConnectAsync(syncClientOptions);
        await _syncClient.SubscribeAsync(new MqttTopicFilterBuilder()
            .WithTopic("#")
            .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.ExactlyOnce)
            .Build());

        // 注册连接和消息处理事件
        server.ClientConnectedAsync += ClientConnectedEventArgs;
        server.InterceptingPublishAsync += ServerInterceptingPublishAsync;

        WriteLine("Press any key to stop the server...");
        ReadLine();

        // 清理资源
        await _syncClient.DisconnectAsync();
        await server.StopAsync();
    }

    private static Task ServerInterceptingPublishAsync(InterceptingPublishEventArgs arg)
    {
        var payloadSegment = arg.ApplicationMessage.PayloadSegment;
        var payload = Encoding.UTF8.GetString(payloadSegment.Array!, payloadSegment.Offset, payloadSegment.Count);

        WriteLine(
            " TimeStamp: {0} -- Message: ClientId = {1}, Topic = {2}, Payload = {3}, QoS = {4}, Retain-Flag = {5}",
            DateTime.Now,
            arg.ClientId,
            arg.ApplicationMessage?.Topic,
            payload,
            arg.ApplicationMessage?.QualityOfServiceLevel,
            arg.ApplicationMessage?.Retain);

        arg.ApplicationMessage.UserProperties.Add(new MqttUserProperty("NodeId", _localNodeId));

        return Task.CompletedTask;
    }

    private static Task ClientConnectedEventArgs(ClientConnectedEventArgs arg)
    {
        WriteLine("New connection: ClientId = {0}, Endpoint = {1}", arg.ClientId, arg.Endpoint);
        return Task.CompletedTask;
    }

    private static async Task ServerOnLoadingRetainedMessageAsync(LoadingRetainedMessagesEventArgs arg)
    {
        var modifiedMessages = arg.LoadedRetainedMessages.Select(msg =>
        {
            var newMsg = MqttApplicationMessageBuilder.FromApplicationMessage(msg.ApplicationMessage)
                .WithUserProperty("NodeId", _localNodeId)
                .Build();
            return new MqttRetainedMessage(msg.ClientId, newMsg);
        });

        var models = modifiedMessages.Select(MqttRetainedMessageModel.Create);
        var buffer = JsonSerializer.SerializeToUtf8Bytes(models);
        await File.WriteAllBytesAsync(storePath, buffer);
        WriteLine("Retained messages saved with node identifier.");
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 17:14:57