ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

ActiveMQ C# 接入实战:从消息丢失到可靠投递的完整 Demo

ActiveMQ C# 接入实战:从消息丢失到可靠投递的完整 Demo 简介这份 ActiveMQ DemoC#资源面向使用 .NET 平台进行消息中间件开发与学习的程序员尤其适合刚接触 ActiveMQ、需要在 WinForm 项目中实践点对点消息收发的开发者。资源包内含发送端与接收端两套完整示例程序可帮助读者理解生产者与消费者之间的通信流程、消息队列的基本用法以及客户端连接配置方式。压缩包共 36 个文件以 19 个 cs 源码文件为核心配合 4 个 resx 资源文件、2 个 csproj 工程文件、1 个 sln 解决方案以及 4 个 dll 与 3 个 exe 等运行依赖整体约 326KB结构紧凑、便于直接编译调试。目前已有 1292 人学习下载说明该示例在同类入门资料中具有一定参考价值。通过阅读与运行这些代码读者可以快速掌握 ActiveMQ 在 C# 环境下的基础集成思路并在此基础上扩展自己的消息通信模块。1. 从一次消息丢失说起ActiveMQ 的 C# 接入到底难在哪线上一个订单同步服务C# 写的跑了大半年没出过事直到某天运维重启了消息中间件第二天对账发现少了三百多笔。翻日志才发现生产者用的是Connection级别的自动确认消费者那边异常退出时消息已经被标记消费但业务逻辑根本没执行完。这不是 ActiveMQ 本身的锅是接入姿势的问题。ActiveMQ 作为老牌 JMS 消息中间件在 Java 生态里资料铺天盖地但落到 C# 这边可参考的完整 Demo 少得可怜很多人第一次接就卡在 NMS 库的版本选择、确认模式和连接工厂参数上。这份 ActiveMQ DemoC#就是冲着这个缺口来的——它把生产者、消费者、事务会话、持久化订阅这几条主线用可运行的代码串起来适合正在做 C# 上位机、后台服务或系统集成、需要可靠消息队列但不想在环境配置上反复翻车的工程师。2. 环境搭建与 NMS 选型为什么不是直接引用 Apache.NMS2.1 先搞清楚 Apache.NMS 和 ActiveMQ 的关系C# 连 ActiveMQ 走的是 NMS.NET Message Service协议它跟 Java 那边的 JMS 是两套 API但底层都通过 OpenWire 协议跟 Broker 通信。NuGet 上搜 ActiveMQ 会出来一堆包核心其实就两个Apache.NMS是接口抽象层Apache.NMS.ActiveMQ才是真正的 ActiveMQ 实现。很多人只装了前者编译能过运行时抛NMSConnectionFactory找不到就是因为缺了实现包。常见做法是三个包一起装# 在项目目录下执行注意版本号要匹配 dotnet add package Apache.NMS --version 2.1.0 dotnet add package Apache.NMS.ActiveMQ --version 2.1.0 dotnet add package Apache.NMS.ActiveMQ.NetCore --version 2.1.0这里有个坑Apache.NMS.ActiveMQ和Apache.NMS.ActiveMQ.NetCore在 .NET Core / .NET 5 环境下需要同时存在前者提供核心实现后者补了 .NET Core 特有的依赖注入和日志适配。版本号必须一致混用 2.0 和 2.1 会在运行时抛TypeLoadException错误信息还特别隐晦只告诉你某个类型加载失败不说是版本冲突。2.2 连接工厂的参数怎么设才不翻车连接工厂是整个链路的入口参数设错后面全白搭。下面这段是 Demo 里生产者的初始化代码using Apache.NMS; using Apache.NMS.ActiveMQ; using System; // 创建连接工厂URI 格式tcp://主机:端口 IConnectionFactory factory new ConnectionFactory(tcp://localhost:61616); // 关键参数用户名密码在 Broker 未开启认证时可省略 // 但生产环境务必显式设置避免默认匿名连接 factory.UserName admin; factory.Password admin; // 连接超时设为 5 秒默认 30 秒在容器环境下容易拖垮启动流程 factory.RequestTimeout TimeSpan.FromSeconds(5); // 开启异步发送提升吞吐但要注意异常回调 IConnection connection factory.CreateConnection(); connection.Start(); // 创建会话第二个参数指定确认模式 ISession session connection.CreateSession(AcknowledgementMode.ClientAcknowledge);AcknowledgementMode是整段代码里最要命的参数。它有五个可选值AutoAcknowledge、ClientAcknowledge、DupsOkAcknowledge、SessionTransacted、TransactionMode。Demo 里默认用ClientAcknowledge意思是消费者必须显式调用message.Acknowledge()才算消费成功。如果你图省事用AutoAcknowledge消息一到客户端就被标记完成业务代码抛异常也不会重投这就是开头那个丢单场景的根因。RequestTimeout也值得单独说默认 30 秒在本地开发没感觉一旦 Broker 部署在另一个网段或者容器里 DNS 解析慢连接建立阶段就会卡满 30 秒日志上看就是启动特别慢实际是超时在等。2.3 队列和主题的创建差异ActiveMQ 里 Queue 和 Topic 是两种投递模型C# 这边的创建方式只差一个方法名// 点对点一条消息只被一个消费者消费 IDestination queue session.GetQueue(Order.Sync.Queue); // 发布订阅一条消息广播给所有订阅者 IDestination topic session.GetTopic(Order.Notify.Topic); // 创建生产者 IMessageProducer producer session.CreateProducer(queue); // 持久化模式非持久化消息在 Broker 重启后丢失 producer.DeliveryMode MsgDeliveryMode.Persistent; // 设置消息优先级0-9默认 4 producer.Priority MsgPriority.Normal; // 消息过期时间这里设 1 小时过期后进入死信队列 producer.TimeToLive TimeSpan.FromHours(1);DeliveryMode和TimeToLive这两个参数在 Demo 里是显式写出来的因为默认值不一定符合业务预期。Persistent模式下消息会落 KahaDBBroker 重启不丢但吞吐会降非持久化适合日志采集这类允许少量丢失的场景。TimeToLive设了之后过期消息不会自动删除而是进ActiveMQ.DLQ死信队列需要单独写消费者去处理否则死信队列会越堆越大最后占满磁盘。3. 生产者与消费者的完整实现从发送到确认的闭环3.1 生产者发送消息的三种写法Demo 里生产者分了同步发送、异步发送和事务发送三种模式对应不同业务场景// 方式一同步发送阻塞直到 Broker 确认 IMessageProducer producer session.CreateProducer(queue); ITextMessage textMessage producer.CreateTextMessage(订单数据 JSON); producer.Send(textMessage); // 方式二异步发送不阻塞主线程异常通过回调处理 producer.Send(textMessage, (result) { if (result.IsCompleted) { Console.WriteLine(发送成功); } else { Console.WriteLine($发送失败{result.Exception?.Message}); } }); // 方式三事务发送多条消息原子提交 using (ITransaction transaction session.BeginTransaction()) { producer.Send(producer.CreateTextMessage(消息1)); producer.Send(producer.CreateTextMessage(消息2)); transaction.Commit(); // 只有 Commit 后消息才真正入队 }同步发送最简单但每条消息都要等 Broker 的确认回执吞吐上不去。异步发送适合高吞吐场景但要注意回调是在 IO 线程上执行的里面不要做耗时操作否则会阻塞后续消息的发送。事务发送是批量场景的首选Commit之前消息都在客户端缓冲区Rollback就全部丢弃适合订单批量导入这种要么全成功要么全失败的逻辑。3.2 消费者确认模式与重投机制消费者这边最容易踩坑的是确认时机。Demo 里用ClientAcknowledge模式代码结构是这样的ISession session connection.CreateSession(AcknowledgementMode.ClientAcknowledge); IMessageConsumer consumer session.CreateConsumer(queue); consumer.Listener (IMessage message) { try { ITextMessage textMessage message as ITextMessage; // 业务处理解析 JSON、写数据库 ProcessOrder(textMessage.Text); // 业务成功后才确认 message.Acknowledge(); } catch (Exception ex) { // 不确认消息会在会话关闭后重新投递 Console.WriteLine($处理失败等待重投{ex.Message}); // 注意这里不能调用 Acknowledge } };关键点在于Acknowledge()的位置。放在业务逻辑之后失败就不确认Broker 会在消费者断开或会话关闭后重新投递。但这里有个隐藏问题如果消费者一直不确认也不断开消息会一直挂在 Broker 的“已投递未确认”列表里默认 30 秒后才会触发重投。这个超时由 Broker 的wireFormat.maxInactivityDuration控制不是客户端参数。Demo 里建议在异常分支里主动session.Recover()或者关闭连接让消息尽快回到队列。3.3 持久化订阅的实现细节Topic 模式下普通订阅者断开后消息就丢了持久化订阅才能保证离线期间的消息不丢// 创建持久化订阅clientId 必须唯一且固定 connection.ClientId OrderService.Client01; ISession session connection.CreateSession(AcknowledgementMode.ClientAcknowledge); ITopic topic session.GetTopic(Order.Notify.Topic); // 第二个参数是订阅名重启后要用同一个名字才能续上 IMessageConsumer consumer session.CreateDurableConsumer(topic, OrderSubscriber, null, false);ClientId和订阅名必须成对固定否则重启后会变成新订阅之前积压的消息就找不回来了。Demo 里把这两个值写在配置文件里而不是硬编码就是为了避免改代码时不小心改掉。另外持久化订阅的消息积压没有上限如果消费者长时间不启动Broker 磁盘会被撑爆生产环境要配合TimeToLive或者定期清理策略。4. 避坑与排查那些文档里不会写的翻车现场4.1 连接数暴涨导致 Broker 拒绝服务现象服务运行一段时间后Broker 日志出现Too many connections新连接全部失败。原因每次发消息都CreateConnection和Close连接池没复用。ActiveMQ 默认最大连接数是 1000但每个连接都会占一个线程实际能撑住的远低于这个数。解决连接和会话要复用Demo 里用单例模式管理IConnectionFactory整个应用生命周期只创建一次连接。如果确实需要多连接用ConnectionFactory的CreateConnection后缓存起来配合心跳检测保活。4.2 消息体过大导致发送超时现象发送 10MB 以上的消息时客户端抛RequestTimeoutException但 Broker 端显示消息已接收。原因ActiveMQ 默认消息大小限制是 100MB但RequestTimeout默认 30 秒大消息在网络传输和磁盘写入上耗时超过这个值客户端等不到确认就超时了。解决要么调大RequestTimeout要么把大消息拆成小块或者改用 Blob 消息把内容存文件系统只传路径。Demo 里建议超过 1MB 的消息就走 Blob 模式避免阻塞生产者线程。4.3 死信队列堆积拖垮磁盘现象Broker 磁盘使用率持续上涨检查发现ActiveMQ.DLQ队列消息数几十万。原因消费者处理失败后消息重投次数超过默认上限6 次自动进入死信队列但没人消费死信队列。解决写一个死信消费者把死信消息落库或者告警而不是让它无限堆积。Demo 里提供了一个简单的死信处理类把消息转存到数据库后确认。另外可以在 Broker 端配置maximumRedeliveries和redeliveryDelay控制重投策略。4.4 持久化订阅重启后收不到离线消息现象消费者重启后离线期间的消息没有收到像是被丢弃了。原因ClientId或订阅名变了Broker 认为是新订阅不会投递旧消息。解决把ClientId和订阅名写死在配置文件里不要用机器名或随机数。Demo 里用OrderService.Client01这种固定格式多实例部署时每个实例用不同的后缀但重启后必须保持一致。4.5 事务会话中 Acknowledge 报错现象在SessionTransacted模式下调用message.Acknowledge()抛InvalidOperationException。原因事务模式下确认由Commit自动完成手动确认是非法操作。解决事务会话里不要调Acknowledge业务成功后直接transaction.Commit()失败就Rollback。Demo 里把两种模式的代码分开写避免混淆。5. 进阶技巧用消息选择器做轻量级路由消息选择器Message Selector是 ActiveMQ 里被低估的功能它让消费者只接收符合条件的消息省掉了应用层的过滤逻辑。Demo 里有一个订单同步的场景多个消费者分别处理不同地区的订单用选择器就能实现// 生产者设置消息属性 ITextMessage message producer.CreateTextMessage(orderJson); message.Properties[Region] North; message.Properties[OrderType] Retail; producer.Send(message); // 消费者只接收华北地区的零售订单 IMessageConsumer consumer session.CreateConsumer(queue, Region North AND OrderType Retail);选择器的语法类似 SQL 的 WHERE 子句支持、、、、BETWEEN、IN、LIKE等操作符。但有几个限制要注意属性值只能是基本类型string、int、bool 等不能是对象选择器是在 Broker 端执行的属性越多、表达式越复杂Broker 的匹配开销越大。我一般建议属性控制在 3 个以内复杂路由逻辑还是放到应用层做。另一个实用技巧是消息分组Message Groups它能保证同一组消息按顺序被同一个消费者处理// 生产者设置分组 ID相同 ID 的消息会路由到同一个消费者 message.Properties[JMSXGroupID] order.CustomerId; producer.Send(message);这个特性在订单场景里特别有用同一个客户的订单必须按顺序处理但不同客户之间可以并行。JMSXGroupID是 JMS 规范里的标准属性ActiveMQ 原生支持。不过要注意如果某个消费者处理特别慢同组的消息会一直排队等它不会切换到其他消费者所以分组粒度不能太细否则并行度上不去。验证消息是否真的按预期投递我习惯在消费者端加一段统计代码private static int _receivedCount 0; private static int _ackCount 0; consumer.Listener (IMessage message) { Interlocked.Increment(ref _receivedCount); try { ProcessOrder((message as ITextMessage).Text); message.Acknowledge(); Interlocked.Increment(ref _ackCount); } catch (Exception ex) { Console.WriteLine($处理失败{ex.Message}); } }; // 定时打印接收和确认的差值差值持续增大说明有消息卡在处理中 Timer timer new Timer(_ Console.WriteLine($接收{_receivedCount}确认{_ackCount}差值{_receivedCount - _ackCount}), null, 0, 5000);这个差值监控帮我抓到过好几次问题有一次是数据库连接池满了业务处理卡住差值从 0 涨到几百及时扩容后恢复还有一次是某个消息体格式不对反序列化一直抛异常差值稳定在 1 不动说明有一条消息反复重投。从那以后我每次接 ActiveMQ 都强制加上这个监控比看 Broker 控制台直观得多。希望帮到你。本文还有配套的精品资源点击获取
返回列表