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.Write,FileShare.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 │ │ │ │ 等待重传块到达→收齐→合并三点关键设计:
- 发送端不等待——发完所有块就Move到
sendedDT,不占着监控目录。补发是被动响应,独立线程处理 - 补发通道独立——新增一条MSMQ队列专门收
type=3请求,新开一个retryThread轮询。与主发送链路完全不耦合 - 原文件在
sendedDT里——只是移动了目录,文件还在。补发时按缺失序号列表重新切块发送,不用额外缓存
思路跟BaoPanTimer银行报文重发是一个道理——发送端只管正常流程,补发是被动的、独立的、异步的。type字段和filelist的设计已经为这个机制打好了基础,只差超时计时器和一条反向队列通道。需要一套完整的确认与重传机制。
十、结语
这套MSMQ文件传输引擎解决了一个政务系统里反复出现的问题——两个不在同一个网络域的服务器之间,怎么可靠地传文件。
拆包分块模式解决了大文件的内存问题——不管文件多大,内存只占2MB。序号合并解决了乱序到达的问题——每块有自己的序号,收齐后按序拼接。文件占用检测解决了"文件还没写完就发送"的竞态问题。
这套代码从2014年跑到系统下线,每天传几百个文件,从来没丢过数据。不是设计得多精妙——是每个可能出问题的环节都做了防御性处理。
另外,目录监控+拆包分块+序号合并这套机制不依赖MSMQ——换成RabbitMQ、Kafka或Redis的List队列同样适用。当时选MSMQ只是因为Windows政务环境下它最省事:系统自带,不用额外装中间件。