惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

J
Java Code Geeks
月光博客
月光博客
aimingoo的专栏
aimingoo的专栏
Google DeepMind News
Google DeepMind News
Recent Announcements
Recent Announcements
MyScale Blog
MyScale Blog
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
S
SegmentFault 最新的问题
Hugging Face - Blog
Hugging Face - Blog
Martin Fowler
Martin Fowler
WordPress大学
WordPress大学
F
Fortinet All Blogs
小众软件
小众软件
D
Docker
U
Unit 42
博客园 - 聂微东
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
爱范儿
爱范儿
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
IT之家
IT之家
云风的 BLOG
云风的 BLOG
博客园 - 司徒正美
有赞技术团队
有赞技术团队
腾讯CDC

博客园 - wzwyc

FreeSql事务型提交 搜索并下载WINDOWS系统补丁 VS2019/2022/2026内网安装闪退的问题 Windows 7开发环境的问题笔记 C#颜色转换 C#把整型数值转16进制字符串 20260507随笔 在Zed中设置Xiaomi MiMo OxyPlot显示Legend图例 apache-iotdb-0.13.3-all-bin在WINDOWS 11系统上部署运行的问题 任务栏上的图标无法正常显示 应用启动时加载XML报错的问题处理 Prism的LoadedCommand命令没有被调用的问题 C#编译报错:可在所有平台上访问此调用站点。"xxxxxxxxxx" 仅在 "Windows" 7.0 及更高版本 上受支持。 C#系统下LINUX系统下组播数据的接收 消除只能在WINDOWS系统上使用的编译报警提示 WpfMultiStyle笔记 Kafka调研笔记 利用HtmlAgilityPack抓取网页的标题 关于Process的使用 Anaconda使用笔记 python下的plotly把图表转换成HTML 利用NAudio实现对音频设备的控制(静音、音量调节) 利用MathNet.Numerics求均方根 .NET 8以上版本,字节数组转字符串 .NET 8项目下载所有依赖到指定目录 异常的处理
RocketMQ调研笔记
wzwyc · 2026-07-01 · via 博客园 - wzwyc

前面有个项目,客户要求使用RocketMQ,前面我们做了调研,使用了RocketMQ.Client这个库。
后面发现客户那边使用的是RocketMQ 4.9.7的版本,而RocketMQ.Client这个库应该是只支持最新的RocketMQ 5.X版本的。
所以需要重新适配。

目前是改用NewLife.RocketMQ。

前面同事反应,使用了NewLife.RocketMQ会有NLog日志无法输出的问题。我这会儿用测试代码试了一下,貌似又没有出现这个问题,貌似还是能正常输出的。
后面跟同事确认了一下,NLog不能正常输出的是RocketMQ.Client这个库。

安装RocketMQ 4.9.7

下载软件包,我是直接下载二进软件包。
到下面这个地址选择对应的版本进行下载即可:
https://dist.apache.org/repos/dist/release/rocketmq/

我下载的相当于是:
https://dist.apache.org/repos/dist/release/rocketmq/4.9.7/rocketmq-all-4.9.7-bin-release.zip

软件包解压缩到指定的目录。

安装Java

RocketMQ 4.X对JAVA版本有要求,太新的版本可能会执行失败,推荐版本是Java 8,我下载的是jdk-8u202-windows-x64.exe。

JAVA官方的网址下载很麻烦,需要注册登录。可以直接用华为的镜像站:
https://mirrors.huaweicloud.com/java/jdk/8u202-b08/

安装完成后需要设置一下系统环境变量:JAVA_HOME和ROCKETMQ_HOME,分别是JAVA的路径和RocketMQ的路径。
设置完成后需要注销一下或重启一下系统。
可以用java -version试一下是否安装成功。

C:\Users\wyc>java -version
java version "1.8.0_202"
Java(TM) SE Runtime Environment (build 1.8.0_202-b08)
Java HotSpot(TM) 64-Bit Server VM (build 25.202-b08, mixed mode)

启动NameServer

打开控制台:

