`

rocketMq 集群(出处:http://blog.csdn.net/zhu_tianwei/article/details/40949523/)

    博客分类:
  • JMS
 
阅读更多
Broker集群部署方式主要有以下几种:(Slave 不可写,但可读)
(1)单个Master
这种方式风险较大,一旦Broker 重启或者宕机时,会导致整个服务不可用,不建议线上环境使用。
(2)多Master模式
一个集群无 Slave,全是 Master,例如 2 个 Master 或者 3 个 Master。
  优点:配置简单,单个Master 宕机或重启维护对应用无影响,在磁盘配置为 RAID10 时,即使机器宕机不可恢复情况下,由与 RAID10 磁盘非常可靠,消息也不会丢(异步刷盘丢失少量消息,同步刷盘一条不丢)。性能最高。
   缺点:单台机器宕机期间,这台机器上未被消费的消息在机器恢复之前不可订阅,消息实时性会受到受到影响。
#先启动 NameServer,例如机器 IP 为:192.168.36.189:9876
nohup ./bin/mqnamesrv >/dev/null 2>&1 &
#在机器 A,启动第一个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-noslave/broker-a.properties >/dev/null 2>&1 &
#在机器 B,启动第二个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-noslave/broker-b.properties >/dev/null 2>&1 &
(3)多Master多Slave模式,异步复制
每个 Master 配置一个 Slave,有多对Master-Slave,HA 采用异步复制方式,主备有短暂消息延迟,毫秒级。
   优点:即使磁盘损坏,消息丢失的非常少,且消息实时性不会受影响,因为 Master 宕机后,消费者仍然可以从 Slave 消费,此过程对应用透明。不需要人工干预。性能同多 Master 模式几乎一样。
   缺点:Master 宕机,磁盘损坏情况,会丢失少量消息。
#先启动 NameServer,例如机器 IP 为:192.168.36.189:9876
nohup ./bin/mqnamesrv >/dev/null 2>&1 &
#在机器 A,启动第一个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-async/broker-a.properties >/dev/null 2>&1 &
#在机器 B,启动第二个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-async/broker-b.properties >/dev/null 2>&1 &
#在机器 C,启动第一个 Slave
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-async/broker-a-s.properties >/dev/null 2>&1 &
#在机器 D,启动第二个 Slave
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-async/broker-b-s.properties >/dev/null 2>&1 &
(4)多Master多Slave模式,同步双写
每个 Master 配置一个 Slave,有多对Master-Slave,HA 采用同步双写方式,主备都写成功,向应用返回成功。
   优点:数据与服务都无单点,Master宕机情况下,消息无延迟,服务可用性与数据可用性都非常高
   缺点:性能比异步复制模式略低,大约低 10%左右,发送单个消息的 RT 会略高。目前主宕机后,备机不能自动切换为主机,后续会支持自动切换功能。
#先启动 NameServer,例如机器 IP 为:192.168.36.189:9876
nohup ./bin/mqnamesrv >/dev/null 2>&1 &
#在机器 A,启动第一个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-sync/broker-a.properties >/dev/null 2>&1 &
#在机器 B,启动第二个 Master
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-sync/broker-b.properties >/dev/null 2>&1 &
#在机器 C,启动第一个 Slave
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-sync/broker-a-s.properties >/dev/null 2>&1 &
#在机器 D,启动第二个 Slave
nohup ./bin/mqbroker -n 192.168.36.189:9876 -c ./conf/2m-2s-sync/broker-b-s.properties >/dev/null 2>&1 &
以上 Broker 与 Slave 配对是通过指定相同的brokerName 参数来配对,Master 的 BrokerId 必须是 0,Slave 的BrokerId 必须是大与 0 的数。另外一个 Master 下面可以挂载多个 Slave,同一 Master 下的多个 Slave 通过指定不同的 BrokerId 来区分。
   除此之外,nameserver也需要集群。
下面以配置一主一备(同步),2个nameserver为例测试。
1、环境两台机器:
192.168.36.54  为主
192.168.36.189 为备
同时在2台机器个启动一个nameserver。安装RocketMq请参考:http://blog.csdn.net/zhu_tianwei/article/details/40948447
2.修改配置
(1)创建目录
mkdir   /home/rocket/alibaba-rocketmq/logs   #创建日志目录
mkdir -p /home/rocket/alibaba-rocketmq/data/store/commitlog  #创建数据存储目录
更改日志目录
cd /home/rocket/alibaba-rocketmq/conf
sed -i  's#${user.home}#${user.home}/alibaba-rocketmq#g' *.xml
(2)修改主配置
vi ./conf/2m-2s-sync/broker-a.properties
[plain] view plain copy print?在CODE上查看代码片派生到我的代码片
brokerClusterName=DefaultCluster 
brokerName=broker-a 
brokerId=0 
namesrvAddr=192.168.36.189:9876;192.168.36.54:9876 
defaultTopicQueueNums=4 
autoCreateTopicEnable=true 
autoCreateSubscriptionGroup=true 
listenPort=10911 
deleteWhen=04 
fileReservedTime=120 
mapedFileSizeCommitLog=1073741824 
mapedFileSizeConsumeQueue=50000000 
destroyMapedFileIntervalForcibly=120000 
redeleteHangedFileInterval=120000 
diskMaxUsedSpaceRatio=88 
 
storePathRootDir=/home/rocket/alibaba-rocketmq/data/store 
storePathCommitLog=/home/rocket/alibaba-rocketmq/data/store/commitlog 
maxMessageSize=65536 
flushCommitLogLeastPages=4 
flushConsumeQueueLeastPages=2 
flushCommitLogThoroughInterval=10000 
flushConsumeQueueThoroughInterval=60000 
 
checkTransactionMessageEnable=false 
sendMessageThreadPoolNums=128 
pullMessageThreadPoolNums=128 
 
brokerRole=SYNC_MASTER 
flushDiskType=ASYNC_FLUSH 
(3)修改备配置
[plain] view plain copy print?在CODE上查看代码片派生到我的代码片
brokerClusterName=DefaultCluster 
brokerName=broker-a 
brokerId=1 
namesrvAddr=192.168.36.189:9876;192.168.36.54:9876 
defaultTopicQueueNums=4 
autoCreateTopicEnable=true 
autoCreateSubscriptionGroup=true 
listenPort=10911 
deleteWhen=04 
fileReservedTime=120 
mapedFileSizeCommitLog=1073741824 
mapedFileSizeConsumeQueue=50000000 
destroyMapedFileIntervalForcibly=120000 
redeleteHangedFileInterval=120000 
diskMaxUsedSpaceRatio=88 
 
storePathRootDir=/home/rocket/alibaba-rocketmq/data/store 
storePathCommitLog=/home/rocket/alibaba-rocketmq/data/store/commitlog 
maxMessageSize=65536 
flushCommitLogLeastPages=4 
flushConsumeQueueLeastPages=2 
flushCommitLogThoroughInterval=10000 
flushConsumeQueueThoroughInterval=60000 
 
checkTransactionMessageEnable=false 
sendMessageThreadPoolNums=128 
pullMessageThreadPoolNums=128 
 
brokerRole=SLAVE 
flushDiskType=ASYNC_FLUSH 
3.启动两台nameserver
nohup ./bin/mqnamesrv >/dev/null 2>&1 &
4.启动2台broker
主:nohup ./bin/mqbroker -c ./conf/2m-2s-sync/broker-a.properties >/dev/null 2>&1 &
备:nohup ./bin/mqbroker -c ./conf/2m-2s-sync/broker-a-s.properties >/dev/null 2>&1 &
ok,配置启动完成。
实例:
1.消费者Consumer.java  ,采用主动拉取方式消费。
[java] view plain copy print?在CODE上查看代码片派生到我的代码片
package cn.slimsmart.rocketmq.demo.cluster; 
 
import java.util.HashMap; 
import java.util.List; 
import java.util.Map; 
import java.util.Set; 
 
import com.alibaba.rocketmq.client.consumer.DefaultMQPullConsumer; 
import com.alibaba.rocketmq.client.consumer.PullResult; 
import com.alibaba.rocketmq.client.exception.MQClientException; 
import com.alibaba.rocketmq.common.message.MessageExt; 
import com.alibaba.rocketmq.common.message.MessageQueue; 
 
//消费者 pull 
public class Consumer { 
 
    // Java缓存 
    private static final Map<MessageQueue, Long> offseTable = new HashMap<MessageQueue, Long>(); 
 
    /**
     * 主动拉取方式消费
     * 
     * @throws MQClientException
     */ 
    public static void main(String[] args) throws MQClientException { 
        /**
         * 一个应用创建一个Consumer,由应用来维护此对象,可以设置为全局对象或者单例<br>
         * 注意:ConsumerGroupName需要由应用来保证唯一 ,最好使用服务的包名区分同一服务,一类Consumer集合的名称,
         * 这类Consumer通常消费一类消息,且消费逻辑一致
         * PullConsumer:Consumer的一种,应用通常主动调用Consumer的拉取消息方法从Broker拉消息,主动权由应用控制
         */ 
        DefaultMQPullConsumer consumer = new DefaultMQPullConsumer("ConsumerGroupName"); 
        // //nameserver服务 
        consumer.setNamesrvAddr("192.168.36.189:9876;192.168.36.54:9876"); 
        consumer.setInstanceName("Consumber"); 
        consumer.start(); 
        // 拉取订阅主题的队列,默认队列大小是4 
        Set<MessageQueue> mqs = consumer.fetchSubscribeMessageQueues("TopicTest1"); 
        for (MessageQueue mq : mqs) { 
            System.out.println("Consume from the queue: " + mq); 
            SINGLE_MQ: while (true) { 
                try { 
                    PullResult pullResult = consumer.pullBlockIfNotFound(mq, null, getMessageQueueOffset(mq), 32); 
                    List<MessageExt> list = pullResult.getMsgFoundList(); 
                    if (list != null && list.size() < 100) { 
                        for (MessageExt msg : list) { 
                            System.out.println(new String(msg.getBody())); 
                        } 
                    } 
                    System.out.println(pullResult.getNextBeginOffset()); 
                    putMessageQueueOffset(mq, pullResult.getNextBeginOffset()); 
                    switch (pullResult.getPullStatus()) { 
                    case FOUND: 
                        break; 
                    case NO_MATCHED_MSG: 
                        break; 
                    case NO_NEW_MSG: 
                        break SINGLE_MQ; 
                    case OFFSET_ILLEGAL: 
                        break; 
                    default: 
                        break; 
                    } 
                } catch (Exception e) { 
                    e.printStackTrace(); 
                } 
            } 
        } 
        consumer.shutdown(); 
    } 
 
    private static void putMessageQueueOffset(MessageQueue mq, long offset) { 
        offseTable.put(mq, offset); 
    } 
 
    private static long getMessageQueueOffset(MessageQueue mq) { 
        Long offset = offseTable.get(mq); 
        if (offset != null) { 
            System.out.println(offset); 
            return offset; 
        } 
        return 0; 
    } 


2.生产者Producer.java ,TransactionMQProducer使用
[java] view plain copy print?在CODE上查看代码片派生到我的代码片
package cn.slimsmart.rocketmq.demo.cluster; 
 
import java.util.concurrent.TimeUnit; 
 
import com.alibaba.rocketmq.client.exception.MQClientException; 
import com.alibaba.rocketmq.client.producer.LocalTransactionExecuter; 
import com.alibaba.rocketmq.client.producer.LocalTransactionState; 
import com.alibaba.rocketmq.client.producer.SendResult; 
import com.alibaba.rocketmq.client.producer.TransactionCheckListener; 
import com.alibaba.rocketmq.client.producer.TransactionMQProducer; 
import com.alibaba.rocketmq.common.message.Message; 
import com.alibaba.rocketmq.common.message.MessageExt; 
 
//生产者 
public class Producer { 
 
    public static void main(String[] args) throws MQClientException, 
            InterruptedException { 
        /**
         * 一个应用创建一个Producer,由应用来维护此对象,可以设置为全局对象或者单例<br>
         * 注意:ProducerGroupName需要由应用来保证唯一,一类Producer集合的名称,这类Producer通常发送一类消息,且发送逻辑一致<br>
         * ProducerGroup这个概念发送普通的消息时,作用不大,但是发送分布式事务消息时,比较关键,
         * 因为服务器会回查这个Group下的任意一个Producer
         */ 
        final TransactionMQProducer producer = new TransactionMQProducer("ProducerGroupName"); 
        //nameserver服务 
        producer.setNamesrvAddr("192.168.36.189:9876;192.168.36.54:9876"); 
        producer.setInstanceName("Producer"); 
 
        /**
         * Producer对象在使用之前必须要调用start初始化,初始化一次即可<br>
         * 注意:切记不可以在每次发送消息时,都调用start方法
         */ 
        producer.start(); 
        //服务器回调Producer,检查本地事务分支成功还是失败 
        producer.setTransactionCheckListener( new TransactionCheckListener() { 
             
            public LocalTransactionState checkLocalTransactionState(MessageExt msg) { 
                System.out.println("checkLocalTransactionState --"+new String(msg.getBody())); 
                return LocalTransactionState.COMMIT_MESSAGE; 
            } 
        }); 
 
        /**
         * 下面这段代码表明一个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消息关键词,多个Key用KEY_SEPARATOR隔开(查询消息使用) 
                            ("Hello MetaQA").getBytes());// body 
                    SendResult sendResult = producer.sendMessageInTransaction(msg,new Producer().new MyTransactionExecuter(),"
$$$"); 
                    System.out.println(sendResult); 
                } 
 
                {    
                    Message msg = new Message("TopicTest2",// topic 
                            "TagB",// tag 
                            "OrderID0034",// key 消息关键词,多个Key用KEY_SEPARATOR隔开(查询消息使用) 
                            ("Hello MetaQB").getBytes());// body 
                    SendResult sendResult = producer.sendMessageInTransaction(msg,new Producer().new MyTransactionExecuter(),"
$$$"); 
                    System.out.println(sendResult); 
                } 
 
                { 
                    Message msg = new Message("TopicTest3",// topic 
                            "TagC",// tag 
                            "OrderID061",// key 
                            ("Hello MetaQC").getBytes());// body 
                    SendResult sendResult = producer.sendMessageInTransaction(msg,new Producer().new MyTransactionExecuter(),"
$$$"); 
                    System.out.println(sendResult); 
                } 
            } catch (Exception e) { 
                e.printStackTrace(); 
            } 
            TimeUnit.MILLISECONDS.sleep(1000); 
        } 
 
        /**
         * 应用退出时,要调用shutdown来清理资源,关闭网络连接,从MetaQ服务器上注销自己
         * 注意:我们建议应用在JBOSS、Tomcat等容器的退出钩子里调用shutdown方法
         */ 
        // producer.shutdown(); 
        Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() { 
            public void run() { 
                producer.shutdown(); 
            } 
        })); 
        System.exit(0); 
    } 
     
    //执行本地事务,由客户端回调 
    public class MyTransactionExecuter implements LocalTransactionExecuter{ 
 
        public LocalTransactionState executeLocalTransactionBranch(Message msg, Object arg) { 
            System.out.println("executeLocalTransactionBranch--msg="+new String(msg.getBody())); 
            System.out.println("executeLocalTransactionBranch--arg="+arg); 
            return LocalTransactionState.COMMIT_MESSAGE; 
        } 
         
    } 

实例代码:http://download.csdn.net/detail/tianwei7518/8138705
分享到:
评论

相关推荐

    cas-3.5.2改造源码

    cas-3.5.2改造源代码。 相关博客: http://blog.csdn.net/zhu_tianwei/article/details/19154891 http://blog.csdn.net/zhu_tianwei/article/details/19156169

    改造CAS单点登录部署文件

    CAS(Central Authentication Service)是Java开发的一个开源的单点登录(Single Sign-On,SSO)系统,主要用于统一认证管理,实现用户在多个应用系统中只需登录一次,即可访问所有相互信任的应用系统。...

    Windows下安装squid的步骤详解

    visible_hostname tianwei-itrus; # 设置代理服务器的可见主机名 cache_mem 64MB; # 设置代理服务器缓存大小为64MB # 定义各种安全端口 acl 命令 acl Safe_ports port 80 ... # 定义HTTP端口 acl Safe_...

    最新ip下载

    "tianwei_acc"可能是该工具的开发商或者产品名,版本号2.6.2.0185表示这是该软件的一个特定版本,可能包含了错误修复、性能提升或新功能的添加。 2. **PC6下载.url**:这通常是一个快捷方式文件,指向一个网站(PC6...

    一蓑烟雨论坛三年精华集(共42M二分卷)分卷二

    感谢帮助UnPacKcN和提供空间的朋友们:TiANWEi、声声慢、闪电狼、nbw、ryoada、小白、xuruifengc、Ivanov 感谢UnPacKcN曾经的版主们,UnPacKcN见证过你们的劳动:lipton、WiNrOOt 感谢UnPacKcN现在的版主们,你们为...

    一蓑烟雨论坛三年精华集(共42M二分卷)分卷一

    感谢帮助UnPacKcN和提供空间的朋友们:TiANWEi、声声慢、闪电狼、nbw、ryoada、小白、xuruifengc、Ivanov 感谢UnPacKcN曾经的版主们,UnPacKcN见证过你们的劳动:lipton、WiNrOOt 感谢UnPacKcN现在的版主们,你们为...

    msxml 6.10.1129.0

    在实际应用中,MSXML6可能会与多种技术结合使用,如ASP.NET、VBScript、JavaScript等,用于处理服务器端或客户端的XML数据。例如,在开发Web应用程序时,可以通过XMLHttpRequest对象利用MSXML来发送异步HTTP请求,...

    Typing Tutor with Dosbox 2021-5-24.rar

    本压缩包内收集是的Tianwei修复了千年虫问题的版本,可以在Win10 64位系统上运行,不过要启动Dosbox环境. 这里附上Tianwei当年修复tt的文章. 另外,为让TT正常运行,附带了DOSBOX运行环境,说明了如何安装运行,详细请...

    中心点

    , 尹天伟,周兴义,PhilippKrähenbühl, arXiv技术报告( ) @article{yin2020center, title={Center-based 3D Object Detection and Tracking}, author={Yin, Tianwei and Zhou, Xingyi and Kr{\"a}henb{\"u}hl...

    tt打字软件

    知乎韦易笑制作, 我搬砖.

    基于中心的3D对象检测和跟踪-Python开发

    基于中心的3D对象检测和跟踪,殷天伟,周兴义,PhilippKrähenbühl,arXiv技术报告(arXiv 2006.11275)@article {yin2020center,title = {基于中心的3D对象检测和跟踪},作者= {Yin,Tianwei和Zhou,Xingyi and ...

    全球倒极电渗析系统市场总体规模,前6强厂商排名及市场份额分析报告.docx

    8. **Shandong Tianwei** 9. **Jiangsu Ritai** 10. **Pure Water Group** #### 市场份额分布 根据2022年的统计数据,全球倒极电渗析系统的前五大厂商占据了大约59.0%的市场份额。这一比例反映了市场竞争格局的...

    Centerpoint_PC

    , 尹天伟,周兴义,PhilippKrähenbühl, arXiv技术报告( ) @article{yin2021center, title={Center-based 3D Object Detection and Tracking}, author={Yin, Tianwei and Zhou, Xingyi and Kr{\"a}henb{\"u}hl...

    KFNet:KFNet:使用卡尔曼滤波学习时间相机的重新定位(CVPR 2020口头)

    如果您发现此项目有用,请引用: @inproceedings{zhou2020kfnet, title={KFNet: Learning Temporal Camera Relocalization using Kalman Filtering}, author={Zhou, Lei and Luo, Zixin and Shen, Tianwei an

    动态效果很棒的欧美风格PPT模板.pptx

    "Microsoft PowerPoint Students University TIANWEI"这部分可能表明该模板适用于学生群体,特别是在校大学生。它可能包含一些适合学术报告、项目展示或课堂讲解的设计元素。 "Love Me Chinese Chongqing ELLIPSE...

    CenterPoint

    @article{yin2021center, title={Center-based 3D Object Detection and Tracking}, author={Yin, Tianwei and Zhou, Xingyi and Kr{\"a}henb{\"u}hl, Philipp}, journal={CVPR}, year={2021}, } 消息 [2021-02...

    STBA:通过运动进行随机捆绑调整以实现有效且可扩展的结构(ECCV 2020)

    author={Zhou, Lei and Luo, Zixin and Zhen, Mingmin and Shen, Tianwei and Li, Shiwei and Huang, Zhuofei and Fang, Tian and Quan, Long}, booktitle={European Conference on Computer Vision (ECCV)}, ...

    在WIN8.1上可以运行的经典DOS版TT

    TT即TypingTutorIV,一个非常经典的DOS游戏,用于练习英文打字极好,当年非常风靡的,这里找的是Tianwei修改过的克服了千年虫问题的版本,可以正常运行保存等。 在DOSBOX环境中运行,本包提供了配置方法。 请详见...

Global site tag (gtag.js) - Google Analytics