拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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

MSMQ实现跨域文件传输——目录监控拆包发送序号合并的完整方案医院和医保、社保和银行——政务系统之间的文件传输往往要跨网段、跨防火墙。这套方案用C#的FileSystemWatcher监控目录变化通过微软MSMQ消息队列跨越网络传输文件支持两种模式小文件的纯文本流传输和大文件的拆包分块序号合并。文章目录MSMQ实现跨域文件传输——目录监控拆包发送序号合并的完整方案一、背景为什么用MSMQ而不是FTP二、整体架构三线程协作三、两种传输模式3.1 模式一纯文本流——适合小文件3.2 模式二拆包分块序号合并——适合大文件四、两种模式的选择五、文件占用检测——发送前必须等写完六、事务支持——发送失败自动回滚七、为什么用MSMQ——跨网段隔离的实际需求八、完整流程总结九、展望——丢失分包的补发机制十、结语一、背景为什么用MSMQ而不是FTP之前写过社保系统和银行之间用FTP传文件。但FTP有几个痛点防火墙穿透难——主动模式下服务器要反向连客户端随机端口文件完整性无保障——文件传一半断了只能整体重传无事务保证——接收端收到文件但没写入磁盘丢数据追不到微软MSMQMicrosoft Message Queuing提供了几个关键能力TCP直连——FormatName:DIRECTTCP: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){StreamReadersrnewStreamReader(sendDt\\filename);MessageQueuemyQueuenewMessageQueue(sendurl);System.Messaging.MessagemyMessagenewSystem.Messaging.Message();// 协议格式文件名文件内容stringcontentsr.ReadToEnd().ToString();sr.Close();StreamstreamnewMemoryStream(Encoding.UTF8.GetBytes(content));myMessage.BodyStreamstream;// 事务发送MessageQueueTransactionmyTransactionnewMessageQueueTransaction();myTransaction.Begin();myQueue.Send(myMessage,myTransaction);myTransaction.Commit();myQueue.Close();// 移动到已发送目录File.Move(sendDt\\filename,sendedDT\\filename);}接收端publicstaticstringReceiveString(){MessageQueuemyQueuenewMessageQueue(localurl);System.Messaging.MessagemyMessagemyQueue.Receive();StreamsrmyMessage.BodyStream;Byte[]btnewByte[sr.Length];sr.Read(bt,0,(int)sr.Length);StringstrEncoding.UTF8.GetString(bt);// 写入到接收目录StreamWriterswnewStreamWriter(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,...,Npublicintid{get;set;}// 本块序号publicstringtype{get;set;}// 消息类型: 1文件, 2回执, 3重发}发送端——拆包publicstaticboolSendMessage(Stringfilename){FileStreamfsReadernewFileStream(sendDt\\filename,FileMode.Open,FileAccess.Read);BinaryReaderbReadernewBinaryReader(fsReader);intbigConvert.ToInt32(signBig);// 每块大小默认2072576字节≈2MBintloops(int)fsReader.Length/big;if(loops*bigfsReader.Length)loops;// 生成块序列表Stringfilelist;for(inti0;iloops;i){filelist(i0)?i.ToString():filelist,i;}// 逐块发送for(inti0;iloops;i){InformationstempbooknewInformations();tempbook.filelistfilelist;tempbook.FileNamefilename;tempbook.idi;tempbook.type1;tempbook.AttBytebReader.ReadBytes(big);System.Messaging.MessagemyMessagenewSystem.Messaging.Message();myMessage.Bodytempbook;myMessage.FormatternewXmlMessageFormatter(newType[]{typeof(Informations)});myQueue.Send(myMessage);}fsReader.Close();bReader.Close();// 移动到已发送File.Move(sendDt\\filename,sendedDT\\filename);}关键点不是把全部文件读进内存再切。bReader.ReadBytes(big)每次只读一块大小——8GB的文件也只占用2MB内存。循环读→发→读→发内存占用恒定。接收端——收块合并publicstaticstringReceiveMessage(){MessageQueuemyQueuenewMessageQueue(localurl);myQueue.FormatternewXmlMessageFormatter(newType[]{typeof(Informations)});System.Messaging.MessagemyMessagemyQueue.Receive();Informationsbook(Informations)myMessage.Body;if(1.Equals(book.type))// 文件消息{// ① 把当前块写入临时文件StringtempNametempDt\\book.FileName.book.idtempExt;// report.pdf.0.extTransByteToFile(book.AttByte,tempName);// ② 检查所有块是否收齐String[]listbook.filelist.Split(,);boolallexiststrue;for(inti0;ilist.Length;i){StringnametempDt\\book.FileName.list[i]tempExt;if(!File.Exists(name)){allexistsfalse;break;}}// ③ 收齐了——按序号合并成完整文件if(allexists){StringcompPathobjDt\\book.FileName;FileStreamfsWritenewFileStream(compPath,FileMode.CreateNew,FileAccess.Write);BinaryWriterbWritenewBinaryWriter(fsWrite);for(inti0;ilist.Length;i){StringnametempDt\\book.FileName.list[i]tempExt;Byte[]attByteTransFileToByte(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){boolinUsetrue;try{// 尝试以独占方式打开——如果别的进程在写这里抛异常FileStreamfsnewFileStream(fileName,FileMode.Open,FileAccess.Read,FileShare.None);fs.Close();inUsefalse;}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打开失败说明文件没写完。六、事务支持——发送失败自动回滚文本流模式用了MessageQueueTransactionMessageQueueTransactionmyTransactionnewMessageQueueTransaction();try{myTransaction.Begin();myQueue.Send(myMessage,myTransaction);myTransaction.Commit();}catch(Exceptionex){myTransaction.Abort();}事务的内容要么消息完整写入队列要么完全不写。接收端看不到一个半截的消息——要么收到一条完整消息要么一条都没有。文件移动在事务提交之后——确保确认入队和标记已发是原子操作。七、为什么用MSMQ——跨网段隔离的实际需求政务网络通常分为互联网区、政务外网区、管理网区等——各区之间物理隔离不能直接通过对端IP传文件。FTP搞不定防火墙拦住主动模式的随机端口共享目录更不可能跨不了网段。MSMQ自带跨网段消息投递能力——应用层只跟本地MSMQ服务打交道不跟对端建立TCP连接。FormatName:DIRECTTCP:192.168.70.3\private$\msmq1是告诉本地MSMQ服务投递目标后续的建连接、序列化、超时重发全是MSMQ底层管应用代码一行都不用写。!-- Conf.xml --sendurlFormatName:DIRECTTCP:192.168.70.3\private$\msmq1/sendurl!-- 发送目标──MSMQ服务负责跨网段投递 --localurlFormatName:DIRECTTCP:192.168.70.3\private$\msmq1/localurl!-- 本地接收队列──从本机MSMQ取消息 --reivurlFormatName:DIRECTTCP: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,找出缺失序号 │ │ │ ← type3(补发请求) ───────────┘ │ {FileName:report.pdf, │ │ filelist:2,4, │ │ type:3} │ │ │ │ 从sendedDT读原文件 │ │ 按缺失序号重切块发送 │ │ 发完继续等…不Move │ │ │ │ 等待重传块到达→收齐→合并三点关键设计发送端不等待——发完所有块就Move到sendedDT不占着监控目录。补发是被动响应独立线程处理补发通道独立——新增一条MSMQ队列专门收type3请求新开一个retryThread轮询。与主发送链路完全不耦合原文件在sendedDT里——只是移动了目录文件还在。补发时按缺失序号列表重新切块发送不用额外缓存思路跟BaoPanTimer银行报文重发是一个道理——发送端只管正常流程补发是被动的、独立的、异步的。type字段和filelist的设计已经为这个机制打好了基础只差超时计时器和一条反向队列通道。需要一套完整的确认与重传机制。十、结语这套MSMQ文件传输引擎解决了一个政务系统里反复出现的问题——两个不在同一个网络域的服务器之间怎么可靠地传文件。拆包分块模式解决了大文件的内存问题——不管文件多大内存只占2MB。序号合并解决了乱序到达的问题——每块有自己的序号收齐后按序拼接。文件占用检测解决了文件还没写完就发送的竞态问题。这套代码从2014年跑到系统下线每天传几百个文件从来没丢过数据。不是设计得多精妙——是每个可能出问题的环节都做了防御性处理。另外目录监控拆包分块序号合并这套机制不依赖MSMQ——换成RabbitMQ、Kafka或Redis的List队列同样适用。当时选MSMQ只是因为Windows政务环境下它最省事系统自带不用额外装中间件。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门