Skip to content

RocketMQ5.X版本API

本示例为普通消息发送。

版本5较之前的版本做了较大的改动,详细可以参考新旧版本比较一文。

并且版本5.x在订阅关系一致性中支持了在同一消费组的消费者可以订阅不同主题的消息(但是同一消费组下同一主题的订阅关系要求一致即同一消费者组同一主题的订阅过滤条件必须一致)。

为什么 5.x API 要先配置 ClientConfiguration

RocketMQ 5.x 新客户端把“连接到哪里”和“创建什么客户端”分开。ClientConfiguration 负责 endpoint、认证、超时等连接层配置;ProducerPushConsumerSimpleConsumer 才负责具体发送和消费。

如果不理解这个拆分,最常见的错误是把 NameServer 地址、Broker 地址、Proxy 地址混在一起。5.x Java Client 示例里常见的 setEndpoints("ip:9081") 通常指向 Proxy gRPC 端口,而不是 4.x 时代的 namesrvAddr=ip:9876

mermaid
flowchart TD
    A["ClientConfiguration<br/>连接配置"] --> B["Producer<br/>发送消息"]
    A --> C["PushConsumer<br/>监听消费"]
    A --> D["SimpleConsumer<br/>主动拉取"]
    B --> E["Proxy / Broker"]
    C --> E
    D --> E

调用流程图

mermaid
flowchart TD
    A["Spring Boot 启动"] --> B["创建 ClientConfiguration\n配置 endpoints"]
    B --> C["创建 Producer Bean"]
    B --> D["创建 PushConsumer Bean"]
    E["Controller 接收请求"] --> F["Service 构造 Message"]
    F --> G["Producer send"]
    G --> H["Broker 持久化消息"]
    H --> I["PushConsumer 根据订阅关系收到消息"]
    I --> J["MessageListener 处理业务"]
    J --> K["返回 ConsumeResult.SUCCESS"]

RocketMQ 5.x Java 客户端的使用重点是:先创建客户端配置,再创建生产者或消费者,生产者负责发送,消费者通过监听器处理消息。

引入依赖

引入RocketMQ 5.x的客户端依赖

xml
<!-- https://mvnrepository.com/artifact/org.apache.rocketmq/rocketmq-client-java -->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client-java</artifactId>
    <version>5.0.5</version>
</dependency>

引入springboot web依赖

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

创建生产者

用于向mq中发送消息。