cd /d E:\Downloads\rocketmq-all-4.9.7-bin-release\bin
mqnamesrv.cmd

启动Broker

打开控制台:

cd /d E:\Downloads\rocketmq-all-4.9.7-bin-release\bin
mqbroker.cmd -n localhost:9876

生产者

using System;
using NewLife.RocketMQ;

namespace ProducerApp
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("RocketMQ 生产者演示");
            Console.WriteLine("连接到 NameServer: localhost:9876");
            Console.WriteLine("Topic: test_topic, Group: producer_group");
            Console.WriteLine();

            // 创建生产者实例
            var producer = new Producer
            {
                Topic = "test_topic",
                NameServerAddress = "localhost:9876",
                Group = "producer_group"
            };

            try
            {
                // 启动生产者
                producer.Start();
                Console.WriteLine("生产者已启动,开始发送消息...");

                // 发送10条测试消息
                for (int i = 1; i <= 10; i++)
                {
                    var message = $"测试消息 #{i} - 发送时间: {DateTime.Now:yyyy-MM-dd HH:mm:ss}";
                    
                    // 同步发送消息
                    var result = producer.Publish(message, "TestTag", $"Key_{i}");
                    
                    Console.WriteLine($"消息 {i}: {message}");
                    Console.WriteLine($"  发送状态: {result.Status}");
                    Console.WriteLine($"  消息ID: {result.MsgId}");
                    Console.WriteLine();
                    
                    System.Threading.Thread.Sleep(1000); // 每秒发送一条
                }

                Console.WriteLine("所有消息发送完成!");
            }
            catch (Exception ex)
            {
                Console.WriteLine($"发送消息时发生错误: {ex.Message}");
                Console.WriteLine(ex.StackTrace);
            }
            finally
            {
                // 停止生产者
                producer.Dispose();
                Console.WriteLine("生产者已停止。");
            }

            Console.WriteLine("按任意键退出...");
            try
            {
                Console.ReadKey();
            }
            catch (InvalidOperationException)
            {
                // 非交互式环境,自动退出
                Console.WriteLine("(非交互式环境,自动退出)");
            }
        }
    }
}

消费者

using System;
using NewLife.RocketMQ;

namespace ConsumerApp
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("RocketMQ 消费者演示");
            Console.WriteLine("连接到 NameServer: localhost:9876");
            Console.WriteLine("Topic: test_topic, Group: consumer_group");
            Console.WriteLine();

            // 创建消费者实例
            var consumer = new Consumer
            {
                Topic = "test_topic",
                Group = "consumer_group",
                NameServerAddress = "localhost:9876"
            };

            // 设置消费回调
            consumer.OnConsume = (queue, messages) =>
            {
                Console.WriteLine($"[{DateTime.Now:yyyy-MM-dd HH:mm:ss}] 收到 {messages.Length} 条消息");
                Console.WriteLine($"  队列: {queue.BrokerName} - {queue.QueueId}");
                
                foreach (var msg in messages)
                {
                    Console.WriteLine($"  消息ID: {msg.MsgId}");
                    Console.WriteLine($"  内容: {msg.BodyString}");
                    Console.WriteLine($"  标签: {msg.Tags}");
                    Console.WriteLine($"  键: {msg.Keys}");
                    Console.WriteLine($"  生产时间: {msg.BornTimestamp:yyyy-MM-dd HH:mm:ss}");
                    Console.WriteLine();
                }
                
                return true; // 返回 true 表示消费成功
            };

            try
            {
                // 启动消费者
                consumer.Start();
                Console.WriteLine("消费者已启动,等待接收消息...");
                Console.WriteLine("按任意键停止消费者并退出。");
                Console.ReadKey();
            }
            catch (Exception ex)
            {
                Console.WriteLine($"消费消息时发生错误: {ex.Message}");
                Console.WriteLine(ex.StackTrace);
            }
            finally
            {
                // 停止消费者
                consumer.Dispose();
                Console.WriteLine("消费者已停止。");
            }
        }
    }
}