RocketMq应用实践

一、使用rocketmq-spring-boot-starter组件

1、添加依赖


<!-- rocketmq依赖 -->

<dependency>

  <groupId>org.apache.rocketmq</groupId>

  <artifactId>rocketmq-spring-boot-starter</artifactId>

  <version>2.2.2</version>

</dependency>

2、添加配置

rocketmq:
  nameServer: xx.xx.xx.xx:9876
  producer:
    group: test-group # 必须指定group
    send-message-timeout: 3000 # 消息发送超时时长,默认3s
    retry-times-when-send-failed: 3 # 同步发送消息失败重试次数,默认2
    retry-times-when-send-async-failed: 3 # 异步发送消息失败重试次数,默认2
  consumer:
    group: test-group
  topic: test-topic

3、配置生产者

常用方法汇总,根据实际需要进行封装

import com.alibaba.fastjson.JSON;
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.beans.factory.annotation.Value;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class MQProducerService {

    @Value("${rocketmq.producer.send-message-timeout}")
    private Integer messageTimeOut;

    // 直接注入使用,用于发送消息到broker服务器
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 发送同步消息,异步消费(阻塞当前线程,等待broker响应发送结果,这样不太容易丢失消息)
     * (msgBody也可以是对象,sendResult为返回的发送结果)
     */
    public SendResult sendMsg(String topic,String msgBody) {
        SendResult sendResult = rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(msgBody).build());
        return sendResult;
    }
    public SendResult sendKeyMsg(String topic,String msgBody,String key){
        SendResult sendResult = rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(msgBody).setHeader("KEYS",key).build());
        return sendResult;
    }

    /**
     * 发送异步消息(通过线程池执行发送到broker的消息任务,执行完后回调:在SendCallback中可处理相关成功失败时的逻辑)
     * (适合对响应时间敏感的业务场景)
     */
    public void sendAsyncMsg(String topic,String msgBody) {
        rocketMQTemplate.asyncSend(topic, MessageBuilder.withPayload(msgBody).build(), new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                // 处理消息发送成功逻辑
            }
            @Override
            public void onException(Throwable throwable) {
                // 处理消息发送异常逻辑
            }
        });
    }
    public void sendAsyncKeyMsg(String topic,String msgBody,String key) {
        rocketMQTemplate.asyncSend(topic, MessageBuilder.withPayload(msgBody).setHeader("KEYS",key).build(), new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                // 处理消息发送成功逻辑
            }
            @Override
            public void onException(Throwable throwable) {
                // 处理消息发送异常逻辑
            }
        });
    }

    /**
     * 发送延时消息(上面的发送同步消息,delayLevel的值就为0,因为不延时)
     * 在start版本中 延时消息一共分为18个等级分别为:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
     */
    public void sendDelayMsg(String topic,String msgBody, int delayLevel) {
        rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(msgBody).build(), messageTimeOut, delayLevel);
    }
    public void sendDelayKeyMsg(String topic,String msgBody, int delayLevel,String key) {
        rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(msgBody).setHeader("KEYS",key).build(), messageTimeOut, delayLevel);
    }

    /**
     * 发送单向消息(只负责发送消息,不等待应答,不关心发送结果,如日志)
     */
    public void sendOneWayMsg(String topic,String msgBody) {
        rocketMQTemplate.sendOneWay(topic, MessageBuilder.withPayload(msgBody).build());
    }
    public void sendOneWayKeyMsg(String topic,String msgBody,String key) {
        rocketMQTemplate.sendOneWay(topic, MessageBuilder.withPayload(msgBody).setHeader("KEYS",key).build());
    }

    /**
     * 发送带tag的消息,直接在topic后面加上":tag"
     */
    public SendResult sendTagMsg(String topic,String msgBody,String tag) {
        return rocketMQTemplate.syncSend(topic + ":"+tag, MessageBuilder.withPayload(msgBody).build());
    }
    public SendResult sendTagKeyMsg(String topic,String msgBody,String tag,String key) {
        return rocketMQTemplate.syncSend(topic + ":"+tag, MessageBuilder.withPayload(msgBody).setHeader("KEYS",key).build());
    }
}

4、配置消费者

@RocketMQMessageListener(topic = "${rocketmq.topic}",nameServer = "${rocketmq.nameServer}",consumerGroup = "${rocketmq.consumer.group}")
@Component
public class ComsumerListener implements RocketMQListener<String> {

    @Override
    public void onMessage(String str) {
        try {
            Thread.sleep(1000);
            System.out.println("睡眠结束");
        }catch (Exception e){
            System.out.println("异常");
        }
        System.out.println("接收信息:"+str);

    }
}

5、测试推送/消费mq

@RestController("/mq")
public class RocketMqTestController {
    @Autowired
    private MQProducerService mqProducerService;

    @GetMapping(value = "/sned")
    @ResponseBody
    public Boolean snedMqMsg(){
        mqProducerService.sendMsg("test-topic","发送mq消息:1111");
        System.out.println("发送成功");
        return true;
    }
}

