RocketMQ最佳实践(一)4.2版本/概念介绍/安装调试/客户端demo

为什么选择RocketMQ

我们来看看官方回答:

“我们研究发现,对于ActiveMQ而言,随着越来越多的使用queues和topics,其IO成为了瓶颈。某些情况下,消费者缓慢(消费能力不足)还会拖慢生产者(造成消息阻塞)。虽然我们做了最大努力进行优化:节流、断路器或者回退,但是并不能进行优雅的扩展。因此我们开始专注于使用时下非常流行的kafka,但是仍然不能满足我们的要求,如低延迟和高可靠性,详情见这里。在这样的背景下,我们决定开发一个新的消息中间件来处理一系列广泛的使用场景,包括从传统的发布/订阅场景到高容量的实时交易系统中不允许消息丢失的场景。”

各位看官也可以撮这里去看看RocketMQ与ActiveMQ以及Kafka的比较

核心概念

生产者(Producer):消息发送方,将业务系统中产生的消息发送到brokers(brokers可以理解为消息代理,生产者和消费者之间是通过brokers进行消息的通信),rocketmq提供了以下消息发送方式:同步、异步、单向。

生产者组(Producer Group):相同角色的生产者被归为同一组,比如通常情况下一个服务会部署多个实例,这多个实例就是一个组,生产者分组的作用只体现在消息回查的时候,即如果一个生产者组中的一个生产者实例发送一个事务消息到broker后挂掉了,那么broker会回查此实例所在组的其他实例,从而进行消息的提交或回滚操作。

消费者(Consumer):消息消费方,从brokers拉取消息。站在用户的角度,有以下两种消费者。

主动消费者(PullConsumer):从brokers拉取消息并消费。

被动消费者(PushConsumer):内部也是通过pull方式获取消息,只是进行了扩展和封装,并给用户预留了一个回调接口去实现,当消息到底的时候会执行用户自定义的回调接口。

消费者组(Consumer Group):和生产者组类似。其作用体现在实现消费者的负载均衡和容错,有了消费者组变得异常容易。需要注意的是:同一个消费者组的每个消费者实例订阅的主题必须相同。

主题(Topic):主题就是消息传递的类型。一个生产者实例可以发送消息到多个主题,多个生产者实例也可以发送消息到同一个主题。同样的,对于消费者端来说,一个消费者组可以订阅多个主题的消息,一个主题的消息也可以被多个消费者组订阅。

消息(Message):消息就像是你传递信息的信封。每个消息必须指定一个主题,就好比每个信封上都必须写明收件人。

消息队列(Message Queues):在主题内部,逻辑划分了多个子主题,每个子主题被称为消息队列。这个概念在实现最大并发数、故障切换等功能上有巨大的作用。

标签(Tag):标签,可以被认为是子主题。通常用于区分同一个主题下的不同作用或者说不同业务的消息。同时也是避免主题定义过多引起性能问题,通常情况下一个生产者组只向一个主题发送消息,其中不同业务的消息通过标签或者说子主题来区分。

消息代理(Broker):消息代理是RockerMQ中很重要的角色。它接收生产者发送的消息,进行消息存储,为消费者拉取消息服务。它还存储消息消耗相关的元数据,包括消费群体,消费进度偏移和主题/队列信息。

命名服务(Name Server):命名服务作为路由信息提供程序。生产者/消费者进行主题查找、消息代理查找、读取/写入消息都需要通过命名服务获取路由信息。

消息顺序(Message Order):当我们使用DefaultMQPushConsumer时,我们可以选择使用“orderly”还是“concurrently”。

orderly:消费消息的有序化意味着消息被生产者按照每个消息队列发送的顺序消费。如果您正在处理全局顺序为强制的场景,请确保您使用的主题只有一个消息队列。注意:如果指定了消费顺序,则消息消费的最大并发性是消费组订阅的消息队列数。

concurrently:当同时消费时,消息消费的最大并发仅限于为每个消费客户端指定的线程池。注意:此模式不再保证消息顺序。

前提:部署rocketmq依赖JDK(1.8+)和MAVEN(3.2.X)

