如何使用C# MQTTnet实现文件的发布与订阅
实现思路
MQTT协议单条消息存在大小限制(多数Broker默认阈值为256KB),文件传输根据体积分两种处理逻辑,你使用的分隔符格式结构化文本文件体积普遍不大,两种场景都可以覆盖:
- 发送端:读取文件原始字节流,文件大小超过单消息阈值时按固定块大小分片,每个分片携带文件名、分片序号、总分片数标识,逐片发送到指定主题;小于阈值直接整文件发送
- 接收端:按文件名缓存收到的分片,收齐所有分片后按序号拼接字节流,写入本地磁盘即可
- 纯二进制字节传输不会和文件内的分隔符、MQTT协议帧产生冲突,比直接转字符串传输可靠性更高,也不需要额外做转义处理
实现示例
以下示例基于MQTTnet 4.x版本编写,使用前先通过NuGet安装MQTTnet包即可。
发送端代码
using MQTTnet; using MQTTnet.Client; using System.IO; // 初始化MQTT客户端 var mqttFactory = new MqttFactory(); using var mqttClient = mqttFactory.CreateMqttClient(); var connectOptions = new MqttClientOptionsBuilder() .WithTcpServer("你的MQTT Broker地址", 1883) // 替换为实际部署的Broker地址和端口 .WithClientId("FileSender_" + Guid.NewGuid().ToString("N")) .Build(); await mqttClient.ConnectAsync(connectOptions); /// <summary> /// 发送文件到指定主题 /// </summary> /// <param name="filePath">本地待发送文件路径</param> /// <param name="targetTopic">传输目标主题</param> /// <param name="chunkSize">单分片大小,默认128KB</param> async Task SendFile(string filePath, string targetTopic, int chunkSize = 128 * 1024) { if (!File.Exists(filePath)) throw new FileNotFoundException("待发送文件不存在", filePath); var fileName = Path.GetFileName(filePath); var fileBytes = await File.ReadAllBytesAsync(filePath); var totalChunks = (int)Math.Ceiling(fileBytes.Length / (double)chunkSize); // 小文件直接整包发送 if (totalChunks == 1) { var message = new MqttApplicationMessageBuilder() .WithTopic(targetTopic) .WithPayload(fileBytes) .WithUserProperty("FileName", fileName) .WithUserProperty("TotalChunks", "1") .WithUserProperty("ChunkIndex", "0") .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .Build(); await mqttClient.PublishAsync(message); return; } // 大文件分片发送 for (int i = 0; i < totalChunks; i++) { var chunk = fileBytes.Skip(i * chunkSize).Take(chunkSize).ToArray(); var message = new MqttApplicationMessageBuilder() .WithTopic(targetTopic) .WithPayload(chunk) .WithUserProperty("FileName", fileName) .WithUserProperty("TotalChunks", totalChunks.ToString()) .WithUserProperty("ChunkIndex", i.ToString()) .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .Build(); await mqttClient.PublishAsync(message); await Task.Delay(10); // 短延迟避免触发Broker限流 } } // 调用示例:发送分隔符格式的发电数据文件 await SendFile(@"C:\本地文件路径\power_gen_data.txt", "mqtt/file/transfer");
接收端代码
using MQTTnet; using MQTTnet.Client; using System.Collections.Concurrent; using System.IO; using System.Text; // 分片缓存:Key为文件名,Value为(总分片数, 已收到分片集合<分片序号, 分片字节>) private static ConcurrentDictionary<string, (int TotalChunks, ConcurrentDictionary<int, byte[]> Chunks)> _fileCache = new(); // 初始化MQTT客户端 var mqttFactory = new MqttFactory(); using var mqttClient = mqttFactory.CreateMqttClient(); var connectOptions = new MqttClientOptionsBuilder() .WithTcpServer("你的MQTT Broker地址", 1883) .WithClientId("FileReceiver_" + Guid.NewGuid().ToString("N")) .Build(); // 注册消息接收处理逻辑 mqttClient.ApplicationMessageReceivedAsync += async e => { if (e.ApplicationMessage.Topic != "mqtt/file/transfer") return; // 读取文件元数据 var fileName = e.ApplicationMessage.UserProperties.FirstOrDefault(p => p.Name == "FileName")?.Value?.ToString(); var totalChunks = int.Parse(e.ApplicationMessage.UserProperties.First(p => p.Name == "TotalChunks").Value.ToString()!); var chunkIndex = int.Parse(e.ApplicationMessage.UserProperties.First(p => p.Name == "ChunkIndex").Value.ToString()!); var chunkData = e.ApplicationMessage.Payload.ToArray(); if (string.IsNullOrEmpty(fileName)) return; // 分片写入缓存 var fileEntry = _fileCache.GetOrAdd(fileName, _ => (totalChunks, new ConcurrentDictionary<int, byte[]>())); fileEntry.Chunks.TryAdd(chunkIndex, chunkData); // 收齐所有分片后拼接保存 if (fileEntry.Chunks.Count == totalChunks) { var completeFileBytes = fileEntry.Chunks.OrderBy(kv => kv.Key) .SelectMany(kv => kv.Value) .ToArray(); var savePath = Path.Combine(@"C:\接收文件保存目录", fileName); await File.WriteAllBytesAsync(savePath, completeFileBytes); _fileCache.TryRemove(fileName, out _); Console.WriteLine($"文件 {fileName} 接收完成,保存路径:{savePath}"); // 此处可直接读取文件内容解析分隔符格式的发电数据 // var fileContent = Encoding.UTF8.GetString(completeFileBytes); } }; await mqttClient.ConnectAsync(connectOptions); await mqttClient.SubscribeAsync("mqtt/file/transfer", MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce); Console.WriteLine("文件接收端已启动,等待传输..."); Console.ReadLine();
注意事项
- 分片大小需要根据实际部署的MQTT Broker单消息最大阈值调整,不要超过
max_message_size配置,否则消息会被Broker直接丢弃 - 公网传输时建议启用MQTT TLS加密,避免文件内容泄露
- 你提供的空格/制表符分隔的电力数据属于纯文本格式,上述方案会原样传输文件字节,不会破坏原有分隔符结构,接收完成后可以直接按原有逻辑解析
- 如果需要更高传输可靠性,可以将消息QoS级别设置为ExactlyOnce(级别2),减少分片丢包概率
- 多客户端同时传输同名文件时,可以将发送端ClientID加入消息元数据,缓存时用「发送端ID+文件名」作为Key,避免分片混淆
内容的提问来源于stack exchange,提问作者chelsea
相关产品推荐
相关产品推荐

