1.目录结构

├── broker.conf
├── data
│   ├── broker
│   └── namesrv
│       └── logs
├── docker-compose.yml

2.docker-compsoe.yml

version: '3.8'

services:
  # RocketMQ 名称服务 (NameServer)
  namesrv:
    image: apache/rocketmq:4.9.4
    container_name: rocketmq-namesrv
    command: sh mqnamesrv
    ports:
      - 9876:9876
    volumes:
      - ./data/namesrv/logs:/home/rocketmq/logs
    environment:
      - JAVA_OPT_EXT=-Xms512m -Xmx512m -Xmn256m
    restart: unless-stopped
    networks:
      - rocketmq_net

  # RocketMQ 代理服务 (Broker)
  broker:
    image: apache/rocketmq:4.9.4
    container_name: rocketmq-broker
    command: sh mqbroker -n namesrv:9876 -c ../conf/broker.conf
    ports:
      - 10909:10909
      - 10911:10911
      - 10912:10912
    volumes:
      - ./data/broker/store:/home/rocketmq/store
      - ./broker.conf:/home/rocketmq/rocketmq-4.9.4/conf/broker.conf
    environment:
      - JAVA_OPT_EXT=-Xms1g -Xmx1g -Xmn512m
      - NAMESRV_ADDR=namesrv:9876
    depends_on:
      - namesrv
    restart: unless-stopped
    networks:
      - rocketmq_net

  # RocketMQ 控制台 (Dashboard)
  dashboard:
    image: apacherocketmq/rocketmq-dashboard:latest
    container_name: rocketmq-dashboard
    ports:
      - 8080:8080
    environment:
      - JAVA_OPTS=-Drocketmq.namesrv.addr=namesrv:9876
    depends_on:
      - namesrv
    restart: unless-stopped
    networks:
      - rocketmq_net

networks:
  rocketmq_net:
    driver: bridge

3.broker.conf

# 所属集群名称
brokerClusterName = DefaultCluster

# broker 名称
brokerName = broker-a

# 0 表示 Master,>0 表示 Slave
brokerId = 0

# 删除文件时间点,默认凌晨4点
deleteWhen = 04

# 文件保留时间,默认48小时
fileReservedTime = 48

# broker 角色
brokerRole = ASYNC_MASTER

# 刷盘方式
flushDiskType = ASYNC_FLUSH

# NameServer地址
namesrvAddr = namesrv:9876

# 存储路径
storePathRootDir = /home/rocketmq/store

# commitLog 存储路径
storePathCommitLog = /home/rocketmq/store/commitlog

# 消费队列存储路径
storePathConsumeQueue = /home/rocketmq/store/consumequeue

# 消息索引存储路径
storePathIndex = /home/rocketmq/store/index

# checkpoint 文件存储路径
storeCheckpoint = /home/rocketmq/store/checkpoint

# abort 文件存储路径
abortFile = /home/rocketmq/store/abort

# Broker IP 地址(如果docker容器有独立IP可以不用设置)
brokerIP1 = 192.168.137.101

# 允许自动创建Topic
autoCreateTopicEnable = true

# 允许自动创建订阅组
autoCreateSubscriptionGroup = true

listenPort = 10911

4.启动

CONTAINER ID   IMAGE                                      COMMAND                   CREATED          STATUS                          PORTS                                                                                                                                NAMES
22eec2e15edd   apache/rocketmq:4.9.4                      "sh mqbroker -n name…"   22 minutes ago   Up 22 minutes                   0.0.0.0:10909->10909/tcp, [::]:10909->10909/tcp, 9876/tcp, 0.0.0.0:10911-10912->10911-10912/tcp, [::]:10911-10912->10911-10912/tcp   rocketmq-broker
d6ccf0ab2297   apacherocketmq/rocketmq-dashboard:latest   "sh -c 'java $JAVA_O…"   22 minutes ago   Up 22 minutes                   0.0.0.0:8080->8080/tcp, [::]:8080->8080/tcp                                                                                          rocketmq-dashboard
26fa0b809f73   apache/rocketmq:4.9.4                      "sh mqnamesrv"            22 minutes ago   Up 22 minutes                   10909/tcp, 0.0.0.0:9876->9876/tcp, [::]:9876->9876/tcp, 10911-10912/tcp                                                              rocketmq-namesrv

5.容器内部验证是否成功

1.进入容器

docker exec -u root -it rocketmq-broker /bin/bash

2.生产者发送消息

./bin/tools.sh org.apache.rocketmq.example.quickstart.Producer

