C#连接ActiveMQ实战:消息队列生产者与消费者Demo详解

发布时间:2026/9/8 12:10:45
C#连接ActiveMQ实战:消息队列生产者与消费者Demo详解 简介这是一份基于 C# 与 WinForm 开发的 ActiveMQ 消息队列 Demo 程序包主要面向刚开始接触消息中间件的 .NET 开发者也适合需要快速评估 ActiveMQ 集成方案的工程师参考。压缩包内共 36 个文件包含 19 个 .cs 源码文件、4 个 .resx 界面资源文件、4 个依赖 DLL、3 个可直接运行的 EXE、2 个工程文件、1 个解决方案文件以及若干 settings 配置文件整体仅 326KB结构清晰可导入 Visual Studio 后直接编译运行也可直接启动 EXE 体验发送与接收效果。程序内置发送端与接收端完整示例通过 WinForm 窗体演示消息生产者与消费者之间的交互覆盖连接、会话、消息目的地、收发处理等关键环节代码注释和模块划分便于跟随调试理解 ActiveMQ 在 C# 中的基本用法和常见实现方式。已有 1290 人学习/下载轻量实用对刚入门消息队列的开发者是一份不错的参考资源。 做了这么多年.NET开发我一直觉得消息队列是C#工程师技能树上比较容易忽略、但一旦用上就回不去的一块。特别是你负责上位机、工业控制系统或者企业内部系统集成的时候经常要面对“A程序告诉B程序某个设备状态变了”“C服务要把数据异步发给D服务处理”这种跨进程、跨网络的通信需求。有人用Socket裸写有人用数据库轮询这些方案不是不能用只是到了复杂场景自己维护连接、重试、缓冲、多客户端协调代码会越来越失控。ActiveMQ是我用过的最适合入门消息队列的服务端之一而且它虽然是Java生态出身C#客户端一样连接得很顺畅。这篇文章我会带你从零跑通一个C#连接ActiveMQ的Demo先讲清楚设计思路和核心概念再给出完整的生产者、消费者代码最后把我实际开发中踩过的一些坑整理出来。适合刚接触消息队列、或者想在C#项目里引入异步通信机制的朋友参考代码可以直接“抄作业”但更重要的是理解它背后的工作原理。1. 整体设计与思路拆解1.1 为什么选ActiveMQ而不是其他消息中间件很多初学者上来就问现在不是有RabbitMQ、Kafka吗为什么还要学ActiveMQ我的看法是ActiveMQ在中小型项目里依然是性价比很高的选择。首先它对C#的支持非常成熟有官方维护的NMS客户端库不像Kafka的.NET客户端那么多兼容性坑。其次ActiveMQ部署成本极低下载解压、启动服务就完事不依赖Erlang这类额外的运行时环境对开发环境和生产环境都很友好。消息队列解决的核心问题是“解耦”和“异步”。举个例子你的上位机读到了设备数据传统做法是直接把数据写入数据库或者通过TCP推给客户端但如果数据库突然变慢、客户端断开了你的数据就丢失或者主线程被阻塞。引入ActiveMQ之后上位机只管把消息丢进队列消费者程序自己去处理后续的逻辑两边互不依赖。哪怕消费者不在线消息也会持久化保存等它上线了再取走这比手动维护Socket连接稳健得多。1.2 理解NMS与OpenWire协议C#要连接ActiveMQ靠的是Apache.NMS.NET Message Service这套类库它对应的Java端叫JMSJava Message Service。NMS是ActiveMQ官方专门给.NET平台做的接口抽象提供统一的Connection、Session、MessageProducer、MessageConsumer等API。也就是说你只要掌握NMS一套写法不管是连接ActiveMQ还是以后换别的兼容NMS的MQ代码迁移成本都比较低。底层传输协议方面默认使用的是OpenWire。这是一个ActiveMQ自有的二进制协议优点是紧凑、性能高适合纯粹的ActiveMQ与ActiveMQ客户端之间通信。如果你有跨语言的异构系统也可以改用AMQP、STOMP、MQTT等协议NMS都支持只是需要在连接字符串里指定transport connector。Demo阶段我们就用默认配置因为安装包自带的连接器就是OpenWire的tcp连接无需额外配置。1.3 队列与主题两种消息模型的选择ActiveMQ支持两种基础消息模型点对点队列Queue和发布订阅主题Topic。这个选择直接影响你的业务设计我简单说下区别。队列模式下一条消息只会被一个消费者消费。生产者和消费者不需要同时在线消费者上线后再取消息即可。适用于任务分发、日志处理这类“每人干一次”的场景。主题模式下一条消息会被所有订阅者收到但普通订阅者下线后会漏掉离线期间的消息除非配置持久化订阅。这适用于广播通知、实时告警这类“所有人必须马上知道”的场景。我们Demo会以队列为主因为队列是最常用、逻辑也最好理解等你掌握了队列主题只是改一行目标地址的事。2. 环境准备与基础配置2.1 下载安装与启动ActiveMQ服务端先去ActiveMQ官网下载最新的稳定版Windows用户直接选Windows发行包zip格式Linux选tar.gz。下载完成后解压到一个路径全程无中文、无空格的目录下比如D:\apache-activemq-x.x.x避免一些奇怪的乱码问题。启动非常简单Windows下进入bin目录双击activemq.bat或者在命令行执行activemq start。启动成功后默认监听61616端口OpenWire和8161端口Web控制台。在浏览器访问http://localhost:8161出现ActiveMQ管理页面就说明服务端没问题了。默认账号密码是admin/admin建议首次登录后就去conf/jetty-realm.properties文件里改掉毕竟Web控制台能直接浏览、操作队列里的消息裸奔在开发机还好在生产环境非常危险。2.2 创建C#项目和引入NMS依赖我用的是.NET 8创建一个控制台项目即可。这里说明一下生产者和消费者可以放在同一个项目里用两个入口运行也可以建两个项目分别运行后者更贴近实际部署环境我下面的代码也是按两个独立控制台项目来组织的。在Visual Studio里右键项目选择“管理NuGet程序包”搜Apache.NMS.ActiveMQ安装最新稳定版即可。这个包会自动依赖Apache.NMS核心库不需要手动再装别的。我用的是1.8.0版本接口稳定资料也好查。如果你用的是.NET Framework老项目同样能装这个包它支持多个目标框架。2.3 连接字符串解析NMS的连接字符串格式跟大家熟悉的数据库连接字符串很像tcp://localhost:61616如果要带用户名密码可以在创建连接时传入也可以在字符串里拼参数。我实际用得比较多的写法是这样tcp://localhost:61616?wireFormatmaxInactivityDuration30000maxInactivityDuration这个参数值得注意它是心跳保活机制的一部分。如果你的程序长时间没有消息收发超过这个时间连接可能会被服务端断开加了参数可以调整超时阈值。开发环境不设也行但生产环境建议显式配置避免防火墙空闲超时把TCP连接干掉。3. 核心代码实现与参数解析3.1 生产者发送消息的完整实现先上代码生产者的核心逻辑就是建立连接、打开会话、声明队列、发送消息。我用一个完整的控制台程序来演示using System; using Apache.NMS; using Apache.NMS.ActiveMQ; class ProducerDemo { static void Main(string[] args) { // 1. 创建连接工厂 IConnectionFactory factory new ConnectionFactory(tcp://localhost:61616); // 2. 建立连接并启动 using (IConnection connection factory.CreateConnection(admin, admin)) { connection.Start(); // 3. 创建会话false表示非事务模式AcknowledgeMode.AutoAcknowledge表示自动确认 using (ISession session connection.CreateSession(AcknowledgementMode.AutoAcknowledge)) { // 4. 声明一个队列队列不存在时会自动创建 IDestination destination SessionUtil.GetDestination(session, queue://demo.queue); // 5. 创建生产者 using (IMessageProducer producer session.CreateProducer(destination)) { // 默认持久化确保服务端重启后消息不丢失 producer.DeliveryMode MsgDeliveryMode.Persistent; for (int i 1; i 10; i) { ITextMessage message session.CreateTextMessage($Hello ActiveMQ, message {i}); producer.Send(message); Console.WriteLine($Sent: {message.Text}); } } } } Console.WriteLine(Producer done. Press any key to exit.); Console.ReadKey(); } }这段代码里有一个容易忽视的细节SessionUtil.GetDestination(session, queue://demo.queue)。queue://前缀表示明确使用队列类型如果把这个前缀换成topic://就变成主题模式了。当然也可以用session.GetQueue(demo.queue)这个更直观的写法效果一样。我习惯用前缀方式因为一个变量就能动态切换队列和主题调试不同消息模型时会方便很多。第二个值得注意的地方是MsgDeliveryMode.Persistent。ActiveMQ有两种投递模式持久化模式和非持久化模式。持久化消息会先写入磁盘日志再返回发送成功性能略低但可靠性高非持久化消息保存在内存中服务端重启就会丢失。Demo里我显式设置成持久化这是为了强调一个原则凡是重要的业务消息一律用持久化。3.2 消费者接收消息与消息确认机制消费者的实现稍微复杂一点涉及消息监听和确认机制。我先贴出最基础的自定义监听器实现using System; using Apache.NMS; using Apache.NMS.ActiveMQ; class ConsumerDemo { static void Main(string[] args) { IConnectionFactory factory new ConnectionFactory(tcp://localhost:61616); using (IConnection connection factory.CreateConnection(admin, admin)) { connection.Start(); using (ISession session connection.CreateSession(AcknowledgementMode.AutoAcknowledge)) { IDestination destination SessionUtil.GetDestination(session, queue://demo.queue); using (IMessageConsumer consumer session.CreateConsumer(destination)) { // 注册消息监听器异步处理 consumer.Listener OnMessage; Console.WriteLine(Consumer started. Press any key to exit.); Console.ReadKey(); } } } } static void OnMessage(IMessage message) { if (message is ITextMessage textMessage) { Console.WriteLine($Received: {textMessage.Text}); } } }这里最关键的是AcknowledgementMode.AutoAcknowledge。自动确认模式下一旦消息从mqtt服务端发给消费者更准确地说是消费端在receive或Listener回调返回时ActiveMQ就认为这条消息已经处理完成会从队列中移除。如果消费者程序在处理消息过程中崩溃了比如数据库操作到一半挂了消息已经确认但业务没落库这条消息就找不回来了。所以在实际项目中我强烈建议用ClientAcknowledge模式处理完业务逻辑后再手动调用message.Acknowledge()确保消息处理的可靠性。下面是一个手动确认的示例片段using (ISession session connection.CreateSession(AcknowledgementMode.ClientAcknowledge)) { using (IMessageConsumer consumer session.CreateConsumer(destination)) { consumer.Listener message { try { // 模拟业务处理比如写数据库、调接口 Console.WriteLine($Processing: {((ITextMessage)message).Text}); // 业务成功后确认消息 message.Acknowledge(); } catch (Exception ex) { // 记录日志不确认消息之后会重新投递 Console.WriteLine($Error: {ex.Message}); } }; Console.ReadKey(); } }注意手动确认模式下如果一直不调用Acknowledge()消息不会丢失客户端重启后它会以红字DLQ或者重新投递的形式再次出现。这里面又延伸出一个坑如果业务处理总是失败消息就会被无限重投形成“毒消息”严重时会挤占队列。正确的姿势是结合死信队列Dead Letter Queue策略把重试多次仍失败的消息转移到单独的队列人工排查。ActiveMQ默认就有死信队列机制处理不了的过期消息会进入ActiveMQ.DLQ。3.3 事务会话与批量发送除了消息确认机制NMS还支持事务会话。把CreateSession的第一个参数设为true就开启了事务模式。事务模式下消息不会立刻提交而是等session.Commit()时才真正发送或接收确认中途出错可以session.Rollback()回滚。我一般在批量导入、批量通知的场景下用它性能提升非常明显。这里给一个批量发送的代码思路using (ISession session connection.CreateSession(AcknowledgementMode.Transactional)) { using (IMessageProducer producer session.CreateProducer(destination)) { for (int i 0; i 1000; i) { ITextMessage message session.CreateTextMessage($Batch message {i}); producer.Send(message); // 每200条提交一次减少事务开销 if (i % 200 0) { session.Commit(); } } session.Commit(); // 最后一批提交 } }分批提交而不是每条都Commit是为了减少磁盘同步和网络往返次数。这个技巧在日志批量上报、数据同步这类高吞吐场景里很常见。需要留意的是事务会持有一段时间的资源批量太大、事务太久服务端的内存占用会明显上升你可以通过调整提交频率找到适合自己业务的平衡点。4. 常见问题与排查技巧实录4.1 连接被拒绝先检查端口和防火墙第一个常见问题就是代码跑起来报SocketException: No connection could be made because the target machine actively refused it。这个错误90%是服务端没启动或者端口不通。先确认61616端口是否在监听Windows上可以用netstat -ano | findstr 61616Linux用ss -lntp | grep 61616。如果端口正常监听再检查本机防火墙是否放行了该端口尤其在公司内网环境防火墙策略经常会拦截到开发机的连接。还有一个容易被忽略的点连接字符串里如果主机名写错或者ActiveMQ服务端配置了安全认证插件也会导致连接失败。建议先把用户名密码去掉测试纯匿名连接能否成功如果能连再逐步加上认证这样能快速定位问题。4.2 消息反序列化失败或类型不匹配C#客户端发消息Java客户端消费或者反过来这是异构系统集成时常遇到的事。如果你在C#端发送的是BytesMessage拿Java端去消费然后用TextMessage强转一定会报类型转换异常。消息内容是跨语言传递时我建议统一用纯文本格式JSON字符串发送方用CreateTextMessage(jsonString)接收方拿到文本后再反序列化成自己的业务对象。这个方法简单粗暴但也是最通用的。千万不要在C#端直接用ObjectMessage序列化一个C#自定义类再发给Java端两个语言的序列化规则完全不同基本必坑。4.3 消费端收不到消息的排查思路如果生产者发送成功控制台显示Sent但消费者就是收不到按照下面顺序排查第一检查目标地址。看消费者订阅的队列名和生产者发送的队列名是否完全一致注意大小写。Demo.Queue和demo.queue在ActiveMQ里是两个不同的队列我之前就因为大小写不统一查了半小时。第二检查消费者是否在生产者发送之前就已经启动。队列模型下消费者不在线不影响消息积压消费者上线后会继续消费但如果你用的是Topic模型普通消费者错过发布时段就收不到了。Demo阶段用队列模式能避免这个坑。第三确认Session和Connection调用了Start()。IConnection.Start()这行代码很容易漏掉漏了以后代码不报任何异常但消息就是静默收不到。这个看似小白的错误实际工作里我见过不止一次。4.4 长时间运行的应用会话与连接资源管理如果你把消费程序做成Windows服务或者后台常驻进程还需要注意一个问题——资源释放。NMS的Connection和Session都实现了IDisposable但它们不是无脑用完就释放的。消费者程序一般会长时间运行一条连接、一个会话可以持续复用不要把using写在消息处理回调里。每次处理消息时重新创建Session、Consumer会造成大量不必要的握手和资源分配在高并发下很容易把服务端连接数打满。正确的做法是启动时创建Connection和Session注册Listener然后程序进入等待状态程序停止时再统一释放资源。生产者的做法则相反短连接、即用即走更合适因为生产者通常是突发性的维护长连接的成本可能大于连接创建的成本。4.5 生产环境必须关注的几个参数最后把我的生产环境配置清单整理出来这些参数在开发环境不需要管但上了生产环境影响很大。参数位置建议值作用maxInactivityDuration连接字符串30000左右调整心跳超时阈值防止空闲断连MemoryLimit服务端conf/activemq.xml根据机器内存设置限制单个broker的消息内存上限持久化适配器服务端conf/activemq.xmlJDBC或KahaDB决定消息存储方式KahaDB适合中小规模JDBC适合要跟业务库联动的大型系统预取大小prefetchSize连接字符串默认1000控制消费者预取消息数量每条消息处理代价大时建议调低避免大量消息堆积在客户端内存提一个跟C#消费端强相关的点prefetchSize默认值会影响确认机制的效果。如果设置成1000消费者本地会预取1000条消息此时即使你在代码里每条都调用Acknowledge()服务端也不会知道你处理得这么频繁它只会认为你批量确认了。对付“每条消息都必须严格确认”的场景比如金额处理建议在连接字符串里加?transport.prefetchSize1相当于每次只取一条处理完确认后再取下一条。牺牲一些吞吐量换来精确的逐条控制代价完全值得。在做Demo或者初学阶段我最推荐的路径是先跑通25行代码的生产者和消费者理解连接、会话、目的地、生产、消费这五个核心动作然后改变消息模型、改变确认模式对比观察行为差别最后再引入事务和持久化配置。我在实际项目中把ActiveMQ用在上位机数据采集、报警通知、工单分发这些场景稳定跑了几年相比最初用TCP直连的方案代码复杂度和故障率都明显下降。如果你正在评估项目里的通信架构或者刚接触消息队列不妨先把这篇文章里的Demo跑起来跑通你就会明白消息队列并不神秘关键在于理解几个核心概念后亲手实践。本文还有配套的精品资源点击获取