RocketMq消息队列使用

简介: 最近在看消息队列框架 ,alibaba的RocketMQ单机支持1万以上的持久化队列,支持诸多特性,目前RocketMQ在阿里集团被广泛应用在订单,交易,充值,流计算,消息推送,日志流式处理,binglog分发等场景比kafka还是有过之无不及,其实kafka文档很丰富但RocketMQ网上的...

最近在看消息队列框架 ,alibaba的RocketMQ单机支持1万以上的持久化队列,支持诸多特性,

目前RocketMQ在阿里集团被广泛应用在订单,交易,充值,流计算,消息推送,日志流式处理,binglog分发等场景

比kafka还是有过之无不及,其实kafka文档很丰富

但RocketMQ网上的文章太少,找不到相关的操作教程

于是研究了下源码 做个单机操作的教程,如果你也对此有兴趣不妨共同研究

下载源码的地址 https://github.com/alibaba/RocketMQ/releases

  • 首选通过在java项目里面Maven依赖方式引用RocketMQ Java SDK

    <dependency>
        <groupId>com.alibaba.rocketmq</groupId>
        <artifactId>rocketmq-client</artifactId>
        <version>3.2.6</version>
    </dependency>

Downloads

在linux 下用wget 下载源码然后解压出来

在runserver.sh里面可以配置 jvm启动的参数 JAVA_OPT_1="-server -Xms4g -Xmx4g -Xmn2g -XX:PermSize=128m -XX:MaxPermSize=320m"  

可以 vi runserver.sh

分别给 mqnamesrv mqbroker play.sh 执行的权限

chmod +x  mqnamersrv 

chmod +x  mqbroker 

chmod +x  play.sh 

下面红线框的这段 命令输入错误了,忽略不用看

通过 nohup sh mqnamesrv& 启动 RocketMq

目前没看到结束的命令,也没找到相关的介绍,

我这里用的 ps -ef|grep rocketmq  查到进程pid

然后kill pid号

或则pkill -9 java [慎用]

用jps -v 查看下java进程的参数

 rocketmq启动后监听 9876端口,这里还是在看源码里面看到的,资料实在是太少了

在防火墙配置里面加上 9876端口,设置iptables对外开放

部署Broker 

nohup sh mqbroker -n "127.0.0.1:9876" -c ../conf/2m-2s-async/broker-a.properties & 

这里ip换成本机的就是单机实例,如果配置主从 这里可以配其他的ip

 Master和Slave的配置文件参考conf目录下的配置文件

 Master与Slave通过指定相同的brokerName参数来配对,Master的BrokerId必须是0,Slave的BrokerId必须是大于0的数

 一个Master下面可以挂载多个Slave,同一Master下的多个Slave通过指定不同的BrokerId来区分

 部署一Master一Slave,集群采用异步复制方式:

 Master: nohup sh mqbroker -n "192.168.1.23:9876" -c ../conf/2m-2s-async/broker-a.properties &  

Slave:   nohup sh mqbroker -n "192.168.1.23:9876" -c ../conf/2m-2s-async/broker-a-s.properties &  

 

 

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
package  com.pgsqlmybatis.common.rocketmq; /*
***************************************************************
* 公司名称    :
* 系统名称    :信用管家专业版
* 类 名 称    :Ios渠道idfa统计,推广统计用
* 功能描述    :
* 业务描述    :
* 作 者 名    :@Author Royal
* 开发日期    :2016-05-15
* Created     :IntelliJ IDEA
***************************************************************
* 修改日期    :
* 修 改 者    :
* 修改内容    :
***************************************************************
*/
 
import  com.alibaba.rocketmq.client.producer.DefaultMQProducer;
import  com.alibaba.rocketmq.client.producer.SendResult;
import  com.alibaba.rocketmq.common.message.Message;
 
