如何为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
相关产品推荐
相关产品推荐

