资讯详情

ActiveMQ C# Demo实战:从Queue到持久化消息的避坑指南

📅 2026/10/11 5:08:49 | 华诺云谱 👁 阅读
ActiveMQ C# Demo实战:从Queue到持久化消息的避坑指南
简介面向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 端口 61616Web 管理台端口 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 startLinux 下执行./activemq start。启动成功之后进程不会立刻退出而是在后台跑一个 broker 实例。为了确认它确实起来了我会顺手做两件事看端口占用再看管理台。cd /opt/activemq/bin ./activemq start netstat -ano | grep 61616 curl http://127.0.0.1:816161616是 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)?startupMaxReconnectAttempts10直连适合调试。tcp://后面跟 IP 和端口broker 不在了就立刻抛异常错误信息干净。failover则是把多个 broker 地址用括号包起来客户端在断线之后会自动尝试重连默认会一直重连所以生产环境普遍用它。这中间有个容易被忽略的细节failover不是 ActiveMQ 特有的魔法它的底层是客户端心跳和通道重建。Demo 里如果只有直连配置说明它只给你演示最小链路不是为了让你直接上生产的。你后续要改的往往是failover:(tcp://host1:61616,tcp://host2:61616)这种高可用写法而不是去改一堆timeout参数。连接串上的参数也属于“少配一个就出大事”的类型。常见的几个我整理在这里参数含义建议startupMaxReconnectAttempts启动时最大重连次数调试设 3 左右生产不要设太小timeoutfailover 重试间隔默认看着合适网络抖动时可调 3000keepAlive是否开启 TCP 心跳长时间空闲连接建议为 truetransport.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 与 TopicDemo 里为什么两套代码长得一样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 nameopenwire uritcp://0.0.0.0:61616?keepAlivetrue/同时C# 端连接串写成failover:(tcp://127.0.0.1:61616)?keepAlivetrue这样即使业务消息稀疏底层 TCP 也会定期发送心跳包避免连接被 broker 误杀。注意keepAlivetrue只是发送心跳不会改变业务消息的投递语义。4.3 消费回调里抛异常消息被反复投递最后刷屏现象Listener回调里处理业务时抛了一个未捕获异常随后日志里出现大量相同消息的接收记录消费者像卡死一样不停处理同一条消息。原因Listener回调在消息没有确认成功的情况下broker 默认会重新投递该消息同时重投次数不受限制。异常抛出来之后确认根本没发生消息就一直在“投递、失败、再投递”的循环里。解决在回调入口就捕获所有异常并限制重投次数var redeliveryPolicy new RedeliveryPolicy { MaximumRedeliveries 3, InitialRedeliveryDelay 2000, UseExponentialBackOff true }; factory.RedeliveryPolicy redeliveryPolicy;然后消费回调里加上 try/catchconsumer.Listener message { try { HandleMessage(message); } catch (Exception ex) { Console.WriteLine(ex.Message); } };MaximumRedeliveries设成 3意味着同一消息最多重投 3 次第 4 次会被送进死信队列。这项配置需要小心MaximumRedeliveries0时消息一次失败就直接进死信设太大又会造成消息延迟。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 里的连接串往生产环境贴每次都要先验证持久化和重连这两个最关键的环节。希望帮到你。本文还有配套的精品资源点击获取
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