RocketMQ怎么集成?

wen python案例 1

本文目录导读:

RocketMQ怎么集成?

  1. 核心概念与准备工作
  2. 集成方式 1:使用 Spring Boot + RocketMQ 官方 Starter (推荐)
  3. 集成方式 2:使用 RocketMQ 原生 Java API (不依赖 Spring)
  4. 关键配置与最佳实践
  5. 常见问题排查

RocketMQ 的集成方式主要取决于你使用的开发语言应用场景(比如是 Spring Boot 项目、纯 Java 项目、还是其他语言)。

下面我将以 Java(最主流)Spring Boot(最常用框架) 为核心,详细说明集成步骤和配置要点。


核心概念与准备工作

在开始集成前,确保 RocketMQ 服务端已经部署并运行,你需要知道以下信息:

  • NameServer 地址168.1.100:9876
  • Topic 名称:生产者发送和消费者订阅的主题。
  • 消费者组 (Consumer Group):消费同一类消息的消费者实例集合。
  • 生产者组 (Producer Group):通常用于事务消息的回查,普通消息可以不特别关注。

集成方式 1:使用 Spring Boot + RocketMQ 官方 Starter (推荐)

这是目前 Java 生态中最简单、最主流的方式,官方提供了 rocketmq-spring-boot-starter

添加 Maven 依赖

在你的 pom.xml 中添加:

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.3</version> <!-- 请使用最新稳定版 -->
</dependency>

配置 application.yml

在配置文件中配置 NameServer 地址:

rocketmq:
  # NameServer 地址,多个地址用分号隔开
  name-server: 192.168.1.100:9876
  # 生产者配置(可选,但推荐配置)
  producer:
    # 生产者组名,用于区分不同业务的生产者
    group: my-producer-group
    # 发送消息超时时间,毫秒
    send-message-timeout: 3000
    # 消息体最大字节数
    max-message-size: 4096
  # 消费者配置(部分在代码注解中配置)
  consumer:
    # 默认的消费者线程数
    listenr-thread-num: 20

编写生产者(发送消息)

注入 RocketMQTemplate,直接调用方法发送消息。

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
@Service
public class OrderProducer {
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    /**
     * 发送普通消息
     */
    public void sendMessage(String topic, String messageContent) {
        rocketMQTemplate.convertAndSend(topic, messageContent);
    }
    /**
     * 发送同步消息(带返回结果)
     */
    public boolean sendSyncMessage(String topic, String messageContent) {
        // SendResult 可以获取发送状态、Queue ID 等
        rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(messageContent).build());
        return true;
    }
    /**
     * 发送带 Key 的消息(便于去重或定位)
     */
    public void sendKeyMessage(String topic, String key, String messageContent) {
        org.springframework.messaging.Message<String> message = MessageBuilder
                .withPayload(messageContent)
                .setHeader("KEYS", key) // 设置消息 Key
                .build();
        rocketMQTemplate.syncSend(topic, message);
    }
}

编写消费者(接收消息)

使用 @RocketMQMessageListener 注解指定消费的主题和消费者组,实现 RocketMQListener 接口。

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
 * 消费者
 * topic: 要订阅的主题
 * consumerGroup: 消费者组名(必须唯一,不同业务用不同组)
 * selectorExpression: 标签过滤,默认 "*" 表示接收所有标签
 */
@Component
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer-group",
    selectorExpression = "*" 
)
public class OrderConsumerListener implements RocketMQListener<String> {
    @Override
    public void onMessage(String message) {
        // 注意:默认情况下,onMessage 执行完后自动返回 CONSUME_SUCCESS
        // 如果抛出异常,则自动返回 RECONSUME_LATER(重试)
        System.out.println("收到消息: " + message);
        // 处理业务逻辑...
    }
}

集成方式 2:使用 RocketMQ 原生 Java API (不依赖 Spring)

如果你的项目不是 Spring Boot,或者需要更底层的控制,可以使用原生 API。