环境:CentOS7,内存4G

1. 下载最新的rocketmq源码文件

下载地址:http://mirrors.hust.edu.cn/apache/rocketmq/4.2.0/rocketmq-all-4.2.0-source-release.zip

1. 下载压缩包

2. 解压缩

3. mvn安装

shell> wget http://mirrors.hust.edu.cn/apache/rocketmq/4.2.0/rocketmq-all-4.2.0-source-release.zip

shell> unzip rocketmq-all-4.2.0-source-release.zip

shell> cd rocketmq-all-4.2.0/shell> mvn -Prelease-all -DskipTests clean install -U

shell> cd distribution/target/apache-rocketmq

2. 启动Name Server

> nohup sh bin/mqnamesrv &

> tail -f ~/logs/rocketmqlogs/namesrv.log查看启动情况。也可以用jps命令查看

启动后输入jps会出现NamesrvStartup表示启动成功

如果服务器内存不够,可以修改runserver.sh脚本(mqnamesrv文件中通过runserver.sh脚本调用Name Server的主函数com.alibaba.rocketmq.namesrv.NamesrvStartup启动Name Server)中的JAVA_OPT_1参数

JAVA_OPT_1="-server -Xms4g -Xmx4g -Xmn2g -XX:PermSize=128m -XX:MaxPermSize=320m"  

3. 启动Broker

> nohup sh bin/mqbroker -n localhost:9876 & 

注:-n localhost:9876表示对应的namesrv

> tail -f ~/logs/rocketmqlogs/broker.log 查看启动命令状态

输入jps命令看到BrokerStartup表示启动成功

无法启动Broker,很可能是配置文件中的内存设置大于服务器内存

解决办法:

修改bin目录下的runbroker.sh,根据本机内存,修改如下部分即可

JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g -Xmn4g"

4停止服务

> sh bin/mqshutdown broker

The mqbroker(36695) is running...

Send shutdown request to mqbroker(36695) OK

> sh bin/mqshutdown

 namesrvThe mqnamesrv(36664) is running..

Send shutdown request to mqnamesrv(36664) OK

生产者

import org.apache.rocketmq.client.exception.MQClientException;  

import org.apache.rocketmq.client.producer.DefaultMQProducer;  

import org.apache.rocketmq.client.producer.SendResult;  

import org.apache.rocketmq.common.message.Message;  


import java.util.concurrent.TimeUnit;  


public class Producer {  

public static void main(String[] args) throws MQClientException,  

            InterruptedException {  

/**

         * 一个应用创建一个Producer,由应用来维护此对象,可以设置为全局对象或者单例

         * 注意:ProducerGroupName需要由应用来保证唯一

         * ProducerGroup这个概念发送普通的消息时,作用不大,但是发送分布式事务消息时,比较关键,

         * 因为服务器会回查这个Group下的任意一个Producer

         */  

DefaultMQProducer producer =new DefaultMQProducer("ProducerGroupName");  

producer.setNamesrvAddr("192.168.56.101:9876");  

producer.setInstanceName("Producer");  

producer.setVipChannelEnabled(false);  


/**

         * Producer对象在使用之前必须要调用start初始化,初始化一次即可

         * 注意:切记不可以在每次发送消息时,都调用start方法

         */  

        producer.start();  


/**

         * 下面这段代码表明一个Producer对象可以发送多个topic,多个tag的消息。

         * 注意:send方法是同步调用,只要不抛异常就标识成功。但是发送成功也可会有多种状态,

         * 例如消息写入Master成功,但是Slave不成功,这种情况消息属于成功,但是对于个别应用如果对消息可靠性要求极高,

         * 需要对这种情况做处理。另外,消息可能会存在发送失败的情况,失败重试由应用来处理。

         */  

for (int i = 0; i < 1; i++) {  

try {  

                {  

Message msg =new Message("TopicTest1",// topic  

"TagA",// tag  

"OrderID001",// key  

("Hello MetaQ").getBytes());// body  

                    SendResult sendResult = producer.send(msg);  

                    System.out.println(sendResult);  

                }  


                {  

Message msg =new Message("TopicTest2",// topic  

"TagB",// tag  

"OrderID0034",// key  

("Hello MetaQ").getBytes());// body  

                    SendResult sendResult = producer.send(msg);  

                    System.out.println(sendResult);  

                }  


                {  

Message msg =new Message("TopicTest3",// topic  

"TagC",// tag  

"OrderID061",// key  

("Hello MetaQ").getBytes());// body  

                    SendResult sendResult = producer.send(msg    );  

                    System.out.println(sendResult);  

                }  

}catch (Exception e) {  

                e.printStackTrace();  

            }  

TimeUnit.MILLISECONDS.sleep(1000);  

        }  


/**

         * 应用退出时,要调用shutdown来清理资源,关闭网络连接,从MetaQ服务器上注销自己

         * 注意:我们建议应用在JBOSS、Tomcat等容器的退出钩子里调用shutdown方法

         */  

        producer.shutdown();  

    }  

}  

