news 2026/8/7 1:16:20

MSMQ实现跨域文件传输——目录监控拆包发送序号合并的完整方案

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
MSMQ实现跨域文件传输——目录监控拆包发送序号合并的完整方案

MSMQ实现跨域文件传输——目录监控+拆包发送+序号合并的完整方案

医院和医保、社保和银行——政务系统之间的文件传输往往要跨网段、跨防火墙。这套方案用C#的FileSystemWatcher监控目录变化,通过微软MSMQ消息队列跨越网络传输文件,支持两种模式:小文件的纯文本流传输和大文件的拆包分块+序号合并。

文章目录

  • MSMQ实现跨域文件传输——目录监控+拆包发送+序号合并的完整方案
    • 一、背景:为什么用MSMQ而不是FTP
    • 二、整体架构:三线程协作
    • 三、两种传输模式
      • 3.1 模式一:纯文本流——适合小文件
      • 3.2 模式二:拆包分块+序号合并——适合大文件
    • 四、两种模式的选择
    • 五、文件占用检测——发送前必须等写完
    • 六、事务支持——发送失败自动回滚
    • 七、为什么用MSMQ——跨网段隔离的实际需求
    • 八、完整流程总结
    • 九、展望——丢失分包的补发机制
    • 十、结语

一、背景:为什么用MSMQ而不是FTP

之前写过社保系统和银行之间用FTP传文件。但FTP有几个痛点:

  • 防火墙穿透难——主动模式下服务器要反向连客户端随机端口
  • 文件完整性无保障——文件传一半断了只能整体重传
  • 无事务保证——接收端收到文件但没写入磁盘,丢数据追不到

微软MSMQ(Microsoft Message Queuing)提供了几个关键能力:

  • TCP直连——FormatName:DIRECT=TCP:192.168.70.3\private$\msmq1,不需要额外端口
  • 事务消息——MessageQueueTransaction,发送失败自动回滚
  • 对象序列化——直接把一个C#对象丢进队列,接收端反序列化出来

这套方案就是基于MSMQ实现的一个文件传输引擎。


二、整体架构:三线程协作

发送端 接收端 ┌─────────────────────┐ ┌─────────────────────┐ 监控线程 │ FileSystemWatcher │ │ receive线程 │ │ 监控 sendDt 目录 │ │ Peek → Receive │ │ 发现新文件→发送 │ │ MSMQ队列 │ └────────┬────────────┘ └──────────┬──────────┘ │ │ ▼ ▼ ┌─────────────────────┐ MSMQ ┌─────────────────────┐ 发送模式 │ SendMessage(拆包) │ ←─ TCP ──→ │ ReceiveMessage(合包) │ │ SendString(文本流) │ 队列 │ ReceiveString(直接写)│ └─────────────────────┘ └─────────────────────┘ │ │ ▼ ▼ ┌─────────────────────┐ ┌─────────────────────┐ │ sendedDT(已发) │ │ objDt(接收完成) │ └─────────────────────┘ │ tempDt(分包暂存) │ └─────────────────────┘
  • 监控线程FileSystemWatcher监控sendDt目录,新文件出现→检查文件是否被其他进程占用→发送
  • 发送线程:从磁盘读文件→拆包或直接流式发送→写入MSMQ队列
  • 接收线程:轮询本地MSMQ→接收消息→写磁盘→收到完整文件后可选转发

三、两种传输模式

3.1 模式一:纯文本流——适合小文件

适用于XML文件、配置文件、报文等文本内容。发送时把文件名和内容拼成一个字符串,接收时按分隔符拆开。

发送端

