☰
ActiveMQ C# Demo实战:从Queue到持久化消息的避坑指南
2026/10/11 5:08:42 网站建设 项目流程

简介:面向C#开发者的ActiveMQ消息中间件入门Demo,采用WinForm实现完整的消息发送与接收流程,适合需要在.NET环境中快速体验消息队列机制的初学者。程序将生产者与消费者集成在同一界面中,通过直观操作展示消息异步传递的核心逻辑,覆盖了ActiveMQ与C#客户端集成的基本配置与调用方式。压缩包共36个文件,以cs源代码、resx界面资源、dll依赖库和exe可执行程序为主,另含csproj工程文件及解决方案文件,项目结构清晰,便于直接编译运行和对照学习。资源整体仅326KB,轻量易用,已有1292人学习下载。除基础收发功能外,还包括GlobalFunction全局函数、ListViewColumnSorter列表排序辅助组件等封装,可帮助读者理解消息应用中界面交互与数据展示的常见处理方式,适合作为消息中间件项目开发的参考起点。

1. 先说结论:这份 ActiveMQ C# Demo,是 .NET 工程师最快落地消息队列的起点

很多 .NET 开发第一次接触 ActiveMQ C# Demo 的时候,以为它就是“发一条消息再收一条消息”的控制台演示;真跑到生产环境才发现,消息队列的坑根本不在 C# 语法上,而在连接串、持久化和确认机制里。我当年接手一个被频繁重启的 Windows 服务,队列消费者明明在跑,但服务一更新,积压的消息全部消失,排查了三天才意识到是发送端消息没有持久化。这套 Demo 的价值在于把 Queue、Topic、持久消息、消费者监听这些基础链路一次跑通,能让刚接触消息中间件的人直观看到消息从生产到落盘的完整过程,也适合老手拿来当验证环境的敲门砖。把 Demo 跑起来,比翻半天官方文档有用得多。

2. 搭建可跑环境:Broker、端口和连接串哪一个都不能错

2.1 一份能跑的 Demo 里,到底装了哪些东西

我拆过不少 ActiveMQ 相关的 C# 示例工程,这一份算是比较标准的。压缩包里通常不是只有一个“Hello World”,而是拆成几个独立项目:一个放服务端配置和启动脚本,两个放生产者和消费者的 .NET 项目,另外还有一个公共类库用来封装连接工厂。使用前先大概看一眼文件结构,能少走很多弯路。

“服务端配置”这一块需要特别留意,因为 ActiveMQ 的默认配置对初学者是不透明的。解压之后你会看到conf/activemq.xml、conf/jetty.xml这些配置文件,它们决定了 broker 监听哪个端口、数据存到哪、管理控制台能不能访问。Demo 里给出的默认值一般是 OpenWire 端口 61616,Web 管理台端口 8161,除非你手动改过activemq.xml,否则端口不要乱动。

我平时拿到一个消息队列 Demo,第一步会整理出下面这张文件清单,对照自己的环境再决定先改哪里:

文件 / 目录作用常见误区
conf/activemq.xmlbroker 核心配置,包括持久化适配器、传输连接器、策略直接改端口但忘记改防火墙,导致 61616 连不上
conf/jetty.xmlWeb 控制台宿主配置8161 起不来往往是 JDK 兼容问题,不是 jetty 配置问题
examples/或src/ProducerC# 生产者示例只看代码不跑 broker,永远等不到消费者回执
src/ConsumerC# 消费者示例忘记调用connection.Start(),Receive 一直超时
Common或Core连接工厂、配置常量的公共封装连接串里的failover被误当成普通协议前缀

很多人在这一步就卡住了。比如启动脚本在 Windows 上需要管理员权限,或者 Java 环境变量没配好,broker 就直接闪退,Demo 代码本身没问题也会跟着“背锅”。所以我一般建议先把服务端跑起来,用管理台打开再看客户端,能一次性排除掉大多数环境问题。

2.2 先把 Broker 拉起来:启动、端口和第一次心跳

