ICode9

精准搜索请尝试: 精确搜索
首页 > 其他分享> 文章详细

Docker安装RocketMQ以及使用

2021-10-12 12:29:57  阅读:195  来源: 互联网

标签:-- broker dzp RocketMQ Docker 220.76 安装 rocketmq


目录

一、RocketMQ简介

二、docker安装RocketMQ

 三、java使用RocketMQ 


一、RocketMQ简介

Apache RocketMQ是一个分布式消息传递和流媒体平台,具有低延迟,高性能和可靠性, 万亿级容量和灵活的可伸缩性。 它由四个部分组成:nameserver,broker,生产者和使用者。

二、docker安装RocketMQ

1、搜索RocketMQ镜像

docker search rocketmq

2、启动NameServer

docker run -d -p 9876:9876 --name rmqserver  foxiswho/rocketmq:server-4.7.0

3、启动broker

编辑broker.conf文件

vim /home/rocketmq/broker.conf

内容为:

# 所属集群名称,如果节点较多可以配置多个
brokerClusterName = DefaultCluster
#broker名称,master和slave使用相同的名称,表明他们的主从关系
brokerName = broker-a
#0表示Master,大于0表示不同的slave
brokerId = 0
#表示几点做消息删除动作,默认是凌晨4点
deleteWhen = 04
#在磁盘上保留消息的时长,单位是小时
fileReservedTime = 48
#有三个值:SYNC_MASTER,ASYNC_MASTER,SLAVE;同步和异步表示Master和Slave之间同步数据的机制;
brokerRole = ASYNC_MASTER
#刷盘策略,取值为:ASYNC_FLUSH,SYNC_FLUSH表示同步刷盘和异步刷盘;SYNC_FLUSH消息写入磁盘后才返回成功状态,ASYNC_FLUSH不需要;
flushDiskType = ASYNC_FLUSH
# 设置broker节点所在服务器的ip地址
brokerIP1 = 192.168.220.76

执行命令:

docker run -d -p 10911:10911 -p 10909:10909\
 --name rmqbroker --link rmqserver:namesrv\
 --privileged=true\
 -e "NAMESRV_ADDR=192.168.220.76:9876" -e "JAVA_OPTS=-Duser.home=/opt"\
 -e "JAVA_OPT_EXT=-server -Xms128m -Xmx128m"\
 -v /home/rocketmq/broker.conf:/etc/rocketmq/broker.conf \
 foxiswho/rocketmq:broker-4.7.0

4、启动rocketmq console

docker run -d --name rmqconsole -p 8080:8080 --link rmqserver:namesrv\
 -e "JAVA_OPTS=-Drocketmq.namesrv.addr=192.168.220.76:9876\
 -Dcom.rocketmq.sendMessageWithVIPChannel=false"\
 -t styletang/rocketmq-console-ng

5、可视化页面,地址为 http://192.168.220.76:8080/

 

 三、java使用RocketMQ 

1、pom.xml中添加依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.0.4</version>
</dependency>

2.创建生产者

// 1 创建消息生产者,指定生成组名
DefaultMQProducer defaultMQProducer = new DefaultMQProducer("dzp-producer-group");
// 2 指定NameServer的地址
defaultMQProducer.setNamesrvAddr("192.168.220.76:9876");
// 3 启动生产者
defaultMQProducer.start();
// 4 构建消息对象,主要是设置消息的主题、标签、内容
Message message = new Message("dzp-topic", "dzp-tag", "dzp-key", ("dzp测试消息发送").getBytes());
// 5 发送消息
SendResult result = defaultMQProducer.send(message);
System.out.println("SendResult-->" + result);
// 6 关闭生产者
defaultMQProducer.shutdown();

3、创建消费者

// 1 创建消费者,指定所属的消费者组名
DefaultMQPushConsumer defaultMQPushConsumer = new DefaultMQPushConsumer("dzp-consumer-group");
// 2 指定NameServer的地址
defaultMQPushConsumer.setNamesrvAddr("192.168.220.76:9876");
// 3 指定消费者订阅的主题和标签
defaultMQPushConsumer.subscribe("dzp-topic", "*");
// 4 进行订阅:注册回调函数,编写处理消息的逻辑
defaultMQPushConsumer.registerMessageListener((List<MessageExt> list, ConsumeConcurrentlyContext context) -> {
    //并且返回ConsumeConcurrentlyStatus.RECONSUME_LATER
    try {
      for (MessageExt messageExt : list) {

          String topic = messageExt.getTopic();
          System.out.println("topic-->" + topic);

          String tags = messageExt.getTags();
          System.out.println("tags-->" + tags);

          String keys = messageExt.getKeys();
          System.out.println("keys-->" + keys);

          String body = new String(messageExt.getBody());
          System.out.println("body-->" + body);
        }

    } catch (Throwable throwable) {
        throwable.printStackTrace();
    }

       return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    });

    // 5 启动消费者
    defaultMQPushConsumer.start();
}

标签:--,broker,dzp,RocketMQ,Docker,220.76,安装,rocketmq
来源: https://blog.csdn.net/DZP_dream/article/details/120718547

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有