3.生产者发送消息成功示例

SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D3A0000, offsetMsgId=C0A8896500002A9F0000000000000000, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=2], queueOffset=0]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D460001, offsetMsgId=C0A8896500002A9F00000000000000BE, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=3], queueOffset=0]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D490002, offsetMsgId=C0A8896500002A9F000000000000017C, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=0], queueOffset=0]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D4B0003, offsetMsgId=C0A8896500002A9F000000000000023A, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=1], queueOffset=0]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D4D0004, offsetMsgId=C0A8896500002A9F00000000000002F8, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=2], queueOffset=1]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D4F0005, offsetMsgId=C0A8896500002A9F00000000000003B6, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=3], queueOffset=1]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D520006, offsetMsgId=C0A8896500002A9F0000000000000474, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=0], queueOffset=1]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D550007, offsetMsgId=C0A8896500002A9F0000000000000532, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=1], queueOffset=1]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D570008, offsetMsgId=C0A8896500002A9F00000000000005F0, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=2], queueOffset=2]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D5A0009, offsetMsgId=C0A8896500002A9F00000000000006AE, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=3], queueOffset=2]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D5C000A, offsetMsgId=C0A8896500002A9F000000000000076C, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=0], queueOffset=2]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D61000B, offsetMsgId=C0A8896500002A9F000000000000082B, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=1], queueOffset=2]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D65000C, offsetMsgId=C0A8896500002A9F00000000000008EA, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=2], queueOffset=3]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D6A000D, offsetMsgId=C0A8896500002A9F00000000000009A9, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=3], queueOffset=3]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D6E000E, offsetMsgId=C0A8896500002A9F0000000000000A68, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=0], queueOffset=3]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D71000F, offsetMsgId=C0A8896500002A9F0000000000000B27, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=1], queueOffset=3]
SendResult [sendStatus=SEND_OK, msgId=7F000001016E3A71F4DD55052D740010, offsetMsgId=C0A8896500002A9F0000000000000BE6, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=2], queueOffset=4]

4.消费者消费消息命令

./bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer

5.消费者消费消息成功示例

Consumer Started.
ConsumeMessageThread_please_rename_unique_group_name_4_4 Receive New Messages: [MessageExt [brokerName=broker-a, queueId=3, storeSize=192, queueOffset=32, sysFlag=0, bornTimestamp=1755432802909, bornHost=/172.18.0.1:41344, storeTimestamp=1755432802910, storeHost=/192.168.137.101:10911, msgId=C0A8896500002A9F0000000000006052, commitLogOffset=24658, bodyCRC=1273902952, reconsumeTimes=0, preparedTransactionOffset=0, toString()=Message{topic='TopicTest', flag=0, properties={MIN_OFFSET=0, MAX_OFFSET=250, CONSUME_START_TIME=1755432937577, UNIQ_KEY=7F000001016E3A71F4DD55052E5D0081, CLUSTER=DefaultCluster, TAGS=TagA}, body=[72, 101, 108, 108, 111, 32, 82, 111, 99, 107, 101, 116, 77, 81, 32, 49, 50, 57], transactionId='null'}]] 
ConsumeMessageThread_please_rename_unique_group_name_4_7 Receive New Messages: [MessageExt [brokerName=broker-a, queueId=3, storeSize=192, queueOffset=35, sysFlag=0, bornTimestamp=1755432802933, bornHost=/172.18.0.1:41344, storeTimestamp=1755432802934, storeHost=/192.168.137.101:10911, msgId=C0A8896500002A9F0000000000006952, commitLogOffset=26962, bodyCRC=326047964, reconsumeTimes=0, preparedTransactionOffset=0, toString()=Message{topic='TopicTest', flag=0, properties={MIN_OFFSET=0, MAX_OFFSET=250, CONSUME_START_TIME=1755432937578, UNIQ_KEY=7F000001016E3A71F4DD55052E75008D, CLUSTER=DefaultCluster, TAGS=TagA}, body=[72, 101, 108, 108, 111, 32, 82, 111, 99, 107, 101, 116, 77, 81, 32, 49, 52, 49], transactionId='null'}]] 
ConsumeMessageThread_please_rename_unique_group_name_4_1 Receive New Messages: [MessageExt [brokerName=broker-a, queueId=1, storeSize=192, queueOffset=59, sysFlag=0, bornTimestamp=1755432803077, bornHost=/172.18.0.1:41344, storeTimestamp=1755432803077, storeHost=/192.168.137.101:10911, msgId=C0A8896500002A9F000000000000B2D2, commitLogOffset=45778, bodyCRC=1353955440, reconsumeTimes=0, preparedTransactionOffset=0, toString()=Message{topic='TopicTest', flag=0, properties={MIN_OFFSET=0, MAX_OFFSET=250, CONSUME_START_TIME=1755432937577, UNIQ_KEY=7F000001016E3A71F4DD55052F0500EF, CLUSTER=DefaultCluster, TAGS=TagA}, body=[72, 101, 108, 108, 111, 32, 82, 111, 99, 107, 101, 116, 77, 81, 32, 50, 51, 57], transactionId='null'}]] 