java
@Component
@Slf4j
public class LpxProducer {
    private Producer producer;
    @PostConstruct
    public void init()
    {
        log.info("生产者初始化");
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        try {
            ClientConfigurationBuilder clientConfigurationBuilder = ClientConfiguration.newBuilder().setEndpoints("158.50.0.2:9081");
            ClientConfiguration clientConfiguration = clientConfigurationBuilder.build();
            //初始化Producer
            producer = provider.newProducerBuilder()
                    // 预绑定Topic(可不配置)
                    .setTopics("TestTopic","TestTopic_0")
                    // 设置通信配置(必须)
                    .setClientConfiguration(clientConfiguration)
                    .build();
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
    public Producer getProducer()
    {
        return producer;
    }
}

创建消费者

消费者1

创建一个在消费组为consumer_group_1的消费者1

java
@Component
@Slf4j
public class Consumer_1 {
    @PostConstruct
    public void init() {
        this.consume();
    }

    public void consume() {
        log.info("消费者_1创建");
        //创建消费者,消费组为consumer_group_1
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder clientConfigurationBuilder = ClientConfiguration.newBuilder().setEndpoints("158.50.0.2:9081");
        ClientConfiguration clientConfiguration = clientConfigurationBuilder.build();
        String topic = "TestTopic";
        FilterExpression filterExpression = new FilterExpression("tagA");
        try {
            provider.newPushConsumerBuilder()
                    .setClientConfiguration(clientConfiguration)
                    .setConsumerGroup("consumer_group_1")
                    .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                    .setMessageListener(messageView -> {
                        log.info("消费者1message body={}", messageView);
                        return ConsumeResult.SUCCESS;
                    }).build();
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

消费者2

创建一个在消费组为consumer_group_2的消费者2

java
@Component
@Slf4j
public class Consumer_2 {
    @PostConstruct
    public void init() {
        this.consume();
    }
    public void consume() {
        log.info("消费者_2创建");
        //创建消费者,消费组为consumer_group_2
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder clientConfigurationBuilder = ClientConfiguration.newBuilder().setEndpoints("158.50.0.2:9081");
        ClientConfiguration clientConfiguration = clientConfigurationBuilder.build();
        String topic = "TestTopic";
        FilterExpression filterExpression = new FilterExpression("tagA");
        try {
            provider.newPushConsumerBuilder()
                    .setClientConfiguration(clientConfiguration)
                    .setConsumerGroup("consumer_group_2")
                    .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                    .setMessageListener(messageView -> {
                        log.info("消费者2message body={}", messageView);
                        return ConsumeResult.SUCCESS;
                    }).build();
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

消费者3

创建一个在消费组为consumer_group_1的消费者3

java
@Component
@Slf4j
public class Consumer_3 {
    @PostConstruct
    public void init() {
        this.consume();
    }

    public void consume() {
        log.info("消费者_3创建");
        //创建消费者,消费组为consumer_group_1
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder clientConfigurationBuilder = ClientConfiguration.newBuilder().setEndpoints("158.50.0.2:9081");
        ClientConfiguration clientConfiguration = clientConfigurationBuilder.build();
        String topic = "TestTopic_0";
        FilterExpression filterExpression = new FilterExpression("*");
        try {
            provider.newPushConsumerBuilder()
                    .setClientConfiguration(clientConfiguration)
                    .setConsumerGroup("consumer_group_1")
                    .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                    .setMessageListener(messageView -> {
                        log.info("消费者3message body={}", messageView);
                        return ConsumeResult.SUCCESS;
                    }).build();
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

消息发送准备

在service层创建消息发送接口

java
public interface IMessageService {
    String sendMsg(JSONObject reqParams);
}

在service层创建消息发送接口实现类

java
@Service
@Slf4j
public class MessageImpl implements IMessageService {
    @Resource
    LpxProducer lpxProducer;
    
    @Override
    public String sendMsg(JSONObject reqParams) {
        String msg = reqParams.getStr("msg");
        String tag = reqParams.getStr("tag");
        String topic = reqParams.getStr("topic");
        MessageBuilder messageBuilder = new MessageBuilderImpl();
        Message message = messageBuilder.setTopic(topic).setBody(msg.getBytes()).setTag(tag).build();
        SendReceipt sendReceipt = null;
        try {
            sendReceipt = lpxProducer.getProducer().send(message);
            log.info("消息发送成功:{}", sendReceipt);
        } catch (Exception e) {
            log.info("消息发送失败");
            e.printStackTrace();
        }

        return "success";
    }
}

在controller层创建处理器

java
@RestController
@RequestMapping("msg")
public class MessageController {
    @Resource
    IMessageService messageService;
    
    @PostMapping("send")
    public String sendMsg(@RequestBody JSONObject reqParams) {
        return messageService.sendMsg(reqParams);
    }
}

条件分析

  1. 修改消费者1和消费者2属于不同消费者组,订阅相同主题(TestTopic)相同过滤标签为tagA

    请求信息:

    json
    {
        "msg": "哈哈哈",
        "tag": "tagA",
        "topic": "TestTopic"
    }

    结果: 消费者1和消费者2都会收到消息。

  2. 修改消费者1和消费者3属于同一消费者组,订阅不同主题不同的过滤标签

    请求信息1:

    json
    {
        "msg": "哈哈哈",
        "tag": "tagA",
        "topic": "TestTopic_0"
    }

    结果: 消费者3收到消息,消费者1不会收到消息。

    请求信息2:

    json
    {
        "msg": "哈哈哈",
        "tag": "tagA",
        "topic": "TestTopic"
    }

    结果: 消费者1收到消息,消费者3不会收到消息。

    总结: 5.x版本的订阅关系一致性包括了同属于同一消费者组的不同消费者订阅不同的主题和过滤标签。

  3. 修改消费者1和消费者3属于同一消费者组,订阅相同主题(TestTopic)相同过滤标签(tagA)

    请求信息:

    json
    {
        "msg": "哈哈哈",
        "tag": "tagA",
        "topic": "TestTopic"
    }

    结果: 消费者1和消费者3只有一个可以接收到消息。

    原因: 同一分组下的多个消费者将按照分组内统一的消费行为和负载均衡策略消费消息。若想消费者1和消费者3都收到消息则需要将消费者1和消费者3分配到不同的消费组。

  4. 修改消费者1和消费者3属于同一消费者组,订阅相同主题(TestTopic)不同的过滤标签(消费者1过滤标签tagA,消费者3过滤标签tagB)

    请求信息:

    json
    {
        "msg": "哈哈哈",
        "tag": "tagA",
        "topic": "TestTopic"
    }

    结果: 消息丢失

    原因: 订阅关系不一致。详情请参考订阅关系

常见风险和排查

现象常见原因排查方向
Producer 初始化失败endpoints 地址或端口错误确认连的是 Proxy gRPC 端口,例如 8081/9081
消息发送成功但消费者收不到Topic、Tag、ConsumerGroup 或订阅关系不匹配打印 Topic、Tag、Key、消费组
同一组两个消费者只有一个收到消息同组负载均衡是正常行为如果都要收到,放到不同 ConsumerGroup
同一组同一 Topic 过滤条件不同订阅关系不一致同一组同一 Topic 保持相同 FilterExpression
业务重复执行消费失败重试或客户端重新投递用业务唯一键做幂等

完整理解 Demo 的关键点

这篇 Demo 不是只为了“发一条消息”。它要让你理解三条规则:

  1. 不同 ConsumerGroup 订阅同一 Topic,同一条消息会被各组各消费一次。
  2. 同一 ConsumerGroup 的多个实例是组内负载均衡,不是每个实例都收到一份。
  3. 同一 ConsumerGroup 下,同一 Topic 的过滤条件必须保持一致,否则会出现订阅关系问题。

如果不会这三条规则,实际项目里最容易把“负载均衡”误判成“消息丢了”,或者把“订阅不一致”误判成“RocketMQ 不稳定”。