Docker安装RocketMQ
·
1、拉取镜像
docker pull foxiswho/rocketmq:4.7.0
2.、启动 NameServer
docker run -d \
-p 9876:9876 \
--name rmqserver \
foxiswho/rocketmq:4.7.0 \
sh mqnamesrv
命令注释
# 启动 RocketMQ NameServer 容器
docker run -d \ # -d 后台运行容器,不阻塞当前终端
-p 9876:9876 \ # 映射端口:宿主机 9876 -> 容器 9876
# 9876 是 NameServer 的默认端口
# Producer/Broker/Consumer 都通过这个端口发现路由
--name rmqserver \ # 给容器命名
foxiswho/rocketmq:4.7.0 \ # 使用的镜像
sh mqnamesrv # 容器启动后执行的命令
# sh mqnamesrv 启动 NameServer 进程
# mqnamesrv 是 RocketMQ 的名字服务程序
3、准备 broker.conf
注意要改 brokerIP1
brokerIP1 一定要改成你宿主机的真实 IP(不能是 127.0.0.1),否则客户端连不上。
sudo mkdir -p /etc/rocketmq
sudo tee /etc/rocketmq/broker.conf <<-'EOF'
# 所属集群名称,多 broker 时要一致
brokerClusterName=DefaultCluster
# broker 名称,同一 master/slave 组要相同
brokerName=broker-a
# 0 表示 Master,>0 表示 Slave
brokerId=0
# 每天几点删除过期 CommitLog 文件(24小时制)
deleteWhen=04
# 消息文件保留时间(小时),默认 48 小时
fileReservedTime=48
# broker 角色:
# ASYNC_MASTER 异步主(性能高,可能丢少量消息)
# SYNC_MASTER 同步主(可靠性高,性能略低)
# SLAVE 从节点
brokerRole=ASYNC_MASTER
# 刷盘方式:
# ASYNC_FLUSH 异步刷盘(快,宕机可能丢消息)
# SYNC_FLUSH 同步刷盘(慢,基本不丢)
flushDiskType=ASYNC_FLUSH
# Broker 对外暴露的 IP
# 必须是宿主机或容器可达 IP
# 不能是 127.0.0.1,否则客户端连不上
brokerIP1=xxxx
# Broker 监听端口(业务端口)
listenPort=10911
# NameServer 地址,多个用 ; 分隔
# 这里用容器名 rmqserver,因为 --link 映射过
namesrvAddr=rmqserver:9876
# 是否允许自动创建 Topic(生产建议关闭)
autoCreateTopicEnable=true
EOF
也可以直接复制粘贴下面的
sudo mkdir -p /etc/rocketmq
cat <<EOF | sudo tee /etc/rocketmq/broker.conf
brokerClusterName=DefaultCluster
brokerName=broker-a
brokerId=0
deleteWhen=04
fileReservedTime=48
brokerRole=ASYNC_MASTER
flushDiskType=ASYNC_FLUSH
brokerIP1=$(hostname -I | awk '{print $1}')
listenPort=10911
namesrvAddr=rmqserver:9876
autoCreateTopicEnable=true
EOF
4、 启动 Broker(依赖 rmqserver)
docker run -d \
-p 10911:10911 \
-p 10909:10909 \
--name rmqbroker \
--link rmqserver:namesrv \
-e "NAMESRV_ADDR=namesrv:9876" \
-e "JAVA_OPTS=-Duser.home=/opt" \
-e "JAVA_OPT_EXT=-server -Xms128m -Xmx128m" \
-v /etc/rocketmq/broker.conf:/etc/rocketmq/broker.conf \
foxiswho/rocketmq:4.7.0 \
sh mqbroker -c /etc/rocketmq/broker.conf
命令注释
# 启动 RocketMQ Broker 容器
docker run -d \ # -d 后台运行容器
-p 10911:10911 \ # 映射 Broker 业务端口(Producer/Consumer 连这个)
-p 10909:10909 \ # 映射 VIP 通道端口(用于高可用,可省略但建议保留)
--name rmqbroker \ # 容器名称,方便后续管理
--link rmqserver:namesrv \ # 连接 NameServer 容器,并给它取别名 namesrv
# 这样 Broker 容器里可以用 namesrv 这个主机名访问 NameServer
-e "NAMESRV_ADDR=namesrv:9876" \ # 设置环境变量:告诉 Broker NameServer 的地址
# 这里的 namesrv 对应上面 --link 设置的别名
-e "JAVA_OPTS=-Duser.home=/opt" \ # 设置 Java 系统属性:指定用户家目录为 /opt
# 主要是让 RocketMQ 把日志等文件写到 /opt 下
-e "JAVA_OPT_EXT=-server -Xms128m -Xmx128m" \ # 扩展 JVM 参数:
# -server:使用 Server 模式 JVM
# -Xms128m:初始堆内存 128MB
# -Xmx128m:最大堆内存 128MB
-v /etc/rocketmq/broker.conf:/etc/rocketmq/broker.conf \ # 挂载配置文件:
# 左边宿主机路径,右边容器内路径
# 这样改配置不用重建容器
foxiswho/rocketmq:4.7.0 \ # 使用的镜像名和版本标签
sh mqbroker -c /etc/rocketmq/broker.conf # 容器启动后执行的命令:
# sh mqbroker 启动 Broker 进程
# -c 指定配置文件路径
5、验证是否启动成功
docker logs -f rmqbroker