6.java程序验证

1.导入依赖

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

2.启动类加上包扫描路径

@SpringBootApplication(scanBasePackages = {"com.cruise.billiards", "org.apache.rocketmq"})

3.application.yml

rocketmq:
  name-server: 192.168.137.101:9876 # Name Server地址,客户端通过这个地址连接到RocketMQ的Name Server
  producer:
    group: DefaultCluster  # 生产者组名称,用于标识一组相关的生产者
    send-message-timeout: 3000 # 消息发送超时时间,单位毫秒
    compress-message-body-threshold: 4096 # 消息体压缩阈值,超过此大小的消息体将被压缩
    retry-times-when-send-failed: 2 # 同步发送消息失败时的重试次数
    retry-next-server: true # 是否在发送失败时尝试下一个Broker
    max-message-size: 4194304 # 最大消息大小,单位字节(4MB)
    vip-channel-enabled: false # 是否启用VIP通道,通常在云环境下需要设置为false
    retry-times-when-send-async-failed: 2 # 异步发送消息失败时的重试次数
    send-latency-fault-enable: true # 是否启用发送延迟故障检测机制
  consumer:
    group: DefaultCluster  # 消费者组名称,用于标识一组相关的消费者
    consume-thread-min: 20 # 消费线程池最小线程数
    consume-thread-max: 64 # 消费线程池最大线程数
    consume-message-orderly: false # 是否顺序消费消息,false表示并发消费
    pull-batch-size: 32 # 单次拉取消息的最大批量大小
    consume-concurrently-max-span: 2000  # 并发消费时允许的最大跨度
    pull-interval: 0 # 拉取消息的时间间隔,0表示立即拉取
    consume-timeout: 15000 # 消费超时时间,单位毫秒

4.创建测试方法

package com.cruise.billiards.mq;

import jakarta.annotation.PostConstruct;
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.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.client.discovery.EnableDiscoveryClient;

import java.util.List;

@SpringBootApplication(scanBasePackages = {"com.cruise.billiards", "org.apache.rocketmq"})
@EnableDiscoveryClient
public class CruiseBilliardsMqApplication {

    public static void main(String[] args) {
        SpringApplication.run(CruiseBilliardsMqApplication.class, args);


    }

    @PostConstruct
    public void rocketMQTest() {
        try {
            // 创建生产者
            DefaultMQProducer producer = new DefaultMQProducer("DefaultCluster");
            producer.setNamesrvAddr("192.168.137.101:9876"); // 根据实际配置修改
            producer.start();

            // 创建消费者
            DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("DefaultCluster");
            consumer.setNamesrvAddr("192.168.137.101:9876"); // 根据实际配置修改
            consumer.subscribe("TopicTest", "*");
            consumer.registerMessageListener(new MessageListenerConcurrently() {
                @Override
                public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                                ConsumeConcurrentlyContext context) {
                    for (MessageExt msg : msgs) {
                        System.out.printf("收到消息: %s %n", new String(msg.getBody()));
                    }
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                }
            });
            consumer.start();

            // 发送测试消息
            Message msg = new Message("TopicTest", "testTag", "Hello RocketMQ from Spring Boot App".getBytes());
            SendResult sendResult = producer.send(msg);
            System.out.printf("发送结果: %s %n", sendResult);

            // 等待消费消息
            Thread.sleep(3000);

            // 关闭资源
            producer.shutdown();
            consumer.shutdown();

        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

5.验证

控制台打印如下报文代表成功

发送结果: SendResult [sendStatus=SEND_OK, msgId=7F0000013D5063947C6B56D248080000, offsetMsgId=C0A8896500002A9F000000000002F08C, messageQueue=MessageQueue [topic=TopicTest, brokerName=broker-a, queueId=1], queueOffset=250] 
收到消息: Hello RocketMQ from Spring Boot App 
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