publicstaticboolSendString(Stringfilename){StreamReadersr=newStreamReader(sendDt+"\\"+filename);MessageQueuemyQueue=newMessageQueue(sendurl);System.Messaging.MessagemyMessage=newSystem.Messaging.Message();// 协议格式:文件名@文件内容stringcontent=sr.ReadToEnd().ToString();sr.Close();Streamstream=newMemoryStream(Encoding.UTF8.GetBytes(content));myMessage.BodyStream=stream;// 事务发送MessageQueueTransactionmyTransaction=newMessageQueueTransaction();myTransaction.Begin();myQueue.Send(myMessage,myTransaction);myTransaction.Commit();myQueue.Close();// 移动到已发送目录File.Move(sendDt+"\\"+filename,sendedDT+"\\"+filename);}

接收端

publicstaticstringReceiveString(){MessageQueuemyQueue=newMessageQueue(localurl);System.Messaging.MessagemyMessage=myQueue.Receive();Streamsr=myMessage.BodyStream;Byte[]bt=newByte[sr.Length];sr.Read(bt,0,(int)sr.Length);Stringstr=Encoding.UTF8.GetString(bt);// 写入到接收目录StreamWritersw=newStreamWriter(objDt+"\\"+filename,false,Encoding.UTF8);sw.WriteLine(str);sw.Close();}

适用场景:文本报文、配置文件、小体量的XML。局限:大文件全量读进内存再发送,signBig建议不超过2MB。

3.2 模式二:拆包分块+序号合并——适合大文件

适用于二进制文件、大体积XML、报表文件。发送端按2MB切块,每块带序号和总块列表。接收端收到每块临时存盘,全部收齐后按序号合并。

核心数据对象——每一块的结构

publicclassInformations{publicbyte[]AttByte{get;set;}// 分块数据publicstringFileName{get;set;}// 原始文件名publicstringfilelist{get;set;}// 总块序列表 "0,1,2,...,N"publicintid{get;set;}// 本块序号publicstringtype{get;set;}// 消息类型: 1=文件, 2=回执, 3=重发}

发送端——拆包

publicstaticboolSendMessage(Stringfilename){FileStreamfsReader=newFileStream(sendDt+"\\"+filename,FileMode.Open,FileAccess.Read);BinaryReaderbReader=newBinaryReader(fsReader);intbig=Convert.ToInt32(signBig);// 每块大小,默认2072576字节≈2MBintloops=(int)fsReader.Length/big;if(loops*big<fsReader.Length)loops++;// 生成块序列表Stringfilelist="";for(inti=0;i<loops;i++){filelist=(i==0)?i.ToString():filelist+","+i;}// 逐块发送for(inti=0;i<loops;i++){Informationstempbook=newInformations();tempbook.filelist=filelist;tempbook.FileName=filename;tempbook.id=i;tempbook.type="1";tempbook.AttByte=bReader.ReadBytes(big);System.Messaging.MessagemyMessage=newSystem.Messaging.Message();myMessage.Body=tempbook;myMessage.Formatter=newXmlMessageFormatter(newType[]{typeof(Informations)});myQueue.Send(myMessage);}fsReader.Close();bReader.Close();// 移动到已发送File.Move(sendDt+"\\"+filename,sendedDT+"\\"+filename);}

关键点:不是把全部文件读进内存再切。bReader.ReadBytes(big)每次只读一块大小——8GB的文件也只占用2MB内存。循环读→发→读→发,内存占用恒定。

接收端——收块+合并

publicstaticstringReceiveMessage(){MessageQueuemyQueue=newMessageQueue(localurl);myQueue.Formatter=newXmlMessageFormatter(newType[]{typeof(Informations)});System.Messaging.MessagemyMessage=myQueue.Receive();Informationsbook=(Informations)myMessage.Body;if("1".Equals(book.type))// 文件消息{// ① 把当前块写入临时文件StringtempName=tempDt+"\\"+book.FileName+"."+book.id+tempExt;// "report.pdf.0.ext"TransByteToFile(book.AttByte,tempName);// ② 检查所有块是否收齐String[]list=book.filelist.Split(',');boolallexists=true;for(inti=0;i<list.Length;i++){Stringname=tempDt+"\\"+book.FileName+"."+list[i]+tempExt;if(!File.Exists(name)){allexists=false;break;}}// ③ 收齐了——按序号合并成完整文件if(allexists){StringcompPath=objDt+"\\"+book.FileName;FileStreamfsWrite=newFileStream(compPath,FileMode.CreateNew,FileAccess.Write);BinaryWriterbWrite=newBinaryWriter(fsWrite);for(inti=0;i<list.Length;i++){Stringname=tempDt+"\\"+book.FileName+"."+list[i]+tempExt;Byte[]attByte=TransFileToByte(name);File.Delete(name);// 合并后删除临时文件bWrite.Write(attByte);}fsWrite.Close();bWrite.Close();}}elseif("2".Equals(book.type))// 回执消息{// 处理回执...}}

三个步骤:收一块→写临时文件→检查filelist是否全到了→全了就按序号从0到N依次合并。

filelist是发送端拼的——"0,1,2,3"。接收端每收到一块,就根据filelist的列表检查对应的文件名.{序号}.ext是否都存在于tempDt。都到了就是收齐,按序号顺序读→拼→写完整文件→删临时文件。


四、两种模式的选择

维度纯文本流(SendString)拆包分块(SendMessage)
适用文本文件、报文、配置文件二进制文件、大文件、报表
内存占用整文件读进内存恒定2MB
断点续传不支持天然支持——丢了某块只需重发那一块
传输协议文件名@内容Informations对象(XML序列化)
传输保证事务无(XmlMessageFormatter 不支持事务)

分开用:公文报文、配置文件走文本流。批量报表、照片打包、二进制数据走拆包分块。


五、文件占用检测——发送前必须等写完

发送端监控目录,但文件不是瞬间写完的——复制一个大文件到sendDt可能需要几秒。如果在文件还在写入时就送出去,收到的就是半截文件。

publicstaticboolIsFileInUse(stringfileName){boolinUse=true;try{// 尝试以独占方式打开——如果别的进程在写,这里抛异常FileStreamfs=newFileStream(fileName,FileMode.Open,FileAccess.Read,FileShare.None);fs.Close();inUse=false;}catch{// 文件被占用,稍后再试}returninUse;// true=正在使用, false=可以发送}

监控事件里用死循环等待:

publicstaticvoidOnChanged(objectsender,FileSystemEventArgse){if(File.Exists(e.FullPath)){while(IsFileInUse(e.FullPath)){// 文件还被占用,等待...}// 文件写入完毕,开始发送bpMessage.SendMessage(e.Name);}}

FileShare.None是关键——要求以独占方式打开。如果其他进程还在FileStream.WriteFileShare.None打开失败,说明文件没写完。


六、事务支持——发送失败自动回滚

文本流模式用了MessageQueueTransaction

MessageQueueTransactionmyTransaction=newMessageQueueTransaction();try{myTransaction.Begin();myQueue.Send(myMessage,myTransaction);myTransaction.Commit();}catch(Exceptionex){myTransaction.Abort();}

事务的内容:要么消息完整写入队列,要么完全不写。接收端看不到一个半截的消息——要么收到一条完整消息,要么一条都没有。文件移动在事务提交之后——确保"确认入队"和"标记已发"是原子操作。


七、为什么用MSMQ——跨网段隔离的实际需求

政务网络通常分为互联网区、政务外网区、管理网区等——各区之间物理隔离,不能直接通过对端IP传文件。FTP搞不定(防火墙拦住主动模式的随机端口),共享目录更不可能(跨不了网段)。

MSMQ自带跨网段消息投递能力——应用层只跟本地MSMQ服务打交道,不跟对端建立TCP连接。FormatName:DIRECT=TCP:192.168.70.3\private$\msmq1是告诉本地MSMQ服务"投递目标",后续的建连接、序列化、超时重发全是MSMQ底层管,应用代码一行都不用写。

<!-- Conf.xml --><sendurl>FormatName:DIRECT=TCP:192.168.70.3\private$\msmq1</sendurl><!-- 发送目标──MSMQ服务负责跨网段投递 --><localurl>FormatName:DIRECT=TCP:192.168.70.3\private$\msmq1</localurl><!-- 本地接收队列──从本机MSMQ取消息 --><reivurl>FormatName:DIRECT=TCP:192.168.203.90\private$\msmq1</reivurl><!-- 回执队列──另一个网段的地址 -->

应用层做的事情就是:监控目录→文件来了→调用myQueue.Send()→完事。对应用来说,就是往一个本地的队列对象里塞了一条消息。至于这条消息怎么穿过互联网区到管理网区——MSMQ自己搞定。


八、完整流程总结

① FileSystemWatcher 监控 D:\msmq2\send 目录 │ ▼ 发现新文件 report.pdf ② IsFileInUse → 等待写入完成 │ ▼ 文件写入完毕 ③ 文件 > 2MB → SendMessage(拆包分块) │ 文件 ≤ 2MB → SendString(文本流) ▼ ④ MSMQ 跨域传输 │ TCP DIRECT:192.168.70.3\private$\msmq1 ▼ ⑤ 接收端轮询 → ReceiveMessage / ReceiveString │ ├── 分块模式 → 写 tempDt → 等收齐 → 按序号合并 → 写入 objDt └── 文本流模式 → 直接写入 objDt │ ▼ ⑥ 可选:转发到下一跳队列 或 发送回执

九、展望——丢失分包的补发机制

当前方案在正常网络环境下工作得很好。但极端情况下——某几块消息在队列传输中丢失、接收端磁盘满了导致写临时文件失败——会出现filelist里的某些序号永远等不到。

Informations类里已经预留了type字段的三个值:1=文件分块,2=回执,3=重发。这意味着设计之初就想到了补发场景,只是没有实现。

合理的补发方案——不阻塞发送端

发送端继续保持无状态——发完所有块就把原文件移到sendedDT,不等确认不阻塞。补充一条独立的反馈通道:

发送端 接收端 │ │ │ 发完所有块→Move到sendedDT │ │ │ │ 收到第1块→启动超时计时器(60秒) │ │ │ 60秒内收齐→合并完成→完成 │ 60秒超时→比对filelist,找出缺失序号 │ │ │ ← type=3(补发请求) ───────────┘ │ {FileName:"report.pdf", │ │ filelist:"2,4", │ │ type:"3"} │ │ │ │ 从sendedDT读原文件 │ │ 按缺失序号重切块发送 │ │ 发完继续等…不Move │ │ │ │ 等待重传块到达→收齐→合并

三点关键设计:

  1. 发送端不等待——发完所有块就Move到sendedDT,不占着监控目录。补发是被动响应,独立线程处理
  2. 补发通道独立——新增一条MSMQ队列专门收type=3请求,新开一个retryThread轮询。与主发送链路完全不耦合
  3. 原文件在sendedDT——只是移动了目录,文件还在。补发时按缺失序号列表重新切块发送,不用额外缓存

思路跟BaoPanTimer银行报文重发是一个道理——发送端只管正常流程,补发是被动的、独立的、异步的。type字段和filelist的设计已经为这个机制打好了基础,只差超时计时器和一条反向队列通道。需要一套完整的确认与重传机制。


十、结语

这套MSMQ文件传输引擎解决了一个政务系统里反复出现的问题——两个不在同一个网络域的服务器之间,怎么可靠地传文件

拆包分块模式解决了大文件的内存问题——不管文件多大,内存只占2MB。序号合并解决了乱序到达的问题——每块有自己的序号,收齐后按序拼接。文件占用检测解决了"文件还没写完就发送"的竞态问题。

这套代码从2014年跑到系统下线,每天传几百个文件,从来没丢过数据。不是设计得多精妙——是每个可能出问题的环节都做了防御性处理。

另外,目录监控+拆包分块+序号合并这套机制不依赖MSMQ——换成RabbitMQ、Kafka或Redis的List队列同样适用。当时选MSMQ只是因为Windows政务环境下它最省事:系统自带,不用额外装中间件。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/7 1:16:12

RT-Thread NetUtils网络组件详解:从核心原理到物联网数据采集实战

1. 从零开始认识 NetUtils&#xff1a;它到底是什么&#xff0c;能帮你解决什么&#xff1f;如果你在嵌入式开发&#xff0c;特别是基于 RT-Thread 的物联网项目中&#xff0c;经常需要处理网络连接、数据收发、协议解析这些“脏活累活”&#xff0c;那你大概率听说过或者已经用…

作者头像 李华
网站建设 2026/8/7 1:14:59

AutoDock Vina完整指南:如何用开源工具加速药物研发

AutoDock Vina完整指南&#xff1a;如何用开源工具加速药物研发 【免费下载链接】AutoDock-Vina AutoDock Vina 项目地址: https://gitcode.com/gh_mirrors/au/AutoDock-Vina 想要快速完成分子对接计算却苦于工具复杂、速度缓慢&#xff1f;AutoDock Vina正是你需要的免…

作者头像 李华
网站建设 2026/8/7 1:14:34

团结引擎微信小游戏打包:WebGL模板配置避坑指南

1. 项目概述&#xff1a;为什么WebGL模板配置是微信小游戏打包的“命门”&#xff1f;如果你正在用团结引擎&#xff08;Unity China&#xff09;开发微信小游戏&#xff0c;并且已经走到了打包这一步&#xff0c;那么恭喜你&#xff0c;也提醒你&#xff1a;最关键的“暗礁”可…

作者头像 李华
网站建设 2026/8/7 1:04:04

Adobe-GenP 3.0终极指南:三步完成Adobe软件激活的完整教程

Adobe-GenP 3.0终极指南&#xff1a;三步完成Adobe软件激活的完整教程 【免费下载链接】Adobe-GenP Adobe CC 2019/2020/2021/2022/2023 GenP Universal Patch 3.0 项目地址: https://gitcode.com/gh_mirrors/ad/Adobe-GenP 如果你正在寻找一款简单高效的Adobe激活工具来…

作者头像 李华
网站建设 2026/8/7 1:00:40

爆仓到稳定翻倍 同花顺期货通指标

今天给大家带来是一款同花顺期货通指标&#xff0c;并且已经上架到同花顺期货通的指标广场上了。喜欢的朋友可以去指标广场安装试用&#xff01;&#xff01;友情提示&#xff1a;&#xff08;指标只是辅助&#xff0c;不作建议&#xff09;拼多多店铺&#xff1a;指标公式编写…

作者头像 李华
网站建设 2026/8/7 1:00:38

RSI变异顶底 同花顺期货通指标

今天给大家带来是一款同花顺期货通指标&#xff0c;并且已经上架到同花顺期货通的指标广场上了。喜欢的朋友可以去指标广场安装试用&#xff01;&#xff01;友情提示&#xff1a;&#xff08;指标只是辅助&#xff0c;不作建议&#xff09;拼多多店铺&#xff1a;指标公式编写…

作者头像 李华