本节要点

  • SpringBoot集成MQ
  • 生产者和消费者
  • 自动ACK
  • 手动ACK(可选)

1、RocketMQ消息发送接收(重点)

生产者发送类型三种

  • 同步发送:同步发送会阻塞当前线程,直到收到 RocketMQ Broker 的确认(或超时)。这能确保你知道消息是否成功送达,但会牺牲一些性能。

  • **异步发送:**默认以异步方式发送消息(当 producer.sync=false 时)。这意味着 send 方法调用后会立即返回,发送结果由 RocketMQ Binder 在后台处理。

  • **单向发送:**只负责发送,不管消息是否发送成功。

    消费者消费消息分两种

  • **拉模式:**消费者主动去 Broker 上拉取消息。

  • **推模式:**消费者等待 Broker 把消息推送过来

使用方式:

  • 命令行
  • api调用
  • spring boot
  • spring cloud alibaba

2、RocketMQ集成SpringBoot

1、生产者

1)加依赖

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

2)加配置

rocketmq:
  name-server: 127.0.0.1:9876  # RocketMQ命名服务器地址
  producer:
    group: my-producer-group   # 生产者组名
    send-message-timeout: 3000 # 发送消息超时时间,毫秒

3)编写发送对象(可选)

也可以发送一个字符串,但一般会用对象来发送,直接就是json到mq里

package com.dw.message;

import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;

import java.io.Serializable;

@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class TestMessage implements Serializable {
    private String msgId;
    private String msgContent;
}

4)编写生产者bean(service)

实际上就是用RocketMQTemplate发送

package com.dw.producer;

import com.dw.model.message.TestMessage;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class MqProducer {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 同步发送消息
     */
    public SendResult sendSyncMessage(String topic, String tag, TestMessage message) {
        // 同步发送消息
        String destination = topic + ":" + tag;
        return rocketMQTemplate.syncSend(destination, MessageBuilder.withPayload(message).build());
    }


    /**
     * 同步发送延时消息
     */
    public SendResult sendSyncDelayMessage(String topic, String message, int delayLevel) {
        return rocketMQTemplate.syncSend(topic,
                MessageBuilder.withPayload(message).build(),
                rocketMQTemplate.getProducer().getSendMsgTimeout(),
                delayLevel);
    }

    /**
     * 异步发送消息
     */
    public void sendAsyncMessage(String topic,String tag, TestMessage message) {
        String destination = topic + ":" + tag;
        rocketMQTemplate.asyncSend(destination, MessageBuilder.withPayload(message).build(), new SendCallback() {
      
            @Override
            public void onSuccess(SendResult sendResult) {
                System.out.println("异步消息发送成功: " + sendResult);
            }

            @Override
            public void onException(Throwable throwable) {
                System.out.println("异步消息发送失败: " + throwable.getMessage());
                // 处理发送失败逻辑
            }
        });
    }

    /**
     * 发送单向消息(不关心发送结果)
     */
    public void sendOnewayMessage(String topic, String message) {
        rocketMQTemplate.sendOneWay(topic, MessageBuilder.withPayload(message).build());
    }
  
  
    /**
     * 顺序发送
     */
    public void sendSyncOrderlyMessage(String topic, TestMessage message) {
        rocketMQTemplate.asyncSendOrderly(topic,
                MessageBuilder.withPayload(message).build(),
                "orderly",  // 这里设置hashKey - 使用业务ID如订单号 接收者也需要设置一样
                new SendCallback() {
                    @Override
                    public void onSuccess(SendResult sendResult) {
                        System.out.println("sendSyncOrderlyMessage-顺序消息-发送成功---");
                    }

                    @Override
                    public void onException(Throwable e) {
                        e.printStackTrace();
                    }
                });
    }
}
延迟级别 延迟时间 延迟级别 延迟时间
1 1秒 10 6分钟
2 5秒 11 7分钟
3 10秒 12 8分钟
4 30秒 13 9分钟
5 1分钟 14 10分钟
6 2分钟 15 20分钟
7 3分钟 16 30分钟
8 4分钟 17 1小时
9 5分钟 18 2小时

broker.conf 默认

# 默认配置,共18个等级
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

# 自定义配置示例:增加一个2天的等级,总共变成19个等级
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h 2d

5)编写controller

package com.dw.controller;

import com.dw.model.message.TestMessage;
import com.dw.producer.MqProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/product")
public class MqTestController {
    @Autowired
    MqProducer mqProducer;

    /**
     * 发送同步消息
     * @return
     */
    @RequestMapping("/send/sync")
    public String sendSyncMessage() {
        TestMessage message = TestMessage.builder().msgId("1").msgContent("hello world").build();
        String str = mqProducer.sendSyncMessage("topic-test","Tag-A", message).toString();
        System.out.println(str);
        return "发送成功---"+str;
    }

    /**
     * 发送异步消息
     * @return
     */
    @RequestMapping("/send/async")
    public String sendAsyncMessage() {
        TestMessage message = TestMessage.builder().msgId("2").msgContent("嘻嘻哈哈").build();
        mqProducer.sendAsyncMessage("topic-test","Tag-B",message);
        return "发送成功";
    }

    /**
     * 发送一次性消息
     * @return
     */
    @RequestMapping("/send/oneway")
    public String sendOnewayMessage() {
        return "发送成功";
    }
}