简介:这份 ActiveMQ Demo(C#)资源面向使用 .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 Demo(C#)就是冲着这个缺口来的——它把生产者、消费者、事务会话、持久化订阅这几条主线用可运行的代码串起来,适合正在做 C# 上位机、后台服务或系统集成、需要可靠消息队列但不想在环境配置上反复翻车的工程师。
2. 环境搭建与 NMS 选型:为什么不是直接引用 Apache.NMS
2.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模式下消息会落 KahaDB,Broker 重启不丢,但吞吐会降;非持久化适合日志采集这类允许少量丢失的场景。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 控制台直观得多。希望帮到你。
本文还有配套的精品资源,点击获取