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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 12:48:09