简介:这套ActiveMQ Demo是面向C#开发者的消息队列示例程序,基于WinForm窗体实现,包含发送端与接收端两个独立工程。它主要解决初学者在搭建ActiveMQ环境后不知如何用C#进行消息投递和消费的痛点,适合正在学习消息中间件、或需要在项目里快速引入队列通信的读者。压缩包共包含36个文件,以19个C#源码文件为核心,配合4个resx界面资源文件、4个dll运行库以及3个exe可执行程序,同时提供解决方案和工程文件,整体仅326KB,布局清晰,便于直接打开和编译运行。资源中,发送程序负责连接队列并推送消息,接收程序负责监听并处理消息,界面操作直观,能清晰看到消息从生产者到消费者的流转过程。此外,源码还封装了MQ连接辅助类、列表排序控件等可复用模块,并保留了完整的WinForm设计文件,方便读者理解界面与后台逻辑的联动方式。目前已有1292人学习下载,对于希望以最小成本掌握ActiveMQ在C#环境下落地方法的学习者来说,是一份实用的参考样例。
1. 用 C# 客户端碰 ActiveMQ,最容易踩的不是 API,而是消息生命周期
很多从 RabbitMQ 或 Redis Stream 转过来的开发者,第一直觉是“ActiveMQ 是 Java 生态的老牌队列,C# 这边可能只是个半残客户端”,等到真用 Apache.NMS 写 Demo 时才发现,API 本身半小时就能跑通,真正让人加班的反而是几个反直觉的设计:连接创建后不主动Start(),消费者就永远收不到消息;订阅一个topic没改持久化,重启一次就把消息全弄丢;消息处理抛异常后重发策略没配置,订单请求一眨眼全进了死信队列。这篇文章不做概念堆砌,直接以“ActiveMQ Demo(C#)本地跑通”为主线,从最小可运行代码一路讲到连接参数、重试策略和线上排障。适合已经在项目里引入 ActiveMQ 但还没稳定运行、或者正准备从零评估这套方案的开发者。
2. 把 ActiveMQ Demo 跑起来:先解决“能连上、能收发”这件事
ActiveMQ 本身是消息中间件,C# 这边的官方客户端是 Apache.NMS.ActiveMQ,它不是一个简单的 HTTP 包,而是实现了消息头、会话、事务、确认模式等完整语义的客户端库。很多新人上来就想着像调用 Web API 一样发一条消息,忽略了两端都有的会话与确认机制,结果收发看起来成功,消息却在 broker 重启或网络抖动时丢失。这一章不做高深原理,先把最小 Demo 在本地完整跑通。
2.1 先起一个 ActiveMQ broker,再建一个 C# 工程
第一步是环境准备。假设你现在手上没有现成 broker,最常见做法是拉一个官方镜像或下载压缩包在本地直接启动,这里不展开部署细节,重点说 C# 工程这边。
dotnet new console -n ActiveMqDemo cd ActiveMqDemo dotnet add package Apache.NMS.ActiveMQ这段命令创建了一个控制台工程,并让 NuGet 把 Apache.NMS.ActiveMQ 拉进项目。这里要注意目标框架的选择:如果你的工程是老式 .NET Framework 4.x,包管理器会自动选兼容版本;如果是 .NET 6/8,默认目标框架也能找到对应的构建。我自己习惯先把工程建好,再检查csproj里的TargetFramework,如果遇到编译时类型定位不到的报错,多数是目标框架过低或 NuGet 源没有配置好。
broker 起来之后,本地默认监听tcp://127.0.0.1:61616。可以用浏览器访问控制台确认61616端口已经可连通,也可以直接用下面的命令测一下:
telnet 127.0.0.1 61616能连上说明 ActiveMQ 进程正常,接下来写第一个最简单的生产者。
2.2 生产者:一条消息从连接到发送要过三关
下面这段代码是一个最小生产者,它做的事情是创建连接工厂、创建连接、创建会话、创建目标、发送一条文本消息。
using Apache.NMS; using Apache.NMS.ActiveMQ; const string brokerUri = "tcp://127.0.0.1:61616"; var factory = new ConnectionFactory(brokerUri); using var connection = factory.CreateConnection(); connection.Start(); using var session = connection.CreateSession(AcknowledgementMode.AutoAcknowledge); var destination = new ActiveMQQueue("demo.orders"); using var producer = session.CreateProducer(destination); producer.DeliveryMode = MsgDeliveryMode.Persistent; var message = session.CreateTextMessage("hello demo, this is order message"); message.Properties.SetString("AppName", "demo-sender"); producer.Send(message); Console.WriteLine("sent");这里有三点必须说明白。第一,CreateConnection之后必须调用connection.Start(),在 ActiveMQ 的 NMS 客户端里,连接未启动时 broker 不会主动向客户端推送消息,很多新手在这里漏掉一句,导致消费者永远收不到消息。第二,CreateSession的入参是确认模式,示例用了AutoAcknowledge,意思是消息被 broker 分发后自动确认;如果消息处理逻辑很轻,这没问题,但如果你要在消费端做业务幂等,通常需要换成ClientAcknowledge。第三,producer.DeliveryMode = MsgDeliveryMode.Persistent表示消息要落盘到 broker 的持久化存储中,而不是只存在内存里。默认情况下 ActiveMQ 队列消息就是持久化的,但显式写出来能避免后面调试时误改配置。
2.3 消费者:同步 Receive 和异步 Listener 二选一
消费者这边同样有创建会话和目标的步骤,但接收消息时有两种完全不同的写法:一种是同步阻塞Receive,一种是异步事件回调Listener。同步写法适合一次性调试,异步写法适合常驻运行的服务。
先看同步写法:
using Apache.NMS; using Apache.NMS.ActiveMQ; const string brokerUri = "tcp://127.0.0.1:61616"; var factory = new ConnectionFactory(brokerUri); using var connection = factory.CreateConnection(); connection.Start(); using var session = connection.CreateSession(AcknowledgementMode.AutoAcknowledge); var destination = new ActiveMQQueue("demo.orders"); using var consumer = session.CreateConsumer(destination); var receiveTime = TimeSpan.FromSeconds(5); var message = consumer.Receive(receiveTime) as ITextMessage; if (message != null) { Console.WriteLine($"received: {message.Text}"); Console.WriteLine($"property AppName = {message.Properties.GetString("AppName")}"); }Receive(TimeSpan)会阻塞当前线程直到超时或收到消息。超时不是错误,返回null表示当前没有可消费消息。同步方式写起来直观,但在真实服务里很少把主线程卡在Receive上,因为消息不多时线程空转,消息多了又不好控制并发。
异步方式更贴近真实项目:
using var consumer = session.CreateConsumer(destination); consumer.Listener += message => { if (message is ITextMessage textMessage) { Console.WriteLine($"received async: {textMessage.Text}"); } }; connection.Start(); Console.ReadLine();注意这里的connection.Start()挪到了注册 Listener 之后。异步回调由 NMS 客户端内部的线程池触发,回调里不应该做耗时操作,如果你需要写数据库或调用外部接口,把消息内容转存到一个线程安全队列,再交给业务线程慢慢处理。如果直接在回调里写数据库,一旦数据库慢,客户端线程池就会被占满,后面消息的消费延迟会被瞬间拉到几十秒。
var pendingMessages = new System.Collections.Concurrent.ConcurrentQueue<ITextMessage>(); consumer.Listener += message => { if (message is ITextMessage textMessage) { pendingMessages.Enqueue(textMessage); } };把回调做成“只入队、不处理”的模式,就能把 ActiveMQ 客户端自身的线程和业务线程彻底解耦。这个习惯能帮你省掉不少排查“为什么消费者卡住”的时间。
3. 连接 URI、确认模式与重发策略:把 Demo 变成能上线的可靠性基座
Demo 能收发只是第一步,真正让 ActiveMQ 在 C# 项目里站稳的,是连接参数和消息确认策略。很多人上线前没调这些参数,等到 broker 重启或网络闪断时才暴露问题。这一章讲三件事:连接 URI 的故障转移写法、会话确认模式的选择、消息重发和死信队列的关系。
3.1 连接 URI:只写 tcp://127.0.0.1:61616 会有单点风险
先看一段实际项目里最常见的连接字符串写法:
failover:(tcp://127.0.0.1:61616,tcp://127.0.0.1:61617)?randomize=false&initialReconnectDelay=1000&maxReconnectDelay=30000这串 URI 由三部分构成:协议failover、括号内的 broker 地址列表、问号后的连接参数。failover是 ActiveMQ 客户端内置的故障转移传输层,它的作用是:如果当前 broker 断开了,客户端会在后台自动尝试列表里的下一个地址,并按配置的延迟时间反复重连。
参数randomize=false表示始终按列表顺序连接,不做随机负载均衡。如果你的多个 ActiveMQ 节点配置了共享存储,可以开randomize=true让客户端分散连接。initialReconnectDelay=1000是第一次重连前的等待毫秒数,maxReconnectDelay=30000是重连延迟上限,配合useExponentialBackOff=true时延迟会指数增长,避免连接频繁抖动打爆 broker。
另一个容易被忽略的参数是sendTimeout。在 failover 模式下,当所有 broker 都不可用时,Send调用会阻塞当前线程。如果生产端不希望消息一直卡死,可以这样写:
failover:(tcp://127.0.0.1:61616)?sendTimeout=5000意思是发送动作最多阻塞 5 秒,超时后抛出异常,由调用方决定丢缓冲还是走降级流程。没有这个参数时,客户端可能无限期等待 broker 恢复,看起来像进程死锁。
3.2 确认模式:Auto、Client 与 Individual,差别在“消费才算数”还是“收到就算数”
ActiveMQ 的确认语义很容易让人混淆。下面这张表列出 NMS 客户端常用的三种确认模式:
| 确认模式 | 行为特征 | 典型场景 |
|---|---|---|
| AutoAcknowledge | 客户端接收消息后自动确认,无需手动调用 | 日志处理、通知推送,允许重复或无需业务级幂等 |
| ClientAcknowledge | 必须调用message.Acknowledge()才算确认 | 数据库写入、余额变更,需要业务成功后再确认 |
| IndividualAcknowledge | 可以逐条确认,也可批量确认 | 批量消费场景,避免一荣俱荣 |
用ClientAcknowledge的代码大概是这个样子:
using var session = connection.CreateSession(AcknowledgementMode.ClientAcknowledge); using var consumer = session.CreateConsumer(destination); var message = consumer.Receive(TimeSpan.FromSeconds(5)) as ITextMessage; if (message != null) { try { // 处理业务 Console.WriteLine($"do business: {message.Text}"); message.Acknowledge(); } catch { // 不确认,ActiveMQ 会重新投递 } }这段代码的核心是:业务处理成功后再确认,否则消息会被 broker 重新投递。需要特别提醒,ClientAcknowledge模式下,同一会话中如果消费了多条消息,调用其中任意一条Acknowledge()可能会让整个会话里前面的消息一起被确认。如果你希望每条消息独立控制,请使用IndividualAcknowledge。这个边界在实际项目中经常引发连锁问题,比如一条失败消息影响了后面的成功消息。
3.3 重发与死信:消息处理失败不能无限循环
消息被重新投递时,默认会不断重试,直到客户端确认成功或超过重发上限。无限重发给业务系统带来的后果往往很严重,所以 ActiveMQ 引入了RedeliveryPolicy。在 NMS 客户端里,连接创建前就可以配置这个策略。
var factory = new ConnectionFactory(brokerUri); factory.RedeliveryPolicy.MaximumRedeliveries = 3; factory.RedeliveryPolicy.InitialRedeliveryDelay = 1000; factory.RedeliveryPolicy.UseExponentialBackOff = true; factory.RedeliveryPolicy.BackOffMultiplier = 5;MaximumRedeliveries = 3表示最多重发 3 次,第 4 次投递失败后会进入死信队列。InitialRedeliveryDelay = 1000是第一次重发前的延迟,BackOffMultiplier = 5会让后续重发间隔乘以 5,也就是 1 秒、5 秒、25 秒。这个配置适合业务处理依赖外部接口的场景,给外部系统留出恢复时间。
死信队列默认叫ActiveMQ.DLQ,在控制台里可以直接看到积压数量。当你发现某个业务消息大量进入死信队列,不要只盯着 ActiveMQ 的控制台看,更要看消费端异常日志。死信队列本身不是问题,问题是不该死信的普通业务错误被当成异常反复重发。常见做法是死信队列只存“最终失败”的消息,并单独配置一个消费者扫描死信队列做人工补偿。
3.4 事务会话:一次批量处理,全部成功或全部不确认
如果你要让多条消息在一个事务里统一确认,NMS 客户端也支持事务会话。关键变化是把CreateSession的第一个参数设为SessionTransacted:
using var session = connection.CreateSession(AcknowledgementMode.SessionTransacted); using var consumer = session.CreateConsumer(destination); var messages = new List<ITextMessage>(); for (int i = 0; i < 5; i++) { var msg = consumer.Receive(TimeSpan.FromSeconds(1)) as ITextMessage; if (msg != null) messages.Add(msg); } // 业务处理 if (messages.Any()) { Console.WriteLine($"process {messages.Count} messages"); } session.Commit();session.Commit()会一次性确认会话中所有未确认消息。如果中途抛异常,你可以调用session.Rollback(),broker 会把当前会话里的消息重新投递。事务模式适合批量对账、批量导入这类场景,不建议在单条消息处理时使用,因为事务开销会比自动确认高不少。
4. ActiveMQ C# Demo 常见问题与避坑:六条真实踩坑记录
这一章直接上场景化问题,每一条都是我自己或身边同事实际踩过、最后靠日志或抓包才定位到根因的。你可以把这章当作排障手册:下次消息丢了、重复了、卡住了,先对照这里的现象。
4.1 消费者已经注册成功,为什么收不到消息
现象是消费者代码看着没问题,broker 控制台里能看到连接和消费者,但消息一直堆积,消费端没有任何输出。
原因大概率是漏了connection.Start()。NMS 客户端设计上把Start独立出来,是为了让你可以在连接真正开始分发消息之前完成资源初始化。如果没调用,连接处于有效但未激活状态,broker 不会把消息发给消费者。解决方法是确保CreateConsumer之后、开始业务逻辑之前调用connection.Start()。如果用了异步 Listener,Start放在注册回调之后也可以,但不要放在CreateConnection之后就不管了。
4.2 订阅 topic 的消息在重启后全部丢失
现象是运行 Demo 时能正常收到 topic 消息,但把程序重新启动一次,再发消息就收不到了。
原因是用普通CreateConsumer订阅的是“非持久订阅”。非持久订阅是即时订阅,消费者离线期间的消息不会保留。解决方法是改用持久订阅,但要注意持久订阅在 ActiveMQ 里需要有一个全局唯一的订阅名称。
string subscriptionName = "demo-durable-sub"; var consumer = session.CreateDurableConsumer( new ActiveMQTopic("demo.events"), subscriptionName, null, false);这里subscriptionName是持久订阅的标识,同一客户端必须保持稳定。持久订阅建立后,broker 会把消息保存在订阅者对应的存储区域,消费者重新上线后继续消费。这个写法对“离线期间不能丢事件”的场景非常关键。
4.3 Listener 回调里同步处理数据库,消息越积越多
现象是程序运行一段时间后,ActiveMQ 控制台队列数量不断增长,操作系统日志里出现线程池耗尽或死锁。
原因是在Listener回调里直接写了数据库操作或调用外部 API。当业务响应慢,回调线程被占用,客户端线程池耗尽,后续消息无法被分发,即使 broker 端消息积压也无法消费。解决方法是把回调里塞进一个线程安全队列,再用独立的消费者线程去处理业务。
var pending = new System.Collections.Concurrent.ConcurrentQueue<ITextMessage>(); consumer.Listener += msg => { if (msg is ITextMessage tm) pending.Enqueue(tm); };这种“回调只入队、业务线程出队”的模式能把客户端线程按纳秒级别释放,业务积压时也只表现为业务线程在等待,而不是整个客户端链路阻塞。
4.4 failover 重连后消息重复投递,消费端出现重复数据
现象是 ActiveMQ 节点重启,客户端自动重连成功后,数据库里出现重复订单或重复通知。
原因是 ActiveMQ 的投递语义是 at-least-once,也就是至少一次。当 broker 把消息发给消费者后因为网络断开,客户端来不及返回确认,重连后 broker 会重新投递这条消息。这不是客户端 bug,而是分布式消息中间件的固有语义。解决办法是消费端做幂等,比如依据消息里的业务唯一键在数据库建唯一索引,或在消费前用 Redis 做一次SETNX。Demo 阶段也要把幂等测试写进去,否则上线后再补很痛苦。
if (message.Properties.GetString("BusinessId") is string businessId) { // 这里通过唯一键判重,防止重复消费 if (HasProcessed(businessId)) return; ProcessBusiness(message); }注意幂等不能依赖 ActiveMQ 自带的JMSMessageID,因为同一个业务消息在客户端重启后可能重新生成消息对象,但业务单据唯一键是稳定不变的。
4.5 消息处理一直抛异常,被无限重发,最终积压爆掉
现象是消费者不断打印异常,ActiveMQ 控制台里队列数据量只增不减,最后磁盘告警。
原因是使用ClientAcknowledge或事务模式时,处理异常后没有调用Acknowledge或Commit,broker 会一直等待确认,并以默认策略持续重发。解决方法是设置RedeliveryPolicy控制最大重发次数,并配置死信队列。
factory.RedeliveryPolicy.MaximumRedeliveries = 3; factory.RedeliveryPolicy.InitialRedeliveryDelay = 500; factory.RedeliveryPolicy.UseExponentialBackOff = true; factory.RedeliveryPolicy.BackOffMultiplier = 4;同时,在业务代码里要区分“可重试错误”和“不可重试错误”。可重试错误比如 HTTP 503、数据库连接超时,可以抛出去让 ActiveMQ 重发;不可重试错误比如参数错误、数据格式错误,应该直接走告警或人工补偿,不要让它进入重发循环。
4.6 死信队列一直有消息,但没人处理
现象是ActiveMQ.DLQ里消息越来越多,主队列看起来没有动静,业务人员反馈订单状态没有更新。
原因是死信队列默认没有消费者,它只负责存放重发超限的消息。如果你没有写死信扫描程序,这些消息就会永远堆在那里。我常用的做法是给死信队列单独建一个消费者,并按固定频率统计和导出异常消息,交给开发人员分析。不要为了图省事把死信队列直接删掉,那是把问题的证据销毁了。
var dlqDestination = new ActiveMQQueue("ActiveMQ.DLQ"); using var dlqConsumer = session.CreateConsumer(dlqDestination); // 这里消费死信队列,记录内容,并通知相关人员5. 给 Demo 加一个故障切换验证:用两个 broker 和一份去重表把坑提前踩完
最后一章讲一个具体收尾动作:在本地起一个包含两个端口的最小 broker 集群,用 C# 客户端做故障切换验证,顺带把幂等消费写进去。这套验证做完,再上生产你会踏实很多。
先准备两个 ActiveMQ broker,分别监听61616和61617,并且在配置里开启共享存储,共享存储保证两个节点能看到同一份队列数据。随后 C# 客户端连接串改成:
failover:(tcp://127.0.0.1:61616,tcp://127.0.0.1:61617)?initialReconnectDelay=1000&maxReconnectDelay=5000验证步骤大致是这样:启动生产者向demo.orders发送 100 条消息,保持消费者运行,然后直接停掉当前正在连接的 broker 进程,观察客户端是否自动切换到第二个端口。正常情况下,消费者端能看到短暂的连接断开日志,然后继续收货,生产者端则因为 failover 自动重连,不需要重启进程。
我习惯在生产者发送循环里记录每条消息的业务编号,并在消费者端用一个HashSet做去重。如果业务编号在集合里已存在,说明发生了重复投递,记录一次重复事件。下面的伪代码演示这个验证逻辑:
var processedIds = new HashSet<string>(); consumer.Listener += msg => { if (msg is ITextMessage textMessage && textMessage.Properties.GetString("BusinessId") is string id) { if (!processedIds.Add(id)) { Console.WriteLine("duplicate detected: " + id); } } };当 broker 断线重连后,如果日志中出现了duplicate detected也不要慌张,它证明你的去重校验是有效的,生产环境可以直接套用。如果日志里没有重复,那一方面说明投递足够顺,另一方面也提醒你不能因此省略幂等,因为网络抖动的概率在生产环境只高不低。
最后我会做一遍“慢消费者”验证:消费者处理消息时Thread.Sleep(2000),然后观察 ActiveMQ 控制台的 prefetch 窗口和积压数量。这个验证能帮你理解prefetchSize、重发策略和死信队列三条链路之间的关系,比读十篇文档都管用。
这一步做完,你的 ActiveMQ Demo 才算真正具备参考价值。我自己最初就是在 Demo 阶段跳过故障切换验证,上线一周后被 broker 重启打了一个措手不及,之后才老老实实把 failover 和幂等写进每个新项目里。希望这篇文章能帮你少走这段弯路。
本文还有配套的精品资源,点击获取