添加 Maven 依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>4.9.6</version> <!-- 使用最新版 -->
</dependency>

生产者代码

import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SimpleProducer {
    public static void main(String[] args) throws Exception {
        // 1. 创建生产者,指定生产者组
        DefaultMQProducer producer = new DefaultMQProducer("my-producer-group");
        // 2. 设置 NameServer 地址
        producer.setNamesrvAddr("192.168.1.100:9876");
        // 3. 启动生产者
        producer.start();
        // 4. 创建消息对象 (topic, tags, keys, body)
        Message msg = new Message("order-topic", "tagA", "orderId_123", "Hello RocketMQ".getBytes());
        // 5. 发送消息
        SendResult sendResult = producer.send(msg);
        System.out.printf("发送结果: %s%n", sendResult);
        // 6. 关闭生产者
        producer.shutdown();
    }
}

消费者代码

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
public class SimpleConsumer {
    public static void main(String[] args) throws Exception {
        // 1. 创建消费者,指定消费者组
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("my-consumer-group");
        // 2. 设置 NameServer
        consumer.setNamesrvAddr("192.168.1.100:9876");
        // 3. 订阅主题和标签
        consumer.subscribe("order-topic", "*");
        // 4. 注册消息监听器
        consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            for (MessageExt msg : msgs) {
                System.out.printf("收到消息: %s %n", new String(msg.getBody()));
            }
            // 消费成功
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            // 消费失败,稍后重试
            // return ConsumeConcurrentlyStatus.RECONSUME_LATER;
        });
        // 5. 启动消费者
        consumer.start();
        System.out.println("消费者启动成功");
    }
}

关键配置与最佳实践

消息类型

消息类型 Spring Starter 方法 原生 API 方法
普通消息 syncSend, asyncSend, oneWaySend producer.send()
顺序消息 syncSendOrderly producer.send(msg, queueSelector, arg)
事务消息 需要实现 RocketMQLocalTransactionListener TransactionMQProducer
延时消息 设置 MessageDELAY 属性 msg.setDelayTimeLevel(3) (1s/5s/10s...)

重试与死信队列

  • 消费重试:默认情况下,如果消费者 onMessage 抛出异常,RocketMQ 会重试 16 次(间隔时间递增),超过后进入死信队列(DLQ)。
  • 死信队列:主题为 %DLQ%${consumerGroup},需要手动处理死信消息(例如重新投递或报警)。

幂等性

RocketMQ 不保证消息只被消费一次(可能重复投递),业务侧必须实现幂等:

  • 去重表:在数据库建立唯一索引(如业务 ID + 消息 ID)。
  • Redis 去重:消费前先检查 Redis 中是否已处理。
  • 业务幂等操作:例如使用 insert ... on duplicate key update

配置建议

  • NameServer 地址:不要写死 IP,建议使用域名或 Nginx 做负载均衡。
  • 消费者组:不同业务逻辑使用不同的消费者组,允许独立消费同一条消息。
  • 标签 (Tag):一个 Topic 下可以用 Tag 区分不同子类型(如 TagA = 下单,TagB = 取消订单),消费者可只订阅需要的 Tag。
  • 消息体大小:官方建议不超过 4MB,避免网络压力和 GC 问题。

常见问题排查

  1. 连接超时:检查防火墙是否开放 9876(NameServer 端口)和 10911(Broker 端口)。
  2. 生产者发送失败:检查 NameServer 地址是否正确,Broker 是否可用。
  3. 消费者收不到消息:检查 Topic 是否存在,Consumer Group 是否与其他程序冲突(同一组内消息会负载均衡)。
  4. 消息重复消费:检查业务代码的幂等处理是否生效。

对于大多数 Java 项目,推荐使用 rocketmq-spring-boot-starter,它能让你用几行注解和配置就完成集成,对于非 Spring 项目或需要精细控制的场景,原生 API 也很清晰,无论哪种方式,业务侧的幂等性处理是必不可少的

抱歉,评论功能暂时关闭!