ActiveMQ 的启动方式非常传统。解压到纯英文目录之后,进入bin目录,Windows 下直接执行activemq.bat start,Linux 下执行./activemq start。启动成功之后,进程不会立刻退出,而是在后台跑一个 broker 实例。为了确认它确实起来了,我会顺手做两件事:看端口占用,再看管理台。

cd /opt/activemq/bin ./activemq start netstat -ano | grep 61616 curl http://127.0.0.1:8161

61616是 OpenWire 协议端口,C# 客户端通过这个端口和 broker 通信;8161是管理控制台端口。启动脚本正常会输出ActiveMQ JMS Message Broker的相关信息。如果61616没有监听,说明 broker 没起来,这时候要看data/activemq.log而不是反复重启。

C# 端对应的准备工作是把客户端库引入项目。如果你用 .NET 6 以上的环境,直接通过包管理器安装即可:

dotnet add package Apache.NMS.ActiveMQ

这套库的名字里带有历史原因,但它就是 ActiveMQ 官方的 .NET 客户端实现,底层走 OpenWire 协议。安装完以后,写一个最简单的连接测试代码,能成功创建连接就说明包引用和 broker 端口都是通的。

var factory = new Apache.NMS.ActiveMQ.ConnectionFactory( "failover:(tcp://127.0.0.1:61616)"); using (var connection = factory.CreateConnection()) { connection.Start(); Console.WriteLine("connected"); }

这里没有调用Start()之前,连接其实已经建立了,但不会主动去接收 broker 推过来的消息。很多新手把CreateConnection()和Start()当成一回事,结果消费者收不到任何消息。Start()的语义是“让连接开始投递消息”,创建连接只是握手成功,握手成功和消息投递是两码事。

2.3 连接串为什么这样写:tcp 直连与 failover 重连

Demo 里最常见到的连接串写法是下面两种,它们不是同一个层面的东西。

第一种是直连:

tcp://127.0.0.1:61616

第二种是带故障转移的写法:

