
简介一份基于C#的ActiveMQ消息中间件Demo面向需要在.NET环境中快速上手消息队列的开发者。项目完整演示了NMS客户端的连接配置、会话创建、生产者/消费者收发消息并覆盖点对点与发布/订阅两种消息模型同时包含消息持久化、高可用等关键特性的实践场景。压缩包共194个文件大小12.71MB以cs源码、dll依赖、xml配置、exe可执行程序等为主辅以resx资源文件和sln工程文件适合对照源码理解ActiveMQ集成细节。已有313人学习资源内包含带图形界面的Form配置可直观设置服务器地址、端口及消息收发参数便于测试和调试ActiveMQ环境。通过该Demo开发者既能掌握在.NET中使用ActiveMQ的基本流程也能了解消息中间件的核心原理为实际项目中的异步通信开发提供参考。 ActiveMQ 这名字在 .NET 圈子里一出现不少人的第一反应是“这玩意不是 Java 生态的吗跟我 C# 有什么关系”。说实话我最早也是这么想的直到在一个上位机项目里被逼着在 C# 客户端和 Java 后端之间倒腾消息才正儿八经把 ActiveMQ 和 C# 的搭配摸了一遍。这篇就把我搭建 ActiveMQ Demo (C#) 的完整过程、踩过的坑、调参经验一次性写清楚给准备在 .NET 环境里接 ActiveMQ 的朋友一条能直接走通的路。这个 Demo 能解决什么问题简单说就是让你在 C# 程序里通过 ActiveMQ 完成消息的发布和订阅实现不同系统之间的异步通信、解耦和削峰。无论你是做上位机需要把设备状态推给后端还是业务系统里要做消息通知、任务分发这套东西都能用得上。适合刚接触消息队列的 .NET 开发者也适合从 Java 转过来想快速上手 C# 客户端的同学。1. 在 C# 项目里引入 ActiveMQ我为什么这么选1.1 消息队列能解决什么实际问题先说个最直观的场景。我以前做的一个产线数据采集上位机设备 PLC 每秒钟产生几百条数据上位机要实时显示还要同步写数据库偶尔还要通知后端系统做联动。如果全部用同步接口硬调上位机界面会卡死数据库压力也扛不住后端服务一重启数据就丢了。消息队列在这个架构里就是中间的那层“缓冲带”。上位机把数据作为消息发到 ActiveMQ后端服务按自己的节奏消费数据库写入可以批量做处理不过来就先堆在队列里服务重启之后还能继续消费不丢数据。C# 这边只管把消息发出去不用关心接收方到底是谁、什么时候处理完这就是解耦。另一个典型场景是系统集成。很多公司的主业务系统是 Java 写的但你负责的模块是 C# 的两边要做数据同步。与其两边各写一套 WebService 互相调不如统一接 ActiveMQJava 端发消息C# 端消费协议统一、接口稳定出问题也好排查。1.2 为什么在 C# 端选择 ActiveMQ 而不是别的队列现在消息队列选择很多RabbitMQ、Kafka、RocketMQ 都有各自的拥趸。我选 ActiveMQ 主要看中三点。第一是协议兼容性好。ActiveMQ 原生支持 OpenWire、AMQP、STOMP、MQTT 多种协议C# 客户端可以用 Apache.NMSNative Messaging Service走 OpenWire 协议也可以用 AMQP 协议接入。这意味着你手里的 C# 代码既可以直接对接 ActiveMQ也可以平滑迁移到其他支持 AMQP 的 Broker比如 RabbitMQ选型余地大。第二是部署轻量。ActiveMQ 是 Java 写的只要机器上有 JRE 就能跑不像 Kafka 那样要依赖 ZooKeeper单机启动一个 Broker 就能用非常适合中小型项目和个人 Demo。我本地开发直接解压运行即可配置也好改。第三是功能够用。消息持久化、事务、死信队列、延迟投递、通配符订阅这些高频功能它都有做生产级应用也不至于捉襟见肘。2. Demo 环境搭建从零把 Broker 跑起来2.1 下载安装与启动 ActiveMQ我用的版本是 ActiveMQ Classic 5.x也就是原来的 ActiveMQ 5.x 系列这个版本稳定、文档多、网上踩坑资料也全适合做 Demo。ActiveMQ Artemis 是后来的新一代 Broker性能更强但 NMS 客户端对 Classic 的支持更成熟所以这里先以 Classic 为例。到 Apache 官网下载压缩包后解压到本地目录比如D:\activemq。进入bin目录里面有两个启动脚本Windows 下用activemq.bat startLinux 下用./activemq start。启动成功后默认监听 61616 端口OpenWire 协议Web 管理控制台跑在 8161 端口。提示本地开发时如果 61616 端口被占用可以修改conf/activemq.xml里的 transportConnector 端口配置。我习惯把端口改成 61617 避免和别的服务冲突。启动后访问http://localhost:8161默认账号密码是admin/admin登录后能在 Queues 和 Topics 页面看到当前 Broker 上的消息情况配合 JMX 还能看更细的指标。这个管理控制台在后面对比消费情况、排查消息堆积时非常有用。2.2 C# 客户端库选型NMS 还是 AMQPC# 接 ActiveMQ 有两条主流路线。第一条是使用 Apache.NMS.ActiveMQ这是官方提供的 .NET 客户端走的是 OpenWire 协议。优点是 API 设计贴近 JMSJava Message Service规范你在 Java 里怎么写 ActiveMQ在 C# 里几乎可以一一对应社区资料也最多。缺点是这个库更新不太频繁但胜在稳定用于生产问题不大。第二条是使用 AMQP 协议的客户端比如 RabbitMQ.Client 或 AMQPNetLite 连接 ActiveMQ 的 AMQP 端口默认 5672。优点是如果你的系统后期要迁移到 RabbitMQ 等支持 AMQP 的 Broker代码改动小缺点是 ActiveMQ 对 AMQP 的支持在某些细节上并不完美比如消息属性和头部的处理需要注意。我这个 Demo 用的是 Apache.NMS.ActiveMQ因为它和 ActiveMQ 配合最顺调试问题也最容易找到参考。用 NuGet 安装Install-Package Apache.NMS.ActiveMQ如果是 .NET Core / .NET 5 项目同样可以用这个包我实测在 .NET 8 下可以正常运行底层依赖是 .NET Standard 2.0兼容性没问题。3. 核心代码实现从生产者到消费者完整跑通3.1 建立连接与创建会话的基本套路ActiveMQ 的 C# 客户端用法非常固定核心对象就五个连接工厂ConnectionFactory、连接IConnection、会话ISession、目的地IDestination、消息生产者/消费者IMessageProducer/IMessageConsumer。先看最基础的连接建立代码using Apache.NMS; using Apache.NMS.ActiveMQ; var connectionFactory new ConnectionFactory(tcp://localhost:61616); using var connection connectionFactory.CreateConnection(); connection.Start(); using var session connection.CreateSession(AcknowledgementMode.AutoAcknowledge);这段代码干了三件事建立到 Broker 的 TCP 连接、启动连接、创建一个自动确认模式的会话。AcknowledgementMode.AutoAcknowledge表示消息消费成功后自动确认适合大多数场景。这里有个细节很多人第一次会忽略connection.Start()一定要调否则消费者连接上了但收不到消息。我一开始就是忘了这行调试了半天发现消息发得出去但收不到后来翻文档才发现 NMS 的设计里连接默认是不启动的必须手动 Start。目的地Destination分两种Queue点对点和 Topic发布订阅。队列消息一个消息只会被一个消费者消费主题消息会被所有订阅者各收到一份。创建方式是对应的// 队列 var queue session.GetQueue(demo.queue); // 主题 var topic session.GetTopic(demo.topic);3.2 生产者发送消息同步与异步的取舍生产者的代码很直接创建一个生产者然后发送文本消息using var producer session.CreateProducer(queue); producer.DeliveryMode MsgDeliveryMode.Persistent; var textMessage producer.CreateTextMessage(Hello ActiveMQ from C#); producer.Send(textMessage);DeliveryMode有两个取值Persistent和NonPersistent。Persistent 模式消息会持久化到 Broker就算 Broker 重启也不会丢代价是性能稍差NonPersistent 模式消息只存在内存里宕机就丢但吞吐量高。我的经验是业务数据必须用 Persistent临时通知类的内容可以用 NonPersistent性能差一个量级。发送消息还可以设置消息属性、延迟时间等比如设置延迟投递var scheduledMessage producer.CreateTextMessage(延迟消息); scheduledMessage.Properties[AMQ_SCHEDULED_DELAY] 5000; // 延迟5秒 producer.Send(scheduledMessage);这个功能在做定时任务分发时很好用但要注意schedulerSupport需要在activemq.xml中开启才能生效否则消息会立即投递。3.3 消费者接收消息同步和异步两种模式消费者的写法有两种。一种是同步阻塞方式适合简单 Demo直接调用Receive()函数没消息就卡住等using var consumer session.CreateConsumer(queue); var message consumer.Receive(TimeSpan.FromSeconds(5)); if (message is ITextMessage textMessage) { Console.WriteLine($收到消息: {textMessage.Text}); }另一种是异步监听方式生产环境推荐用这个不阻塞主线程var consumer session.CreateConsumer(queue); consumer.Listener message { if (message is ITextMessage textMessage) { Console.WriteLine($收到消息: {textMessage.Text}); } };异步监听模式要注意线程安全问题。Listener 回调是在 NMS 的线程池里执行的如果你的消费逻辑里要操作 UI 控件需要 Invoke 到 UI 线程否则 WinForms / WPF 里会直接抛异常。对于 Topic 的消费还有个持久化订阅的问题。如果消费者先退出再重新订阅 Topic普通订阅会错过离线期间的消息。想要不丢需要创建持久化订阅var durableConsumer session.CreateDurableConsumer(topic, client-id-001, null, false);同时连接工厂这边要设置clientIdvar connectionFactory new ConnectionFactory(tcp://localhost:61616); connectionFactory.ClientId client-id-001;3.4 完整 Demo 的代码组织一个完整的 Demo我通常会写成一个简单的控制台应用结构如下Program.cs入口启动时根据参数决定是生产者还是消费者Producer.cs封装消息发送逻辑Consumer.cs封装消息接收逻辑appsettings.json配置 Broker 地址、队列/主题名称// Program.cs Console.WriteLine(请输入角色: 1-生产者 2-消费者); var role Console.ReadLine(); if (role 1) { var producer new Producer(tcp://localhost:61616); while (true) { Console.WriteLine(输入消息内容(exit退出):); var input Console.ReadLine(); if (input exit) break; producer.SendText(demo.queue, input); } } else { var consumer new Consumer(tcp://localhost:61616); consumer.ConsumeQueue(demo.queue, msg Console.WriteLine($收到: {msg})); Console.ReadLine(); }这样组织的好处是代码清晰后续想扩展成 WinForms 或者 WPF 上位机版本只需要把控制台输入换成界面按钮就行核心的消息收发逻辑不用大改。4. 进阶功能与性能调优Demo 之外的实战细节4.1 与 Spring Boot 项目协同跨语言消息互通热搜词里出现了“springboot整合activemq”说明不少人的实际场景是 Java 后端和 C# 客户端共存。我在一个项目里就遇到过Java 后端负责业务逻辑C# 做的 WinForms 上位机负责展示和操作两边通过 ActiveMQ 通信。Java 端用 Spring Boot 整合 ActiveMQ 发布消息到demo.queueC# 消费者这边不需要做任何特殊处理因为发送到 Queue 的消息对协议层来说是透明的。反过来 C# 生产者发消息Java 端用JmsListener就能订阅。唯一要注意的是消息格式约定比如统一用 TextMessageJSON 字符串或者用 BytesMessage 传二进制数据两端约定好属性名和序列化规则。跨语言协作时最容易出的问题就是消息体类型不匹配。Spring Boot 默认用 SimpleMessageConverter 把对象序列化成字节流这个格式 C# 端直接当 TextMessage 读会读到乱码。解决办法是 Java 端生产消息时用convertAndSend(String)或者自己控制消息体为纯 JSON 字符串C# 端也只处理 TextMessage两边都别玩花活。4.2 JMX 监控与 TopicSubscriptionViewMBean 的用法热搜词里“activemq jmx 查询 topicsubscriptionviewmbean”也是个高频问题这属于线上排查的进阶操作。ActiveMQ 启动时默认开启了 JMXJava Management Extensions可以通过 JConsole 或者 JMX 客户端查看 Broker 内部状态。TopicSubscriptionViewMBean这个 MBean 在org.apache.activemq:typeBroker,brokerNamelocalhost,destinationTypeTopic,destinationNamedemo.topic,endpointConsumer,clientIdxxx,consumerIdxxx路径下。通过它你可以查到某个 Topic 的订阅者状况是否有活跃订阅者、消息积压了多少、消费速率是多少。在 C# 里查询 JMX 需要额外引入System.Management或者通过连接到 JMX 的 RMI 端口默认 1099用 javax.management API 的桥接方式。实操中我更推荐直接用 JConsole 图形化看省时省力在conf/activemq.xml中确认useJmxtrue默认开启运行 JConsole连接到本地 Java 进程即 activemq 进程进入 MBean 面板展开org.apache.activemq节点找到对应的 TopicSubscriptionViewMBean查看PendingQueueSize、DequeuedCount、DispatchedCount等属性这个操作在排查“消息发出去但消费者没收到”这类问题时非常关键。你可以一眼看出消息是卡在 Broker 上没有被消费还是已经被消费了但显示端没刷新。4.3 几个关键调优参数性能调优在 Demo 阶段用不上但提前了解一下能避免上线后手足无措。C# 客户端这边有几个参数值得关注。PrefetchSize预取大小决定了消费者一次从 Broker 拉多少消息到本地缓存。默认值是 1000对于高吞吐场景很友好但如果你的消费逻辑比较重一次性拉太多反而容易造成消息堆积在客户端内存里。可以设置更小的值var prefetchPolicy new ConnectionPrefetchPolicy { QueuePrefetch 100, TopicPrefetch 50 }; connectionFactory.PrefetchPolicy prefetchPolicy;DispatchAsync异步分发设置是否异步分发消息给消费者。默认是 false在部分场景下改成 true 能提升吞吐量但消息到达的时序性会略受影响。消费者这边的ReceiveTimeout也很重要。用consumer.Receive()不带参数会无限期阻塞程序退出时可能出现线程挂住的情况。我一般给个 3-5 秒的超时配合循环判断退出标志。4.4 事务会话与消息确认机制如果你追求“不丢消息”事务会话是必须掌握的。事务的意义在于一批消息要么全部成功要么全部回滚。using var session connection.CreateSession(AcknowledgementMode.Transactional); using var producer session.CreateProducer(queue); try { producer.Send(producer.CreateTextMessage(第一条)); producer.Send(producer.CreateTextMessage(第二条)); session.Commit(); // 提交事务 } catch { session.Rollback(); // 出错回滚消息不会投递 throw; }对于消费者事务会话表示在Commit()之前你消费到的消息不会被确认如果消费者崩溃这些消息会被重新投递给其他消费者。这在实现“至少一次”投递语义时很实用。不过事务会显著降低吞吐量而且要注意事务的作用范围是整个 Session不是单条消息。如果一个 Session 上的消费者处理某条消息报错了回滚之后这条消息会再次被投递可能造成死循环代码里要做好防御。5. 高频问题排查我把踩过的坑都列在这5.1 消息发出去但消费者收不到步骤一先看管理控制台的 Queues 页面确认消息是否已经进入队列。如果队列里消息数为 0说明生产者的消息根本没有到达 Broker检查连接地址和网络。如果队列里有消息但消费者收不到说明消费者注册有问题。常见原因有三种。一是忘记调connection.Start()前面提到过。二是消费者在生产者发送之前还没有注册到 BrokerQueue 消息不会补发给不在线的消费者Topic 则更严格非持久化订阅者只能收到订阅期间发送的消息。三是消费者线程被阻塞了比如 UI 线程里同步调用Receive()导致消息处理不过来。5.2 消费者重复消费消息如果消费者处理消息后异常崩溃或者事务回滚消息会被重新投递这是消息队列的常见特性。确认机制没设置好也会导致重复消费。比如AcknowledgementMode.ClientAcknowledge模式下你必须手动调用message.Acknowledge()否则消息永远不会确认重启后会重新投递。处理重复消费的最佳实践是在消费端做幂等记录消息 ID处理过就不再处理。ActiveMQ 的消息带有NMSCorrelationID和NMSMessageId可以把NMSMessageId存到本地缓存或者数据库里做去重。5.3 连接不稳定或经常断开首先检查网络和防火墙61616 端口是否开放。其次 NMS 客户端默认有心跳KeepAlive机制避免 Broker 在长时间没有消息交互时误判连接失效。如果连接被防火墙切断可以通过设置连接的KeepAliveInterval和RequestTimeout来优化var connectionFactory new ConnectionFactory(tcp://localhost:61616?keepAlivetruekeepAliveInterval30000);生产环境我还会在程序里加一个重连机制用 while 循环包裹连接建立逻辑检测到ConnectionClosedException时自动重连并让消费者重新注册监听避免服务因一次性网络抖动就掉线。5.4 消息堆积但不消费消息堆积最常见的原因是消费者处理速度跟不上生产速度。用管理控制台或者 JMX 看消费者是否在线、是否有大量消息处于 Pending 状态。解决办法是增加消费者数量并发消费或者优化单条消息的处理逻辑。如果是消费者数量固定但消息仍然堆积检查有没有消费者线程挂死。在 C# 里如果消费回调抛出未处理异常NMS 的线程可能直接崩掉外部看起来连接还在但消费停了。所以回调里必须写 try-catch把异常记录下来再继续。5.5 序列化导致的接收乱码C# 端发送字符串消息Java 端收到乱码反过来也一样。遇到这类问题先确认两边都用的 TextMessage UTF-8 编码字符串。ActiveMQ 对字符串消息的默认编码是基于平台编码的可能会在跨平台时踩坑。稳妥做法是在发送方显式指定 UTF-8var textMessage producer.CreateTextMessage(input); textMessage.Properties[encoding] UTF-8;接收方看到这个属性就按 UTF-8 解码。如果是二进制消息最好直接在消息里附带格式说明、版本号字段方便兼容演进。6. 尾声一个小技巧想起来最开始做这个 Demo 时还发生过一个有趣的问题WinForms 上位机里消息回调需要更新 UI我直接在 Listener 事件里写了this.Text msg结果程序狂闪退。后来才知道 NMS 的消息线程和 UI 线程不是一回事得用Invoke或者SynchronizationContext切线程。这个问题后来成了一个固定笔记接 ActiveMQ 的上位机项目消息回调里所有操作 UI 的地方统一走一个TaskScheduler.FromCurrentSynchronizationContext()或者Control.Invoke千万别图省事直接写界面调用。还有一次消息队列里出现了一堆死信消息原因是消费者处理消息时抛异常事务不断回滚重投递最后超过最大重投次数进入死信队列ActiveMQ 默认的死信队列前缀是ActiveMQ.DLQ。后来我给 Demo 加了一个异常策略消息处理失败时先检查重投次数如果超过 3 次就手动发送到单独的error.queue由人工事后处理而不是让系统无限重试。这些都是文档里不会主动告诉你但一上线就会让你加班的细节。希望这篇基于 C# 的 ActiveMQ Demo 实战记录能让你少走几步弯路。如果你也在自己的项目里摸索消息队列欢迎多试多踩只有实际跑起来才能感受到消息队列带来的架构红利。本文还有配套的精品资源点击获取