测试结果:


image.png

二、使用rocketmq-client组件

1、添加依赖


<dependency>

  <groupId>org.apache.rocketmq</groupId>

  <artifactId>rocketmq-client</artifactId>

  <version>4.9.1</version>

</dependency>

2、添加配置

rocketmq:
  config:
    topic: test-topic
    nameServerAddress: xx.xx.xxx.xxx:9876;xx.xx.xxx.xxx:9876
    consumerGroupId: test-group

添加配置类

@Component
public class RocketMqProducerConfig {
    private static String nameServerAddress;
    private static String topic;
    private static String groupId;



    @Value("${rocketmq.nameServerAddress}")
    private String appAddress;
    @Value("${rocketmq.config.topic}")
    private String appTopic;
    @Value("${rocketmq.esSync.consumerGroupId}")
    private String appGroupId;
    

    @PostConstruct
    public void getConfig() {
        topic = this.appTopic;
        nameServerAddress = this.appAddress;
        groupId = this.appGroupId;
    }

    public static String getNameServerAddress() {
        return nameServerAddress;
    }
    public static String getTopic() {
        return topic;
    }
    public static String getGroupId() {
        return groupId;
    }
 
}

添加消息dto

public class MqMsgDto implements Serializable {
    private String pushId;

    public String getPushId() {
        return pushId;
    }
    public void setPushId(String pushId) {
        this.pushId = pushId;
    }

}

3、配置生产者

@Service
@DependsOn("rocketMqProducerConfig")
public class CalendarPushMqServiceImpl implements ICalendarPushMqApi {
    protected static final Logger logger = LoggerFactory.getLogger(CalendarPushMqServiceImpl.class);
    private static DefaultMQProducer rocketProducer = null;

static {
        try {
            // 实例化消息生产者Producer
            rocketProducer = new DefaultMQProducer(RocketMqProducerConfig.getGroupId());
            // 设置NameServer的地址
            rocketProducer.setNamesrvAddr(RocketMqProducerConfig.getNameServerAddress());
            rocketProducer.start();
            logger.error("rocketMq生产者初始化成功");
        } catch (Exception e) {
            logger.error("rocketMq生产者初始化失败:", e);
        }
  }

    public void pushMsg(MqMsgDto mqMsgDto) {
        try {
            Message message = new Message();
            message.setTopic(RocketMqProducerConfig.getTopic());
            message.setBody(JsonUtil.json(mqMsgDto).getBytes());
            message.setKeys(calendarMqMsgDto.getCalendarId());
            message.setTags(Long.toString(new Date().getTime()));
            //延时消费,第二个等级-5s
            message.setDelayTimeLevel(2);
            rocketProducer.send(message);
        } catch (Exception e) {
            logger.error("推送mq出错:", e);
        }
    }
}

4、配置消费者

@Configuration
public class RocketMqConfiguration {
    protected static final Logger logger = LoggerFactory.getLogger(RocketMqConfiguration.class);
    /**
     * -推送开放搜索-普通消费队列
     * */
    @Bean
    public DefaultMQPushConsumer testConsumer() {
        try {
            // 实例化消费者
            DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(RocketMqProducerConfig.getGroupId());
            // 设置NameServer的地址
            consumer.setNamesrvAddr(RocketMqProducerConfig.getNameServerAddress());
            // 订阅一个或者多个Topic,以及Tag来过滤需要消费的消息
            consumer.subscribe(RocketMqProducerConfig.getTopic(), "*");
            consumer.registerMessageListener(new RocketMqMessageListener());
            consumer.start();
            logger.info("初始化消费者成功");
            return consumer;
        }catch (Exception e){
            logger.error("初始化消费队列失败:",e);
        }
        return null;
    }
}

监听器配置

@Component
public class RocketMqMessageListener implements MessageListenerConcurrently {

    private static final Logger logger = LoggerFactory.getLogger(RocketMqMessageListener.class);

    private Test test= SpringBeanUtil.getBean(Test.class);
    @Override
    public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
        if(CollectionUtils.isEmpty(list)){
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
        for (MessageExt message : list) {
            String mdmqMessage = null;
            MqMsgDto dto = null;
            try {
                mdmqMessage = new String(message.getBody(), StandardCharsets.UTF_8);
                dto = JSON.parseObject(mdmqMessage, MqMsgDto.class);
                if (Objects.isNull(dto)){
                    continue;
                }
               logger.error("接收信息:"+JSON.toJSONString(dto));
               test.update(dto);
            }catch (Exception e){
                logger.error("接收无效消息,自动过滤",e);
                continue;
            }
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
}
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 215,133评论 6 497
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 91,682评论 3 390
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 160,784评论 0 350
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 57,508评论 1 288
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 66,603评论 6 386
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 50,607评论 1 293
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,604评论 3 415
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,359评论 0 270
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,805评论 1 307
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,121评论 2 330
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,280评论 1 344
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,959评论 5 339
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,588评论 3 322
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,206评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,442评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,193评论 2 367
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,144评论 2 352