failover:(tcp://127.0.0.1:61616)?startupMaxReconnectAttempts=10

直连适合调试。tcp://后面跟 IP 和端口,broker 不在了就立刻抛异常,错误信息干净。failover则是把多个 broker 地址用括号包起来,客户端在断线之后会自动尝试重连,默认会一直重连,所以生产环境普遍用它。

这中间有个容易被忽略的细节:failover不是 ActiveMQ 特有的魔法,它的底层是客户端心跳和通道重建。Demo 里如果只有直连配置,说明它只给你演示最小链路,不是为了让你直接上生产的。你后续要改的往往是failover:(tcp://host1:61616,tcp://host2:61616)这种高可用写法,而不是去改一堆timeout参数。

连接串上的参数也属于“少配一个就出大事”的类型。常见的几个我整理在这里:

参数含义建议
startupMaxReconnectAttempts启动时最大重连次数调试设 3 左右,生产不要设太小
timeoutfailover 重试间隔默认看着合适,网络抖动时可调 3000
keepAlive是否开启 TCP 心跳长时间空闲连接建议为 true
transport.connectTimeout建连超时时间跨机房环境建议加大到 10000

把连接串理解到这个程度,再去看 Demo 里的ConnectionFactory就不会一脸茫然了。

3. Queue 模式拆解:生产者、消费者与消息确认

3.1 生产者代码:发给“队列”的消息到底去哪了

Queue 模式下,生产者发出的每一条消息都会被 broker 保存下来,直到某个消费者取走并确认。先不看复杂的事务和确认策略,Demo 里最基础的生产者代码大概长这样:

public class DemoProducer { public void Send(string message) { var factory = new Apache.NMS.ActiveMQ.ConnectionFactory( "failover:(tcp://127.0.0.1:61616)"); using (var connection = factory.CreateConnection()) using (var session = connection.CreateSession( AcknowledgementMode.AutoAcknowledge)) { var queue = session.GetQueue("demo.queue"); using (var producer = session.CreateProducer(queue)) { producer.DeliveryMode = MsgDeliveryMode.Persistent; producer.Priority = MsgPriority.Normal; producer.Send(session.CreateTextMessage(message)); } } } }

这段代码里最有价值的是那两行参数设置。DeliveryMode决定消息是否落到磁盘,Persistent表示持久化,broker 重启后消息还在;如果改成NonPersistent,消息只存在内存里,broker 一重启就没了。Priority是消息优先级,范围 0 到 9,默认是 4,优先级越高越容易被提前消费。

这里要纠正一个常见误解:在持久化模式下,Send()返回并不代表消息已经写入物理磁盘。ActiveMQ 默认开启了异步发送,Send()只是把消息交给底层的传输通道,真正的写入动作在后台完成。如果一定要等 broker 确认落盘,需要在连接工厂上关闭异步发送:

var factory = new Apache.NMS.ActiveMQ.ConnectionFactory( "tcp://127.0.0.1:61616"); factory.AsyncSend = false;

把这个开关设为 false 之后,每次Send()都会同步等待 broker 的写盘确认,吞吐量会下降,但可靠性更高。Demo 里通常不会写这个开关,因为它会拖慢演示速度;但你在做订单类系统时,关闭异步发送是更稳妥的选择。

3.2 消费者代码:Receive 轮询和 Listener 回调怎么选

队列的消费方式有两种,Demo 里一般都会展示一种。第一种是Receive()阻塞轮询,适合定时任务或者测试脚本:

var factory = new Apache.NMS.ActiveMQ.ConnectionFactory( "failover:(tcp://127.0.0.1:61616)"); using (var connection = factory.CreateConnection()) using (var session = connection.CreateSession( AcknowledgementMode.AutoAcknowledge)) { var queue = session.GetQueue("demo.queue"); var consumer = session.CreateConsumer(queue); connection.Start(); var message = consumer.Receive(TimeSpan.FromSeconds(5)); if (message is ITextMessage textMessage) { Console.WriteLine(textMessage.Text); } }

第二种是注册Listener事件,消息一到就触发回调,适合持续运行的服务:

var consumer = session.CreateConsumer(queue); consumer.Listener += OnMessage; connection.Start(); Thread.Sleep(TimeSpan.FromSeconds(30)); static void OnMessage(IMessage message) { if (message is ITextMessage textMessage) { Console.WriteLine($"received: {textMessage.Text}"); } }

这两种方式的本质区别是“拉”和“推”。Receive()每调用一次就从本地缓冲区拿一条消息,拿不到就等超时;Listener则是由 broker 主动推消息过来,客户端在回调里处理。由于回调是单线程的,OnMessage里尽量不要做耗时超过几十毫秒的操作,否则后续消息会被堵住。

很多 Demo 跑不通都是同一个原因:connection.Start()放在了消费者创建之前,或者干脆漏掉了。没有Start()时,连接虽然建立了,但消息不会进入回调,Receive()会一直超时返回 null。把Start()的位置记牢,比记住一堆理论参数有用得多。

3.3 Queue 与 Topic:Demo 里为什么两套代码长得一样

ActiveMQ 里 Queue 和 Topic 的 C# 代码几乎一模一样,唯一区别是GetQueue换成GetTopic。但语义差别很大:

维度QueueTopic
消息分发一条消息只能被一个消费者消费一条消息广播给所有在线订阅者
离线消息消费者离线消息也在队列里,上线后继续消费默认离线收不到
持久化订阅不需要额外参数需要持久订阅者标识
典型场景任务分发、削峰填谷事件广播、通知

Demo 里如果同时出现这两类代码,我建议优先把 Queue 跑通,再切 Topic。原因很简单:Topic 有一个非常容易踩的坑,消费者必须在自己上线之后、生产者发布之前完成订阅,否则消息会在 broker 里找不到目标订阅者而直接丢弃。Queue 则没有这种时序问题。

一旦需要“消费者离线时也要收到 Topic 消息”,就得用持久订阅者:

var topic = session.GetTopic("demo.topic"); var consumer = session.CreateDurableConsumer( topic, "consumer-id-01", null, false);

consumer-id-01是订阅者的唯一标识,broker 会为这个标识保存离线消息。这里要特别注意,这个 ID 一旦改了,broker 认为是新的订阅者,旧的订阅者消息就不会继续投递。很多人在代码里动态生成这个 ID,也导致消息“丢失”。同一个消费者必须使用固定的持久订阅标识。

4. 避坑排查:重启丢消息、连接被断开、死信循环

4.1 重启 broker 后队列消息全没了

现象:生产者已经成功发送,消费者也确认收到,但 broker 重启之后,队列里积压的消息全部消失。

原因:发送端没有把消息设为持久化,或者 broker 的持久化适配器没生效。很多 Demo 为了演示方便,消息都默认走内存存储,进程一重启,内存里的队列数据就清了。

解决:先把生产者里的producer.DeliveryMode改成MsgDeliveryMode.Persistent,再检查activemq.xml里的持久化配置。最常见的正确写法是:

<persistenceAdapter> <kahaDB directory="${activemq.data}/kahadb"/> </persistenceAdapter>

这个配置告诉 broker 使用 KahaDB 作为存储引擎。改完配置以后一定要重启 broker 再验证一次,先发几条消息,重启 broker,再启动消费者,看消息是否还在。

4.2 客户端空闲一段时间后被强制断开连接

现象:消费者服务挂在那里,日志一片平静,几个小时后突然抛Connection reset,重连之后又开始正常接收。

原因:ActiveMQ 的传输连接器启用了不活跃检查,客户端和 broker 之间如果长时间没有数据传输,broker 会认为连接已经死了,主动把它关掉。C# 监听模式下,如果消息本身不发不送,连接上没有任何心跳,就很容易触发这个机制。

解决:在 broker 的 transport connector 上开启 keepAlive,或者在连接串上带心跳参数,二选一即可。我一般两个都会加:

<transportConnector name="openwire" uri="tcp://0.0.0.0:61616?keepAlive=true"/>

同时,C# 端连接串写成:

failover:(tcp://127.0.0.1:61616)?keepAlive=true

这样即使业务消息稀疏,底层 TCP 也会定期发送心跳包,避免连接被 broker 误杀。注意,keepAlive=true只是发送心跳,不会改变业务消息的投递语义。

4.3 消费回调里抛异常,消息被反复投递最后刷屏

现象:Listener回调里处理业务时抛了一个未捕获异常,随后日志里出现大量相同消息的接收记录,消费者像卡死一样不停处理同一条消息。

原因:Listener回调在消息没有确认成功的情况下,broker 默认会重新投递该消息,同时重投次数不受限制。异常抛出来之后,确认根本没发生,消息就一直在“投递、失败、再投递”的循环里。

解决:在回调入口就捕获所有异常,并限制重投次数:

var redeliveryPolicy = new RedeliveryPolicy { MaximumRedeliveries = 3, InitialRedeliveryDelay = 2000, UseExponentialBackOff = true }; factory.RedeliveryPolicy = redeliveryPolicy;

然后消费回调里加上 try/catch:

consumer.Listener += message => { try { HandleMessage(message); } catch (Exception ex) { Console.WriteLine(ex.Message); } };

MaximumRedeliveries设成 3,意味着同一消息最多重投 3 次,第 4 次会被送进死信队列。这项配置需要小心,MaximumRedeliveries=0时消息一次失败就直接进死信;设太大又会造成消息延迟。Demo 阶段设 2 到 3 次比较合理。

4.4 Java 端发的 ObjectMessage 在 C# 里反序列化失败

现象:消息已经能收到,但强转ITextMessage时报类型转换异常,控制台输出一堆InvalidCastException,消息队列状态显示消息被“卡死”在队列里反复投递。

原因:ActiveMQ 的ObjectMessage使用 Java 序列化格式,C# 客户端无法直接解析。这不是 C# 库的功能缺失,而是跨语言传输时选错了消息类型。很多团队早期用 Java 生产者发对象消息,后来接了一个 C# 消费者,就碰到这种问题。

解决:跨语言场景尽量统一用TextMessage或BytesMessage,不要在 C# 和 Java 之间混用ObjectMessage。如果对方已经用ObjectMessage发了,C# 端能做的只能是把消息体当成原始字节取出来,然后按约定的序列化协议解析,比如把内部字段封装成 JSON 字符串。这种问题在 Demo 里看不出来,因为 Demo 的发和收都是同一种语言。

4.5 持久订阅者改了客户端 ID,消息立刻丢失

现象:使用了CreateDurableConsumer之后,第一次运行能收到 Topic 消息,第二次改了一个订阅者 ID,离线期间的消息就收不到了。

原因:持久订阅者的离线消息是绑定在“客户端 ID + 订阅名称”上的。ID 一旦变化,broker 认为这是一个全新订阅,重新开始记录,旧订阅下的消息就成了无人认领的孤儿消息。

解决:把持久订阅者 ID 提升为配置项,写入配置文件或环境变量,不要写死在代码里。固定同一个 ID,消费者重启、离线、断网重连都不会丢 Topic 消息。

5. 进阶验证:把 Demo 改成事务会话,让消息可回滚

5.1 Transactional Session:确认权从 broker 交回给业务代码

Demo 里的会话通常使用AutoAcknowledge,消息一经接收就自动确认,业务代码没有第二次机会。想给消息加一道“后悔药”,就把会话改成事务模式:

using (var session = connection.CreateSession( AcknowledgementMode.Transactional)) { var producer = session.CreateProducer(queue); try { producer.Send(session.CreateTextMessage("tx-1")); producer.Send(session.CreateTextMessage("tx-2")); session.Commit(); } catch { session.Rollback(); throw; } }

Commit()是一次性把当前事务内所有发送动作全部提交,只要有任何一步失败,可以Rollback()把这一批消息全部撤掉。这个机制在批量处理场景里非常有用,比如从数据库读一批订单,先发三条消息再确认,如果第三条失败,前两条也不会发给消费者。

消费者端同样可以结合事务,在回调中手动确认。常见做法是把 AutoAcknowledge 换成 Transactional,处理完业务再执行session.Commit()。这样消息只有 commit 了才会被标记为已消费,如果处理失败,回滚后消息还能重新进入队列。

5.2 用 Web 控制台验证消息真的落盘了

事务模式跑通之后,建议到管理台页面的Queues里做一次目测巡检。队列行会显示Number Of Pending Messages、Number Of Consumers和Number Of Messages Dequeued。发持久消息后,页面里能看到 pending 数量增加;消费者收完,dequeued 数量上来了才算完整链路。

想验证“落盘”这个动作,更直观的方法是看 KahaDB 目录下的文件大小变化。发送几十条持久消息之后,kahadb目录里的日志文件大小会明显增长,然后消费者清空队列,存储文件会回落。如果文件大小始终是零,大概率持久化适配器没生效。

如果是 Linux 环境,直接执行:

du -sh /opt/activemq/data/kahadb/

消息发之前记录一次大小,发完再看一次大小,这个对比比任何日志都可靠。

5.3 我现在的固定验证流程

现在每拿到一个新的 ActiveMQ 环境,我都会强制走一遍同样的验证流程:先把 broker 启动,用直连串写一个最小生产者发一条持久消息,立刻重启 broker 再启动消费者,确认队列消息还在;接着把连接串换成 failover,模拟断开网络,观察重连是否恢复;最后才改事务会话和持久订阅。这三步跑完,一份 Demo 才算真正吃透了。从那以后,我再也不敢直接拿着 demo 里的连接串往生产环境贴,每次都要先验证持久化和重连这两个最关键的环节。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询