消费者

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;  

import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;  

import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;  

import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;  

import org.apache.rocketmq.client.exception.MQClientException;  

import org.apache.rocketmq.common.message.MessageExt;  


import java.util.List;  

public class PushConsumer {  

/**

     * 当前例子是PushConsumer用法,使用方式给用户感觉是消息从RocketMQ服务器推到了应用客户端。

     * 但是实际PushConsumer内部是使用长轮询Pull方式从MetaQ服务器拉消息,然后再回调用户Listener方法

     */  

public static void main(String[] args) throws InterruptedException,  

            MQClientException {  

/**

         * 一个应用创建一个Consumer,由应用来维护此对象,可以设置为全局对象或者单例

         * 注意:ConsumerGroupName需要由应用来保证唯一

         */  

DefaultMQPushConsumer consumer =new DefaultMQPushConsumer(  

"ConsumerGroupName");  

consumer.setNamesrvAddr("192.168.56.101:9876");  

consumer.setInstanceName("Consumber");  

/**

         * 订阅指定topic下tags分别等于TagA或TagC或TagD

         */  

consumer.subscribe("TopicTest1", "TagA || TagC || TagD");  

/**

         * 订阅指定topic下所有消息

         * 注意:一个consumer对象可以订阅多个topic

         */  

consumer.subscribe("TopicTest2", "*");  

consumer.registerMessageListener(new MessageListenerConcurrently() {  

/**

             * 默认msgs里只有一条消息,可以通过设置consumeMessageBatchMaxSize参数来批量接收消息

             */  

@Override  

public ConsumeConcurrentlyStatus consumeMessage(  

                    List msgs, ConsumeConcurrentlyContext context) {  

                System.out.println(Thread.currentThread().getName()  

+" Receive New Messages: " + msgs.size());  

MessageExt msg = msgs.get(0);  

if (msg.getTopic().equals("TopicTest1")) {  

// 执行TopicTest1的消费逻辑  

if (msg.getTags() != null && msg.getTags().equals("TagA")) {  

// 执行TagA的消费  

System.out.println(new String(msg.getBody()));  

}else if (msg.getTags() != null  

&& msg.getTags().equals("TagC")) {  

// 执行TagC的消费  

}else if (msg.getTags() != null  

&& msg.getTags().equals("TagD")) {  

// 执行TagD的消费  

                    }  

}else if (msg.getTopic().equals("TopicTest2")) {  

System.out.println(new String(msg.getBody()));  

                }  

return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;  

            }  

        });  

/**

         * Consumer对象在使用之前必须要调用start初始化,初始化一次即可

         */  

        consumer.start();  

System.out.println("Consumer Started.");  

    }  

}  

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 205,236评论 6 478
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 87,867评论 2 381
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 151,715评论 0 340
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 54,899评论 1 278
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 63,895评论 5 368
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,733评论 1 283
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 38,085评论 3 399
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,722评论 0 258
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 43,025评论 1 300
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,696评论 2 323
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,816评论 1 333
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,447评论 4 322
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 39,057评论 3 307
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 30,009评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 31,254评论 1 260
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 45,204评论 2 352
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 42,561评论 2 343

推荐阅读更多精彩内容