6、安装RocketMQ Console
docker pull styletang/rocketmq-console-ng
docker run -d \
--name rmqconsole \
-p 8180:8080 \
--link rmqserver:namesrv \
-e "JAVA_OPTS=-Drocketmq.namesrv.addr=namesrv:9876 -Dcom.rocketmq.sendMessageWithVIPChannel=false" \
styletang/rocketmq-console-ng
命令注释
# 第一步:拉取 RocketMQ Web 管理界面的镜像
docker pull styletang/rocketmq-console-ng
# 第二步:启动 RocketMQ Web 管理界面容器
docker run -d \ # -d 后台运行容器,不占用当前终端
--name rmqconsole \ # 容器名称,方便后续管理(stop/start/logs)
-p 8180:8080 \ # 端口映射:宿主机 8180 -> 容器内 8080
# 容器内 Console 默认在 8080 提供 Web 服务
# 映射到 8180 避免与其他服务端口冲突
# 浏览器访问 http://宿主机IP:8180
--link rmqserver:namesrv \ # 连接 NameServer 容器,并给它取别名 namesrv
# rmqserver:你之前启动的 NameServer 容器名
# namesrv:在 Console 容器里访问 NameServer 用的主机名
# 这样 Console 就知道去哪里查 Broker 的路由信息
-e "JAVA_OPTS=\ # 设置 Java 系统属性,传递给 Console 的 Spring Boot 应用
-Drocketmq.namesrv.addr=namesrv:9876 \ # 指定 NameServer 地址:namesrv:9876
# namesrv 来自上面 --link 设置的别名
# 9876 是 NameServer 的默认端口
# Console 通过这个地址连接 NameServer 获取集群信息
-Dcom.rocketmq.sendMessageWithVIPChannel=false" \ # 禁用 VIP 通道(端口 10909)
# 4.x Broker 默认开启 VIP 通道
# 但有些客户端/工具不支持
# 设为 false 后走普通端口 10911
# 不加这个,某些操作可能超时或失败
styletang/rocketmq-console-ng # 使用的镜像名称
# 这是一个第三方的 RocketMQ Web 管理界面
# 功能:查看集群状态、Topic 列表、消费者情况、消息轨迹等
浏览器访问http://宿主机ip:8180/#/message就会显示console页面,点击右上角可以切换成中文
7、代码测试
引入依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>
配置文件
server:
port: 6723
rocketmq:
# NameServer 地址(改成你的实际 IP)
name-server: xxxx:9876
producer:
# 生产者组名,随便起但不能为空
group: springboot-producer-group
生产者(发送消息)
package com.azure.controller;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class RocketMQProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 发送消息接口
* 浏览器访问:http://localhost:6723/send?msg=你好
*/
@GetMapping("/send")
public String sendMessage(@RequestParam(defaultValue = "Hello RocketMQ") String msg) {
// 发送消息到 TestTopic,tag 为 TagA
rocketMQTemplate.convertAndSend("TestTopic:TagA", msg);
return "发送成功:" + msg;
}
}
消费者(接收消息)
package com.azure.rocketmqdemo;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
topic = "TestTopic", // 监听的 Topic
selectorExpression = "*", // 监听所有 Tag,也可以指定 "TagA || TagB"
consumerGroup = "springboot-consumer-group" // 消费者组名
)
public class RocketMQConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
// 收到消息后的处理逻辑
System.out.println("========== 收到消息 ==========");
System.out.println("消息内容:" + message);
System.out.println("================================");
}
}
postman或者浏览器发送消息测试
url:http://localhost:6723/send?msg=你好
控制台打印消息
console页面也会显示,我是多点了一次,所以显示两条消息
更多推荐
所有评论(0)