public  class  Producer {
     public  static  void  main(String[] args) {
         DefaultMQProducer producer =  new  DefaultMQProducer( "Producer" );
         producer.setNamesrvAddr( "xxxxxxxxxx:9876" );
         try  {
             producer.start();
 
             String pushMsg= "kafka activeMq rocketMq 消息队列使用1" ;
             Message msg =  new  Message( "PushTopic" , "push" , "1" ,
                     pushMsg.getBytes( "UTF-8" ));
 
             SendResult result = producer.send(msg);
             System.out.println( "id:"  + result.getMsgId() +
                     " result:"  + result.getSendStatus());
 
             String pushMsg2= "海量级消息记录单机测试2" ;
             msg =  new  Message( "PushTopic" , "push" , "2" ,pushMsg2.getBytes( "UTF-8" ));
 
             result = producer.send(msg);
             System.out.println( "id:"  + result.getMsgId() +
                     " result:"  + result.getSendStatus());
 
             String pushMsg3= "海量级消息记录单机测试3" ;
             msg =  new  Message( "PullTopic" , "pull" , "1" ,pushMsg3.getBytes());
 
             result = producer.send(msg);
             System.out.println( "id:"  + result.getMsgId() +
                     " result:"  + result.getSendStatus());
         catch  (Exception e) {
             e.printStackTrace();
         finally  {
             producer.shutdown();
         }
     }
}

  

启动生成者

 

启动消费者

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
package  com.pgsqlmybatis.common.rocketmq; /*
***************************************************************
* 公司名称    :
* 系统名称    :信用管家专业版
* 类 名 称    :Ios渠道idfa统计,推广统计用
* 功能描述    :
* 业务描述    :
* 作 者 名    :@Author Royal
* 开发日期    :2016-05-15
* Created     :IntelliJ IDEA
***************************************************************
* 修改日期    :
* 修 改 者    :
* 修改内容    :
***************************************************************
*/
 
import  java.io.UnsupportedEncodingException;
import  java.util.List;
 
import  com.alibaba.rocketmq.client.consumer.DefaultMQPushConsumer;
import  com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import  com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import  com.alibaba.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import  com.alibaba.rocketmq.common.consumer.ConsumeFromWhere;
import  com.alibaba.rocketmq.common.message.Message;
import  com.alibaba.rocketmq.common.message.MessageExt;
 
public  class  Consumer {
     public  static  void  main(String[] args){
         DefaultMQPushConsumer consumer =
                 new  DefaultMQPushConsumer( "PushConsumer" );
         consumer.setNamesrvAddr( "xxxxxxxxxxxx:9876" );
         try  {
             consumer.subscribe( "PushTopic" "push" );
             /**
              * 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费<br>
              * 如果非第一次启动,那么按照上次消费的位置继续消费
              */
             consumer.setConsumeFromWhere(
                     ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
             consumer.registerMessageListener(
                     new  MessageListenerConcurrently() {
                         public  ConsumeConcurrentlyStatus consumeMessage(
                                 List<MessageExt> list,
                                 ConsumeConcurrentlyContext Context) {
                             Message msg = list.get( 0 );
                             System.out.println(msg.toString());
                             String recString=  null ;
                             try  {
                                 recString =  new  String(msg.getBody() , "UTF-8" );
                             catch  (UnsupportedEncodingException e) {
                                 e.printStackTrace();
                             }
                             System.out.println(recString);
                             return  ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                         }
                     }
             );
             consumer.start();
         catch  (Exception e) {
             e.printStackTrace();
         }
     }
}

   

 

以上为单机实例配置

如果你遇到什么问题可以私信我,如果觉得此文对你很有帮助,点下赞推荐下额^_^ 

参考:http://blog.csdn.net/a19881029/article/details/34446629

        http://sofar.blog.51cto.com/353572/1540874

        http://blog.csdn.net/loongshawn/article/details/51086876

        RocketMq最佳实践

       《RocketMQ原理简介》

       分布式开放消息系统(RocketMQ)的原理与实践      

       《RocketMQ用户指南》

相关实践学习
快速体验阿里云云消息队列RocketMQ版
本实验将带您快速体验使用云消息队列RocketMQ版Serverless系列实例进行获取接入点、创建Topic、创建订阅组、收发消息、查看消息轨迹和仪表盘。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
相关文章
|
4月前
|
存储 人工智能 运维
1949AI 轻量化 AI 自动化 本地自动化工具浏览器自动化 Agent 自动化工具 自动化运维状态监测与消息推送技术实践
1949AI是一款轻量化AI自动化工具,专注本地化、低资源、零配置运维实践。支持浏览器自动化监测、状态智能判定、本地日志存储与消息推送,适配低配电脑与个人/小型团队,安全合规、开箱即用。(239字)
|
6月前
|
消息中间件 运维 监控
别只盯着充电枪:聊聊一个真正“能赚钱、能扩展、能运维”的智慧充电桩系统架构
别只盯着充电枪:聊聊一个真正“能赚钱、能扩展、能运维”的智慧充电桩系统架构
391 8
|
6月前
|
存储 弹性计算 人工智能
2026年阿里云服务器租用价格表整理:配置体系、收费标准与优惠政策
阿里云服务器租用价格覆盖从 38 元 / 年的轻量配置到数万元 / 年的高性能实例,用户需根据业务负载(并发量、算力需求)与使用周期选择:个人 / 小微团队优先轻量服务器或经济型 e 实例,中小企业核心业务选通用算力型 u1 实例,高负载场景适配第八代高性能实例,AI / 渲染需求选择 GPU 服务器。选型时需结合优惠政策(如续费同价、长期折扣)与附加成本(带宽、存储),通过阿里云官方价格计算器精准核算,确保配置与需求匹配,实现性能与成本的平衡。
|
6月前
|
缓存 JavaScript 开发者
微信小游戏开发的技术难点
微信小游戏开发在2026年面临五大技术挑战:高性能模式下的内存管控、WASM性能瓶颈、4MB主包极速启动、跨平台渲染一致性及开放数据域通信限制。开发者需在严苛环境下实现流畅体验,考验极致优化能力。#微信小游戏 #游戏外包
|
11月前
|
传感器 运维 监控
AR眼镜在工业运维的场景应用和方案说明
AR眼镜通过虚实融合技术,革新工业运维模式。从设备巡检、故障维修到员工培训,AR实现远程协作、实时数据叠加与沉浸式教学,大幅提升效率与准确性,推动智能工厂发展。
|
弹性计算 运维 搜索推荐
幻兽帕鲁内存溢出怎么办,一键设置定时重启,修改虚拟内存,定时清理,轻松解决卡顿!再也不怕爆内存了!
幻兽帕鲁的内存溢出问题,玩久了确实会变卡。这里给出三个解决思路:第一种方法是定时进行内存清理(装个软件就可以),网上也有很多教程,我会把下载地址放在文章后面,大家可以去下载。第二种方法是调大虚拟内存,这个可以一键设置。第三种方法是定时重启游戏服务,这个也可以一键设置。这三种方法我下面都会教给大家,可以有效解决内存增长过快的问题,避免游戏卡顿甚至崩溃。
1593 3
|
数据可视化 数据挖掘 atlas
地图不只是导航:DataV Atlas 揭示地理数据的深层价值
随着互联网场景的快速衍生,打车、外卖、智能驾驶等领域的空间数据爆发式增长,海量数据分析成为日常需求。然而,传统地图服务面临性能、安全和成本挑战。为此,我们推出「DataV Atlas 地理数据服务」,提供高效、安全、易用的地理数据解决方案。通过简单的 SQL 查询即可生成专业地理服务,支持多源数据整合、实时更新与分析,确保数据安全,并深度集成 DataV Board 数据看板,实现一键上屏和交互式分析。适用于大屏展示、城市规划等多种场景,助力企业轻松挖掘空间数据价值。
1063 6
地图不只是导航:DataV Atlas 揭示地理数据的深层价值
|
人工智能 自然语言处理 API
研究大模型门槛太高?不妨看看小模型SLM,知识点都在这
大型语言模型(LLM)在文本生成、问答等领域表现出色,但也面临资源受限环境应用难、领域知识不足及隐私问题等挑战。为此,小型语言模型(SLM)逐渐受到关注,其具备低延迟、成本效益高、易于定制等优点,适合资源受限环境和领域知识获取。SLM可通过预训练、微调和知识蒸馏等技术增强性能,在自然语言处理、计算机视觉等领域有广泛应用潜力。然而,SLM也存在复杂任务表现有限等问题,未来研究将进一步提升其性能与可靠性。 论文链接:https://arxiv.org/abs/2411.03350
752 5
|
JSON Java 数据格式
【小知识】Windows下ElasticSearch 安装与配置
【小知识】Windows下ElasticSearch 安装与配置
1026 0
【小知识】Windows下ElasticSearch 安装与配置

相关产品

  • 云消息队列 MQ