`

RocketMQ学习入门

阅读更多

一、部署
1.从https://github.com/alibaba/RocketMQ下载安装包。
2.tar -xf ***.tar 解压tar包。
安装git yum install git
3.echo $JAVA_HOME 确认安装java环境变量。
4.export JAVA_HOME='*****' 设置环境变量。
5.安装nameserver,打开安装路径的bin目录,执行 nohup mqnamesrv & 命令。
6.设置环境nameserver环境变量,export NAMESRV_ADDR=192.168.0.1:9876。
7.设置RocketMQ的的安装位置环境变量ROCKATMQ_HOME
8.安装broker,打开安装路径的bin目录,运行 mqbroker -n "192.168.0.1:9876"(如果设置了环境变量,-n参数可以省略)。
ps:通过nohup.out可以查看安装启动日志。
二、 Broker集群部署
   推荐的几种 Broker 集群部署方式,这里的Slave 不可写,但可读,类似与 Mysql 主备方式。
1.单个 Master
   这种方式风险较大,一旦Broker 重启或者宕机时,会导致整个服务不可用,不建议线上环境使用。
2.多 Master 模式
   一个集群无 Slave,全是 Master,例如 2 个 Master 或者 3 个 Master
   优点:配置简单,单个Master 宕机或重启维护对应用无影响,在磁盘配置为 RAID10 时,即使机器宕机不可恢复情况下,由与 RAID10 磁盘非常可靠,消息也不会丢(异步刷盘丢失少量消息,同步刷盘一条不丢)。性能最高。
   缺点:单台机器宕机期间,这台机器上未被消费的消息在机器恢复之前不可订阅,消息实时性会受到受到影响。
   ###  先启动 NameServer,例如机器 IP 为:192.168.1.1:9876
nohup sh mqnamesrv &
   ###  在机器 A,启动第一个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-noslave/broker-a.properties &
   ###  在机器 B,启动第二个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-noslave/broker-b.properties &
3.多 Master 多 Slave 模式,异步复制
   每个 Master 配置一个 Slave,有多对Master-Slave,HA 采用异步复制方式,主备有短暂消息延迟,毫秒级。
   优点:即使磁盘损坏,消息丢失的非常少,且消息实时性不会受影响,因为 Master 宕机后,消费者仍然可以从 Slave 消费,此过程对应用透明。不需要人工干预。性能同多 Master 模式几乎一样。
   缺点:Master 宕机,磁盘损坏情况,会丢失少量消息。
   ###  先启动 NameServer,例如机器 IP 为:192.168.1.1:9876
nohup sh mqnamesrv &
   ###  在机器 A,启动第一个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-async/broker-a.properties &
   ###  在机器 B,启动第二个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-async/broker-b.properties &
   ###  在机器 C,启动第一个 Slave
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-async/broker-a-s.properties &
  ###  在机器 D,启动第二个 Slave
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-async/broker-b-s.properties &
4.  多 Master 多 Slave 模式,同步双写
   每个 Master 配置一个 Slave,有多对Master-Slave,HA 采用同步双写方式,主备都写成功,向应用返回成功。
   优点:数据与服务都无单点,Master宕机情况下,消息无延迟,服务可用性与数据可用性都非常高
   缺点:性能比异步复制模式略低,大约低 10%左右,发送单个消息的 RT 会略高。目前主宕机后,备机不能自动切换为主机,后续会支持自动切换功能。
   ###  先启动 NameServer,例如机器 IP 为:192.168.1.1:9876
nohup sh mqnamesrv &
   ###  在机器 A,启动第一个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-sync/broker-a.properties &
   ###  在机器 B,启动第二个 Master
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-sync/broker-b.properties &
   ###  在机器 C,启动第一个 Slave
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-sync/broker-a-s.properties &
   ###  在机器 D,启动第二个 Slave
nohup sh mqbroker -n 192.168.1.1:9876 -c$ROCKETMQ_HOME/conf/2m-2s-sync/broker-b-s.properties &
   以上 Broker 与 Slave 配对是通过指定相同的brokerName 参数来配对,Master 的 BrokerId 必须是 0,Slave 的BrokerId 必须是大与 0 的数。另外一个 Master 下面可以挂载多个 Slave,同一 Master 下的多个 Slave 通过指定不同的 BrokerId 来区分。
   $ROCKETMQ_HOST 指的 RocketMQ 安装目录,需要用户自己设置此环境变量
三、启动
1.启动nameserver。 nohup sh mqnamesrv &
2.启动broker。nohup sh mqbroker &
四、使用
1.构造消息的生成者producer和消息的消费者consumer。
2.在maven中添加如下dependency.

    <dependency>
    <groupId>com.alibaba.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>3.0.2</version>
</dependency>
<dependency>
    <groupId>com.alibaba.rocketmq</groupId>
    <artifactId>rocketmq-remoting</artifactId>
    <version>3.0.2</version>
</dependency>
<dependency>
    <groupId>com.alibaba.rocketmq</groupId>
    <artifactId>rocketmq-broker</artifactId>
    <version>3.0.4-open</version>
</dependency>
<dependency>
    <groupId>com.alibaba.rocketmq</groupId>
    <artifactId>rocketmq-common</artifactId>
    <version>3.0.2</version>
</dependency>
<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty</artifactId>
    <version>3.8.0.Final</version>
</dependency>
<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-common</artifactId>
    <version>4.0.7.Final</version>
</dependency>
<dependency>
    <groupId>org.apache.httpcomponents</groupId>
    <artifactId>httpclient</artifactId>
    <version>4.0.1</version>
</dependency>
<dependency>
        <groupId>com.qq.sdk</groupId>
        <artifactId>qzone-sdk</artifactId>
        <version>1.0.0</version>
      </dependency>
<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-buffer</artifactId>
    <version>4.0.7.Final</version>
</dependency>
<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-all</artifactId>
    <version>4.0.10.Final</version>
</dependency>
<dependency>
  <groupId>com.alibaba</groupId>
  <artifactId>fastjson</artifactId>
  <version>1.1.41</version>
</dependency>

 
3.Producer代码如下所示:

 

package rocketmq;

import com.alibaba.rocketmq.client.exception.MQClientException;
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) throws MQClientException, InterruptedException {
        /**
         * 一个应用创建一个Producer,由应用来维护此对象,可以设置为全局对象或者单例<br>
         * 注意:ProducerGroupName需要由应用来保证唯一<br>
         * ProducerGroup这个概念发送普通的消息时,作用不大,但是发送分布式事务消息时,比较关键,
         * 因为服务器会回查这个Group下的任意一个Producer
         */
        DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupName");
        producer.setNamesrvAddr("10.10.0.102:9876");
        /**
         * Producer对象在使用之前必须要调用start初始化,初始化一次即可<br>
         * 注意:切记不可以在每次发送消息时,都调用start方法
         */
        producer.start();

        /**
         * 下面这段代码表明一个Producer对象可以发送多个topic,多个tag的消息。
         * 注意:send方法是同步调用,只要不抛异常就标识成功。但是发送成功也可会有多种状态,<br>
         * 例如消息写入Master成功,但是Slave不成功,这种情况消息属于成功,但是对于个别应用如果对消息可靠性要求极高,<br>
         * 需要对这种情况做处理。另外,消息可能会存在发送失败的情况,失败重试由应用来处理。
         */
        for (int i = 0; i < 10; 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();
            }

        /**
         * 应用退出时,要调用shutdown来清理资源,关闭网络连接,从MetaQ服务器上注销自己
         * 注意:我们建议应用在JBOSS、Tomcat等容器的退出钩子里调用shutdown方法
         */
        producer.shutdown();
    }
}

 
4.Consumer代码如下所示:

package rocketmq;

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.client.exception.MQClientException;
import com.alibaba.rocketmq.common.consumer.ConsumeFromWhere;
import com.alibaba.rocketmq.common.message.MessageExt;


public class PushConsumer {

    /**
     * 当前例子是PushConsumer用法,使用方式给用户感觉是消息从RocketMQ服务器推到了应用客户端。<br>
     * 但是实际PushConsumer内部是使用长轮询Pull方式从Broker拉消息,然后再回调用户Listener方法<br>
     */
    public static void main(String[] args) throws InterruptedException, MQClientException {
        /**
         * 一个应用创建一个Consumer,由应用来维护此对象,可以设置为全局对象或者单例<br>
         * 注意:ConsumerGroupName需要由应用来保证唯一
         */
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CID_001");
         consumer.setNamesrvAddr("10.10.0.102:9876");
        // consumer.setNamesrvAddr("127.0.0.1:9876");

        /**
         * 订阅指定topic下tags分别等于TagA或TagC或TagD
         */
        consumer.subscribe("TopicTest1", "TagA || TagC || TagD");
        /**
         * 订阅指定topic下所有消息<br>
         * 注意:一个consumer对象可以订阅多个topic
         */
        consumer.subscribe("TopicTest2", "*");
        consumer.subscribe("TopicTest3", "*");

        /**
         * 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费<br>
         * 如果非第一次启动,那么按照上次消费的位置继续消费
         */
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);

        consumer.registerMessageListener(new MessageListenerConcurrently() {

            /**
             * 默认msgs里只有一条消息,可以通过设置consumeMessageBatchMaxSize参数来批量接收消息
             */
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                    ConsumeConcurrentlyContext context) {
                System.out.println(Thread.currentThread().getName() + " Receive New Messages: " + msgs);

                MessageExt msg = msgs.get(0);
                if (msg.getTopic().equals("TopicTest1")) {
                    // 执行TopicTest1的消费逻辑
                    if (msg.getTags() != null && msg.getTags().equals("TagA")) {
                        // 执行TagA的消费
                    System.out.println("TagA开始。");
                    }
                    else if (msg.getTags() != null && msg.getTags().equals("TagC")) {
                    System.out.println("TagC开始。");
                        // 执行TagC的消费
                    }
                    else if (msg.getTags() != null && msg.getTags().equals("TagD")) {
                        // 执行TagD的消费
                    System.out.println("TagD开始。");
                    }
                }
                else if (msg.getTopic().equals("TopicTest2")) {
                    // 执行TopicTest2的消费逻辑
                System.out.println("TopicTest2");
                }

                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        /**
         * Consumer对象在使用之前必须要调用start初始化,初始化一次即可<br>
         */
        consumer.start();

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

 

RocketMQ开发文档见附件。

 

分享到:
评论
1 楼 Chrome_1 2017-05-12  
 

相关推荐

    day09-分布式消息系统RocketMQ的入门.zip

    分布式消息系统RocketMQ是阿里巴巴开源的一款高性能、高可用、稳定且易用的消息中间件,它在JavaEE开发中扮演着重要角色...学习过程中,通过实际操作和案例分析,将深入理解RocketMQ的工作原理及其在实际项目中的应用。

    消息中间件+RocketMq+入门文档+用于学习

    【消息中间件RocketMQ入门】 消息中间件(MQ)是一种在分布式系统中用于通信的技术,它通过消息队列作为消息的容器,使得系统组件之间的交互变得更加高效和解耦。RocketMQ是由阿里巴巴开源的高性能消息中间件,已被...

    Apache RocketMQ 从入门到实战1

    1.3 RocketMQ学习环境搭建指南篇: 本章指导读者如何在本地或虚拟环境中搭建RocketMQ,包括安装Java运行环境、下载RocketMQ源码或二进制包、配置环境变量、启动各个组件(NameServer、Broker、Producer、Consumer)...

    分布式消息系统RocketMQ的入门.docx

    在学习或测试环境中,可以适当调整配置,减少内存需求。 RocketMQ中的**Topic**是用来区分不同种类消息的逻辑概念,例如User、Order等。消息实际存储在**Message Queue**中,每个Topic可以有多个Message Queue,...

    RocketMQ入门实战及源码解析.7z

    RocketMQ是中国阿里巴巴开源的一款分布式消息中间件,它在大规模分布式系统中扮演着重要的角色,用于处理高并发、低延迟的...提供的"RocketMQ源码解析.pdf"和"RocketMQ入门指南.pdf"将是深入理解RocketMQ的宝贵资源。

    RocketMq学习视频

    一、rocketmq入门到精通视频教程目录大纲 001-001_RocketMQ_简介 002-002_RocketMQ_核心概念详解 003-003_RocketMQ_集群构建模型详解(一) 004-004_RocketMQ_集群构建模型详解(二) 005-005_RocketMQ_双主模式集群...

    【RocketMQ windows小白安装教程】

    本教程将详细讲解如何在Windows操作系统上进行RocketMQ的安装与配置,适合初学者快速入门。 一、 RocketMQ的组成部分 RocketMQ主要包含四个核心组件:NameServer、Broker、Producer和Consumer。NameServer是服务...

    全面解剖RocketMQ和项目实战视频教程

    视频详细讲解,需要的小伙伴自行百度网盘下载,...2. RocketMQ快速入门 3. RocketMQ集群搭建 4. 消息发送样例 5. 案例介绍 6. 技术分析 7. 环境搭建 8. 下单业务 9. 支付业务 10. 整体联调 11. 高级功能 12. 源码分析

    RocketMQ实战与原理解析.zip

    《RocketMQ实战与原理解析》是一本专为IT专业人士准备的深度学习Apache RocketMQ的指南。这本书分为两大部分:基础入门和源码分析,旨在帮助读者从零开始熟悉RocketMQ,并逐步深入到其核心机制的理解。 在基础入门...

    rocketmq-demo.zip

    这个"rocketmq-demo.zip"压缩包提供了一个入门级的示例,帮助开发者理解RocketMQ的基本工作原理和使用方法。以下是对RocketMQ及其相关代码示例的详细解释。 首先,RocketMQ的核心功能是作为一个消息队列,它在生产...

    rocketmq-all-4.8.0-source-release.zip

    总之,`rocketmq-all-4.8.0-source-release.zip`提供了一个深入了解和定制RocketMQ的平台,无论是对于学习分布式消息队列的原理,还是在实际项目中优化RocketMQ的使用,都是极其有价值的资源。通过阅读源码和实践...

    rocketmq消息中间件.zip

    - **快速入门**:安装部署 RocketMQ,编写 Producer 和 Consumer 示例代码,理解基本操作流程。 - **消息过滤**:学习使用 SQL92 进行消息过滤,实现动态路由。 - **消息轨迹追踪**:了解如何开启消息轨迹,用于...

    rocketmq小示例项目及Linux下的编译安装说明

    在"RocketMQ学习笔记(1)_Linux下安装RocketMQ.pdf"中,应该详细阐述了这些步骤,并可能包括一些最佳实践和常见问题解决方案。 项目中的`.classpath`和`.project`文件是Eclipse IDE的项目配置文件,它们定义了项目的...

    RocketMQ 的 jar 文件 开发文档等资源

    综上所述,这个压缩包文件包含了开发 RocketMQ 应用所需的核心 JAR 包和文档资源,对于想要学习和使用 RocketMQ 的开发者来说,是非常宝贵的资料。通过深入理解 RocketMQ 的设计理念、组件、API 使用及社区支持,...

    rocketmq工具以及官方源码4.2.0版本

    RocketMQ是阿里巴巴开源的一款分布式消息中间件,广泛...通过深入学习源码,开发者可以更好地理解和定制RocketMQ,以适应各种复杂的企业级应用场景。同时,提供的压缩包文件可以方便地下载和部署,为开发者提供了便利。

    SpringCloud入门 nacos 、sentinel、rocketMQ、dubbo、

    本教程将围绕Spring Cloud的几个关键组件进行入门讲解,包括Nacos、Sentinel、RocketMQ以及Dubbo。 **Nacos** Nacos是阿里巴巴开源的动态服务发现和配置管理平台,它可以帮助开发者快速地构建分布式应用系统。Nacos...

    Tedu5阶段RocketMQ中核心概念

    在RocketMQ的学习过程中,创建一个简单的项目来进行实践是非常有益的。 ##### 2.3.1 创建项目 - 使用Maven或者Gradle创建一个新的Java项目。 - 添加RocketMQ的依赖:`&lt;groupId&gt;org.apache.rocketmq&lt;/groupId&gt;` 和 `...

    rocketmq-learning:结合应用场景介绍如何使用rocketmq,尝试打造rocketmq领域最佳的入门示例,并包含40+篇原创博文,两本免费电子书

    火箭学习结合应用场景介绍如何使用rocketmq,尝试打造rocketmq领域最佳的入门示例代码介绍代码组织将遵循功能,每一个功能都会结合特定使用场景。 transaction事务消息发送。RocketMQ免费电子书获取方式:关注...

    阿里云 专有云企业版 V3.12.0 消息队列 RocketMQ 版 MQTT用户指南 20200623.pdf

    快速入门部分通常会涵盖如何创建RocketMQ实例、配置MQTT客户端连接、发布和订阅消息的基本步骤。用户需要了解如何在阿里云控制台上操作,包括实例创建、 Topic管理、ACL设置等。此外,还会提供MQTT客户端SDK的使用...

    Java大数据学习笔记.zip

    Java大数据学习笔记大数据专题JVM春天書目《Java编程思想》《Mybatis从入门到精通》《深入分析Java Web技术内幕》《Java设计模式》《Java EE框架技术》《自顶向下方法》《Spark机器学习进阶实战》《Java编程从入门到...

Global site tag (gtag.js) - Google Analytics