ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

C# ActiveMQ Demo实战:从连接配置到生产者消费者完整教程

C# ActiveMQ Demo实战:从连接配置到生产者消费者完整教程 简介ActiveMQ DemoC#是一套面向消息队列初学者的可运行示例工程演示了ActiveMQ在.NET环境下的发送与接收流程。工程采用WinForm图形界面将连接配置、消息发送区域与消息列表整合在同一界面中适合C#开发人员快速上手也可用于教学演示或作为企业内部消息服务的前期验证。代码包含消息生产者与消费者两端完整实现涉及队列模式下的消息创建、发送、监听、接收与界面绑定并提供了GlobalFunction、MQ封装类等辅助模块帮助理解ActiveMQ客户端API的调用方式同时利用ListView控件展示消息记录并支持列排序便于观察消息投递与消费结果。压缩包共36个文件以19个cs源码文件为主辅以resx界面资源、settings配置、sln解决方案以及编译生成的dll与exe整体仅326KB结构简洁便于对照学习。目前已有1292人学习下载。作者bodybo还撰写了配套介绍文章结合Demo阅读可更系统地掌握ActiveMQ的核心机制并在此基础上扩展出更适合自身项目的消息处理逻辑。1. 为什么 C# 项目里要自己搭一个 ActiveMQ Demo很多 C# 团队在第一次面对消息队列选型时第一反应是 RabbitMQ 或者 KafkaActiveMQ 常常被当成“老古董”跳过。但实际一查ActiveMQ 的 JMS 语义最完整、部署最简单、对 Windows 环境最友好尤其是当你的下游是政府项目、工业上位机、物流 WCS 这类不能随便上重型组件的场景时ActiveMQ 几乎是唯一能在半天内跑通、并且不引入额外运维负担的选择。这个标题要解决的就是一件事在你的 C# 工程里用最小代价把 ActiveMQ 的生产者、消费者、队列、主题全部跑通得到一个可以继续往上加业务的模板。适合读这篇文章的人有三类第一类是上位机或桌面端开发需要在设备和后端之间塞一层消息缓冲第二类是刚接触消息队列的 C# 新手想找一个能复现、能打断点、能看明白的 Demo第三类是已经在用 RabbitMQ 但被交换机绑定搞烦了想看看 JMS 风格是不是更直接。我不会给你贴一个“开箱即用”的完整工程而是把每一步的代码和参数讲清楚让你能亲手搭出来后面出问题也知道去哪里看。2. ActiveMQ 与 C# 客户端选型先搞懂消息代理、连接协议和 NMS2.1 ActiveMQ 在消息链路里到底扮演什么角色ActiveMQ 本质上是一个独立的 Java 进程它不跟你抢内存也不寄宿在 IIS 里而是以 Broker 的身份单独存在。所谓 Broker就是生产者和消费者之间的中转站生产者把消息发到 BrokerBroker 帮你存着消费者什么时候来取都行。这样设计带来的第一个好处是解耦——生产者和消费者不需要知道对方在哪里、是否在线只要认 Broker 的地址就行第二个好处是削峰——下游处理不过来的时候消息在 Broker 里排队不会直接压垮下游数据库或接口。在 C# 的语境里你要意识到 ActiveMQ 是不关心你的业务代码用什么语言写的它只认协议。默认情况下它开放了两个端口61616 是 OpenWire 协议端口给 Java、C#、C 这类原生客户端用8161 是 Web 控制台端口浏览器访问用来查看队列深度、消费者数量、消息内容。很多 C# 新手第一次连接失败就是因为把地址写成了tcp://localhost:8161这个端口是 HTTP 的不是消息协议端口我用tcp://localhost:61616才能连上。2.2 用 Apache.NMS.ActiveMQ 还是 NFactomcat.ActiveMQC# 生态里连接 ActiveMQ 有两条主流路线一条是 Apache 官方维护的 NMS.NET Message Service库另一条是国产开源的 NFactomcat.ActiveMQ。我做项目用的是官方 NMS原因有三一是它跟 ActiveMQ Broker 的版本兼容性最稳ActiveMQ 5.15、5.16、5.18 都验证过二是它的 API 风格贴近 JMS 原始语义你以后转 Java 或转 Pythonstomp.py时认知不用切换三是官方库的命名空间清晰ConnectionFactory、IConnection、ISession、IMessageProducer、IMessageConsumer看名字就知道干嘛的。NFactomcat.ActiveMQ 我也简单看过它的优点是纯 C# 实现、不依赖 IKVM某些低版本 .NET Framework 项目里好使但它对 OpenWire 协议的实现版本偏旧遇到 ActiveMQ 5.17 以上版本时偶尔出现连接后无响应的情况。我的建议是.NET 6 以上无脑选 Apache.NMS.ActiveMQ还在维护 .NET Framework 4.5 老项目的先测 NFactomcat 的握手不行再切回 NMS。提示NuGet 包名是Apache.NMS.ActiveMQ不要下错成Apache.NMS核心库前者是协议实现后者只是抽象接口两个都要装但实际连接用的是前者它会自动带上 NMS 核心库依赖。2.3 Windows 上安装 ActiveMQ Broker17 版和 5.18 版的差别如果你只是本地验证 Demo我推荐直接下载 Apache ActiveMQ 5.18.x 的 Windows 压缩包解压即用不需要安装。下载地址从 ActiveMQ 官网的 Download 页面进找到apache-activemq-5.18.x-bin.zip解压到D:\activemq这种纯英文路径。5.18 之前的版本比如 5.15、5.16默认免认证控制台直接能打开5.17 开始 Jetty 版本升级控制台的行为有变化5.18 默认也需要留意登录配置但消息端口 61616 的开放协议默认是开的不影响 C# 连接。启动方式是在解压目录里打开命令行执行cd /d D:\activemq\bin activemq start启动完成后看到控制台输出ActiveMQ WebConsole available at http://127.0.0.1:8161/就是成功了。5.18 版本登录控制台的默认用户名密码是admin/admin新版如果提示认证失败去conf/jetty-realm.properties里确认一下账号是否存在。生产者和消费者的 Demo 不依赖控制台控制台只是用来观察队列深度和消息内容。到这里Broker 已经在你机器上跑起来了。下一步我带你创建一个 C# 控制台工程用 NuGet 把 NMS 客户端拉进来然后分别写生产者和消费者。3. C# 控制台工程跑通 ActiveMQ 最小链路从 NuGet 到收发消息3.1 创建工程和 NuGet 依赖我会用一个控制台工程演示因为它最能暴露问题——没有 ASP.NET Core 的依赖注入遮挡每一个对象都是你亲手 new 出来的你才能真正理解连接的生命周期。打开 Visual Studio 或 Rider创建.NET 6的Console App项目名称建议叫ActiveMQ.Demo。创建完成后在项目文件里加包引用或者直接在命令行操作dotnet new console -n ActiveMQ.Demo cd ActiveMQ.Demo dotnet add package Apache.NMS.ActiveMQ --version 2.1.0需要说明的是2.1.0 不是最新版本但它是经过最多生产环境验证的版本ActiveMQ 5.15 到 5.18 都兼容。如果你要用最新的 2.x 版本直接dotnet add package Apache.NMS.ActiveMQ不加版本号也行但我强调一下2.1.0 足够稳定不要追新。添加完后dotnet restore然后确认obj/project.assets.json里有Apache.NMS.ActiveMQ/2.1.0的记录。3.2 写一个最简生产者连接、会话、消息、发送下面是完整的生产者代码我把它写在一个Producer.cs文件里。它做的事情是建立 TCP 连接到 61616 → 创建会话 → 创建生产者 → 往名叫demo.queue的队列发一条文本消息 → 关闭资源。using Apache.NMS; using Apache.NMS.ActiveMQ; // 1. 创建连接工厂指定 Broker 的 OpenWire 端口 var connectionFactory new ConnectionFactory(tcp://127.0.0.1:61616); // 2. 建立连接 using var connection connectionFactory.CreateConnection(); connection.Start(); // 3. 创建会话第一个参数 AcknowledgementMode第二个参数是否事务 using var session connection.CreateSession(AcknowledgementMode.AutoAcknowledge); // 4. 创建队列对象和目标生产者 var destination session.GetQueue(demo.queue); using var producer session.CreateProducer(destination); // 5. 发送一条文本消息 var message session.CreateTextMessage(Hello from C# at DateTime.Now); producer.Send(message); Console.WriteLine(Sent: message.Text);这里的核心逻辑在CreateSession的两个参数上。AcknowledgementMode.AutoAcknowledge表示消息发出去后Broker 不用等接收方确认就认为发送成功false表示不走事务每条消息独立提交。这个搭配适合 Demo 和大多数日志上报、缓存刷新场景。GetQueue(demo.queue)是关键——ActiveMQ 的 Queue 是虚拟的你不需要在 Broker 里预先创建第一次发送消息时它自动创建。这也是 JMS 和 RabbitMQ 的一个大区别RabbitMQ 的队列得先声明ActiveMQ 的队列直接发就行。如果你连的是 Topic把GetQueue换成GetTopic即可但行为完全不同后面我用一整节讲。提示connection.Start()不能省略。在 NMS 里连接默认是未启动状态只有显式调用Start()后消费者才能真正收到消息但生产者不调也能发这是很多新手困惑的“为什么我能发但不能收”的原因之一。3.3 写一个最简消费者同步接收和异步监听两种方式消费者的写法有两种第一种是同步轮询适合命令行工具、定时任务consumer.Receive()会阻塞当前线程直到有消息到达。第二种是异步监听适合 WinForms、WPF、上位机这类需要保持 UI 线程响应的场景。我先贴同步版本using Apache.NMS; using Apache.NMS.ActiveMQ; var connectionFactory new ConnectionFactory(tcp://127.0.0.1:61616); using var connection connectionFactory.CreateConnection(); connection.Start(); using var session connection.CreateSession(AcknowledgementMode.AutoAcknowledge); var destination session.GetQueue(demo.queue); using var consumer session.CreateConsumer(destination); // 同步阻塞接收等待消息 var message consumer.Receive(TimeSpan.FromSeconds(10)); if (message is ITextMessage textMessage) { Console.WriteLine($Received: {textMessage.Text}); } else { Console.WriteLine(No message within 10s.); }Receive(TimeSpan.FromSeconds(10))是带超时的接收10 秒内没有消息就返回 null。这里有个值得注意的地方consumer.Receive()不传参时会无限阻塞如果你的程序是被定时任务调起来的无限阻塞会导致第二次任务触发时线程池被占满。异步监听版也不复杂核心是注册一个Listener委托它会在消息到达时在 NMS 内部线程池里被调用using var consumer session.CreateConsumer(destination); consumer.Listener message { if (message is ITextMessage textMessage) { Console.WriteLine($Async received: {textMessage.Text}); } }; connection.Start(); Console.ReadKey(); // 保持进程存活注意异步版的connection.Start()要放在Listener注册之后虽然 NMS 内部对注册顺序没有强校验但从语义上理解先注册回调再启动连接才能保证连接启动后第一时间就能路由消息。控制台程序里不加Console.ReadKey()的话进程会立刻退出消息还没收完就没了这是控制台 Demo 最容易出现的“看起来没收到”的假象。3.4 消息确认机制AutoAcknowledge 之外的两个选择AcknowledgementMode.AutoAcknowledge是最省心的但你要知道还有ClientAcknowledge和DupsOkAcknowledge因为真正的生产环境大概率要用到它们。ClientAcknowledge表示消息从 Broker 拉下来后不算完成必须你手动调用message.Acknowledge()才算确认。这个模式解决什么问题解决“处理了一半崩了”的问题。比如你收到一条消息先写数据库、再调接口如果写数据库成功、调接口失败在 AutoAcknowledge 模式下这条消息已经算消费成功了接口调失败了也没人管。但在ClientAcknowledge下你可以把Acknowledge()放在接口调用成功之后失败就不确认消息会由 Broker 重新投递。DupsOkAcknowledge比较特殊它允许消息被重复投递但不会把未确认的消息存盘适合追求吞吐、能容忍重复消费的场景。它的语义用大白话说就是我告诉你收到了但你也可以再给我发一次我不介意。在 C# 里配置这个模式只需要改CreateSession的第一个参数using var session connection.CreateSession(AcknowledgementMode.ClientAcknowledge); // 在处理成功后手动调 textMessage.Acknowledge();我自己的经验是内部系统之间传通知用 AutoAcknowledge 就行涉及订单、计费、库存这类想对账的必须用ClientAcknowledge加业务补偿因为消息中间件的“不丢失”只保证到 Broker不保证你的业务代码执行完整。4. Queue 与 Topic 的选择消息模型取决于你的业务是“分发”还是“广播”4.1 Point-to-Point 和 Publish/Subscribe 的本质区别ActiveMQ 支持两种消息模型GetQueue对应点对点模型GetTopic对应发布订阅模型。点对点的核心是竞争消费一条消息只会被一个消费者拿走多个消费者同时监听同一个 Queue 时Broker 会把消息轮流分给它们这就是负载均衡的基础。发布订阅的核心是广播Topic 上的每一条消息会复制给所有当前在线的订阅者谁不在线谁就错过。怎么选看业务场景。如果你有一个订单创建事件下游有库存服务、积分服务、通知服务要各自处理这是典型的 Topic 场景如果你有一批待处理的图片压缩任务有 10 个工作线程并发处理这就是 Queue 场景。注意在 ActiveMQ 里 Queue 和 Topic 不是同一个东西你不能把一个消费者先 GetQueue 再 GetTopic 去收同一条消息。用 C# 创建 Topic 消费者时唯一的变化就是session.GetTopic(demo.topic)。但这里有一个隐蔽的坑Topic 消费者如果在消息发送之前还没有Start()并且CreateConsumer那发送的消息它收不到。Topic 的消息默认不持久化而且没有订阅者在线的时候消息直接丢弃。很多第一次接触 JMS 的人在这里翻车先运行了生产者再启动消费者结果消费者什么都收不到然后怀疑代码有问题。其实不是代码问题是模型问题——Topic 对“离线者”没有补偿。要补偿就得请出持久化订阅。4.2 持久化订阅让 Topic 消费者错过消息也能补收CreateDurableConsumer是 JMS 规范里的概念NMS 也支持但用法和普通消费者略有不同。它需要客户端提供一个唯一的订阅名称Broker 会把这个订阅者的状态保存下来消费者离线期间发到 Topic 的消息会被存起来下次这个消费者上线时补发。using var consumer session.CreateDurableConsumer( destination, // 必须是 Topic pay-service-subscriber, // 订阅者唯一标识 null, // 消息筛选器 false); // noLocal是否不接收本连接发布的消息这里有两个关键参数。订阅者名称必须全局唯一如果你用两台机器部署同一个服务订阅名称写死成同一个那它们会变成同一个订阅者消息只会补发给其中一个另一个永远收不到全量消息——这是个非常隐蔽的生产事故。noLocal参数设为true时消费者不会收到同一连接上生产者发到该 Topic 的消息大多数场景设false就行。持久化订阅在 Demo 里不太显眼但做上位机或服务端时很值得提前设计。比如 WCS 系统里有一个“任务下发”Topic分拣控制器离线 5 分钟重新上线后如果不做持久化订阅这 5 分钟的任务单就丢了而上位机又没有简单的方式来重新拉取任务。用持久化订阅Broker 会帮你攒着上线自动补。4.3 Waiting Message Count 和消费者数量的关系在 ActiveMQ 控制台的 Queues 页面你会看到Pending列的数字一直涨然后有消费者接入后降下来。很多 C# 开发者在这时候会困惑我开了 5 个消费者为什么Pending不下降Consumer Count却显示 5原因通常是消费者创建了但 Session 没启动或者connection.Start()没调用。控制台的Consumer Count表示有 5 个消费者注册到了 Broker 上但它们没有活跃接收消息不会流向它们Pending自然不降。这里我建议你把“注册消费者”和“启动消费”当成两件事CreateConsumer只是注册connection.Start()Receive()/Listener才是真正开始消费。另外要注意 Queue 模式下的预取Prefetch行为。默认情况下NMS 客户端每次向 Broker 批量拉取一定数量的消息放到本地缓冲区这个数量默认是 1000。即使你代码里只调用了Receive()客户端也可能已经预取了 1000 条到本地这在控制台看Pending是下降的但你还没处理完程序一崩溃这 1000 条就丢了。想减少丢消息窗口可以手动把预取数量调小var connection connectionFactory.CreateConnection() as Connection; if (connection ! null) { connection.PrefetchPolicy.QueuePrefetch 10; }这个参数在 Demo 里无所谓但一旦你要做真正的业务系统这就是“少丢几条”和“多丢几百条”的分水岭。后面说坑的时候我还会提一次。5. ActiveMQ Demo 踩坑记录连接失败、消息丢失、内存溢出这些坑我都替你趟过5.1 端口连不通8161 是 Web 控制台61616 才是消息端口现象C# 代码里ConnectionFactory写的是tcp://127.0.0.1:8161程序一跑就抛Apache.NMS.NMSConnectionException提示连接被拒绝。很多人第一反应是 ActiveMQ 没启动反复启动服务也没用。原因8161 端口绑定的是 Jetty Web 容器它只提供 HTTP 协议不认 OpenWire。OpenWire 端口默认是 61616。NMS 只实现了 OpenWire 和 STOMP 协议不能跟 HTTP 端口通信。解决把连接地址改成tcp://127.0.0.1:61616。改完还连不上的话在 Windows 上执行netstat -ano | findstr 61616看端口是否处于监听状态同时确认conf/activemq.xml里没有把 transportConnector 的端口改掉。5.18 默认配置不会动它但如果你是从老项目拷贝的配置里面可能是 61613STOMP或 61616 被注释掉了。5.2 生产者发送成功消费者收不到Topic 离线消息直接丢弃现象先运行消费者再运行生产者Queue 模式下消费者收到消息了但把GetQueue换成GetTopic后同样的代码消费者收不到任何消息。原因Queue 会保存消息直到被消费Topic 只在消费者在线时推送。如果你先启动生产者消息发出去Topic 上没有订阅者这条消息立刻蒸发。这个行为是 JMS 规范定义的不是 ActiveMQ 的 bug。解决确认业务确实需要广播模型时用持久化订阅兜底。在 Demo 阶段最简单的排查方式是按“先消费者后生产者”的顺序运行Topic 就能收到。如果还想验证消息是否真的发到了 Broker去控制台 Topics 页面看Enqueue Count是否增长只看Pending是看不出来的因为 Topic 的 Pending 只对持久化订阅者有意义前面已说明。5.3 控制台 Queue 显示有消息但 C# 消费者一启动就“全没了”现象Receive()调用后拿到一条消息但业务处理失败消息找不回来了。控制台Dequeue Count已经增加Pending清零。原因AutoAcknowledge 模式下消息只要从 Broker 发到客户端就算确认了哪怕你还没处理完。C# 程序崩溃、断电、异常抛出消息都不会重投。解决对不能丢的业务改用AcknowledgementMode.ClientAcknowledge并在业务代码成功结束后再调Acknowledge()。有一点必须提醒如果消息在Acknowledge()之前抛异常消息会重新进入队列但 ActiveMQ 默认重投次数不限这会导致一条坏消息反复消费、卡住整个队列。要限制重投次数可以在消费者上设置消息属性或使用死信队列配置——在activemq.xml里对IndividualDeadLetterStrategy做配置把重投多次的消息转到一个DLQ队列存起来。5.4 Prefetch 预取导致消息堆积在客户端内存现象控制台Pending降为 0但 C# 程序的处理速度明显低于预期甚至内存占用持续升高。业务日志里消息是一条一条在处理但总觉得客户端“偷偷多拿了很多”。原因NMS 默认的 Queue 预取策略是批量拉取到本地PrefetchPolicy.QueuePrefetch默认 1000。消费者一启动就把 1000 条消息拉到内存里Broker 认为它们已经投递出去了但你的代码才刚刚处理到第 1 条。解决根据业务调整预取数量。需要逐条精处理的任务把QueuePrefetch降到 1 或 10消息变成“拉一条、确认一条”的模式高吞吐场景不需要动但要在文档里明确标注“进程崩溃可能会丢预取消息”。注意这个参数要在connection.Start()之前设置否则不会生效。5.5 多次调用 connection.Start() 会不会出问题现象有人写完代码后在CreateConsumer前面调了Start()在Listener注册后又调了一次Start()心虚来问会不会重复消费。原因Start()不是“打开开关”它是把你已经注册的消费者都激活。多次调用没有副作用它不会重启消费者也不会导致消息重复投递。但如果你在Start()之后再创建新消费者新消费者也能收到消息——这点和后创建的消费者必须先Start()才能收到消息的经验正好相反实际上Start()之后新建的消费者是立即激活的不需要再调一次。解决把connection.Start()统一放在所有消费者创建之后。这个顺序在 Demo 里没影响但在多消费者工程里顺序乱了会导致部分消费者先收消息、部分后收日志的时间线对不起来排查问题时会困惑。6. 进阶验证与调优用控制台确认消息轨迹用事务会话做批量推送6.1 用 ActiveMQ 控制台二次确认 Demo 是否真的通了Demo 跑通后先别急着写业务打开http://127.0.0.1:8161用admin/admin登录进 Queues 页面看一眼。你能看到队列demo.queue的Enqueue Count是你发送的总条数Dequeue Count是消费者消费的条数两个数字相等说明链路是完整的。如果Enqueue Count很高但Dequeue Count是 0说明消费者虽然显示在线但实际没有真正拉消息——检查connection.Start()是否在CreateConsumer之后。Topics 页面则要复杂一点除非你用了持久化订阅否则离线消费者的消息你看不到Pending只能看Enqueue Count。验证 Topic 链路是否通我习惯先跑消费者再跑生产者然后去 Topics 页面看Enqueue Count从 0 涨到 N同时 C# 控制台打印出 N 条消息两者对得上才放心。6.2 事务会话批量发送和批量确认的正确姿势CreateSession的第二个参数是事务开关前面演示的代码写的都是false。如果你要批量发消息比如一次导入 1000 条数据每条都开一次事务性能会很差正确做法是开启事务会话批量发送后统一提交using var session connection.CreateSession(AcknowledgementMode.AutoAcknowledge, true); var destination session.GetQueue(batch.queue); using var producer session.CreateProducer(destination); for (int i 0; i 1000; i) { var message session.CreateTextMessage($Message {i}); producer.Send(message); } session.Commit(); // 全部确认这里有个细节CreateSession的第一个参数在事务模式下其实被忽略了因为事务本身就包含了确认语义——Commit()成功后消息才会真正进入队列Rollback()则丢弃。注意session.Rollback()在事务模式里是有用的比如循环里处理到第 500 条发现数据格式不对你可以Rollback()回滚整个批次。但别指望它只回滚一条JMS 事务的最小粒度是“当前未提交的所有消息”不是单条。6.3 我的收尾习惯把连接生命周期封装成工具类最后说一个我自己的工程习惯也是我认为比 Demo 本身更有价值的地方不要在生产项目里到处new ConnectionFactory。写一个MqHelper静态类把连接工厂、会话创建、重试策略封装起来业务代码只传队列名和消息内容。public static class MqHelper { private static readonly ConnectionFactory Factory new ConnectionFactory(tcp://127.0.0.1:61616); public static void SendText(string queueName, string text) { using var connection Factory.CreateConnection(); connection.Start(); using var session connection.CreateSession( AcknowledgementMode.AutoAcknowledge, false); var destination session.GetQueue(queueName); using var producer session.CreateProducer(destination); producer.Send(session.CreateTextMessage(text)); } }这样做的好处是连接工厂是线程安全的可以在多线程环境共享using块保证连接和会话一定会释放掉将来要加认证信息、故障转移地址、SSL 配置只需要改这一个类。别小看这个封装我见过太多项目里每个页面各自 new 一个连接工厂线上一高并发Broker 连接数直接爆掉控制台一堆红色告警。最后说一个我的个人习惯每次改完连接相关的代码我都会把消费者进程杀掉然后看控制台Dequeue Count是否停止增长如果还在涨说明有别的消费者在抢消息——这种排查思路能帮你省下很多跟“消息被谁吃了”纠缠的时间。希望这篇 Demo 笔记能帮你把 ActiveMQ 从“听说过”变成“跑通了”下一步不管是接上位机还是做服务解耦你都有一个能兜底的模板了。本文还有配套的精品资源点击获取
返回列表