李锋镝的博客

  • 首页
  • 时间轴
  • 说说
  • 每日心情
  • Now
  • 系列文章
  • 论坛
  • 左邻右舍
    • 左邻右舍
    • 博友圈
  • 留言
    • 留言
    • 走心评论
  • 关于
    • 关于本站
    • 网站地图
    • 网站统计
    • 另一个网站
    • 我的导航站
    • 赞助
  • 🚇开往
!Destiny
惟坚韧者始能遂其志
  1. 首页
  2. 原创
  3. 正文

使用RocketMQ时,服务启动过程中,Consumer在服务未启动时消费消息问题处理

2022年6月23日 约 614 字3 分钟 149 0 0
本文最后更新于 2022年6月23日,距今已 1552 天,其中的信息可能已经发生变化,请注意甄别。

背景

我们使用RocketMQ时,一般Consumer启动都是使用的@PostConstruct注解。(@PostConstruct:用于在执行任何初始化时执行依赖注入后需要执行的方法。),或者使用bean的方式配置。

配置如下:

生产者配置

在配置类中配置所有生产者,在业务中注入使用,将生产者的启动和销毁绑定到 Bean 的初始化和销毁上:

@Configuration
public class MQProducerConfig {

    // 第一个生产者
    @Bean(initMethod = "start", destroyMethod = "shutdown")
    public DefaultMQProducer demo1MQProducer() {
        DefaultMQProducer defaultMQProducer = new DefaultMQProducer();
        defaultMQProducer.setNamesrvAddr("<nameServer>");
        defaultMQProducer.setProducerGroup("<group>");
        defaultMQProducer.setInstanceName("<instanceName>");
        // 其他生产者配置
        return defaultMQProducer;
    }

    // 第二个生产者
    @Bean(initMethod = "start", destroyMethod = "shutdown")
    public DefaultMQProducer demo2MQProducer() {
        DefaultMQProducer defaultMQProducer = new DefaultMQProducer();
        defaultMQProducer.setNamesrvAddr("<nameServer>");
        defaultMQProducer.setProducerGroup("<group>");
        defaultMQProducer.setInstanceName("<instanceName>");
        // 其他生产者配置
        return defaultMQProducer;
    }

    // ......
}

消费者配置

在配置类中配置所有消费者,将消费者的启动和销毁绑定到 Bean 的初始化和销毁上:

@Configuration
public class MQConsumerConfig {

    // 第一个消费者
    @Bean(initMethod = "start", destroyMethod = "shutdown")
    public DefaultMQPushConsumer demo1Consumer() throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<group>");
        consumer.setNamesrvAddr("<nameServer>");
        consumer.setInstanceName("<instanceName>");
        consumer.subscribe("<topic>", "<tag>");
        consumer.setMessageListener(<listener>);
        // 其他消费者配置
        return consumer;
    }

    // 第二个消费者
    @Bean(initMethod = "start", destroyMethod = "shutdown")
    public DefaultMQPushConsumer demo2Consumer() throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<group>");
        consumer.setNamesrvAddr("<nameServer>");
        consumer.setInstanceName("<instanceName>");
        consumer.subscribe("<topic>", "<tag>");
        consumer.setMessageListener(<listener>);
        // 其他消费者配置
        return consumer;
    }

    // ......
}

优化

上述配置在项目启动 Bean 加载的时候就会启动生产者和消费者,导致项目启动慢,并且会在项目还未启动完,就会有大量消息涌入,所以可以使用 ApplicationRunner 或 CommandLineRunner 接口在项目启动成功后再执行 MQ 的启动。同时,去掉Consumer的@PostConstruct注解。

注解如下:

// 去掉initMethod配置
@Bean(destroyMethod = "shutdown")

将上述配置中的消费者启动不再绑定到 Bean 初始化阶段。

新的消费者配置如下:

@Configuration
public class MQConsumerConfig {

    // 第一个消费者
    @Bean(destroyMethod = "shutdown")
    public DefaultMQPushConsumer demo1Consumer() throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<group>");
        consumer.setNamesrvAddr("<nameServer>");
        consumer.setInstanceName("<instanceName>");
        consumer.subscribe("<topic>", "<tag>");
        consumer.setMessageListener(<listener>);
        // 其他消费者配置
        return consumer;
    }

    // 第二个消费者
    @Bean(destroyMethod = "shutdown")
    public DefaultMQPushConsumer demo2Consumer() throws Exception {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<group>");
        consumer.setNamesrvAddr("<nameServer>");
        consumer.setInstanceName("<instanceName>");
        consumer.subscribe("<topic>", "<tag>");
        consumer.setMessageListener(<listener>);
        // 其他消费者配置
        return consumer;
    }

    // ......
}

在项目启动完成节点统一启动消费者:

@Slf4j
@Component
public class MqConsumerApplicationRunner implements ApplicationRunner {

    @Autowired
    private Map<String, DefaultMQPushConsumer> defaultMQPushConsumerMap;

    @Override
    public void run(ApplicationArguments args) {
        if (CollectionUtils.isEmpty(defaultMQPushConsumerMap)) {
            return;
        }
        defaultMQPushConsumerMap.forEach((bean, consumer) -> {
            try {
                consumer.start();
            } catch (MQClientException e) {
                log.error("Consumer bean:[{}] start error.", bean, e);
            }
        }
    }
}
除非注明,否则均为李锋镝的博客原创文章,转载必须以链接形式标明本文链接

本文链接:https://www.lifengdi.com/article/3888

本作品采用 知识共享署名-非商业性使用-相同方式共享 4.0 国际许可协议 进行许可
分享到

使用RocketMQ时,服务启动过程中,Consumer在服务未启动时消费消息问题处理

也可使用浏览器菜单中的「分享」功能

微信扫一扫分享

标签: JAVA MQ RocketMQ SpringBoot
最后更新:2022年6月23日

岁月同一天 9 月 22 日

回望过去的今天,你在写什么

  • 3 年前 2023年9月22日
    《人生海海》读后感

    前段时间读完了麦家的《人生海海》,起初对人生海海这个词不太理解,记得第一次听到这个词,应该是在《欢喜就好》这首歌里,里面…

  • 4 年前 2022年9月22日
    OHCache使用

    OHCache介绍 缓存框架OHC基于Java语言实现,并以类库的形式供其他Java程序调用,是一种以单机模式运行的堆外…

相关文章
  • UUID太长怎么办?快来试试NanoId2022年3月31日
  • Java 灵魂拷问 13 个为什么,你都会哪些?2025年5月20日
  • SpringBoot整合MongoDB2022年9月2日
  • 我要狠狠的反驳“公司禁止使用 Lombok ”的观点!2021年3月3日
  • SpringBoot 实现 RSA+AES 自动接口解密2025年5月26日

李锋镝

既然选择了远方,便只顾风雨兼程。

打赏 点赞
< 上一篇
下一篇 >
1234567891112131415161718192021222324252627282930313233343536373839404142434446474849505152535455575859606162636465666769727476777879808182858687909293949596979899
取消回复
…

文章评论

还没有评论,快来抢沙发吧~

路曼曼其修远兮,吾将上下而求索。

听点儿音乐吧 朋友~
文章目录
最新 热点 随机
最新 热点 随机
让WordPress静态化之Rocket‑Nginx WordPress下一代默认主题Ipsum预览 C++之父重磅发声:AI编程正在毁掉一代程序员 支撑全网40%网站的WordPress正在重新拥抱PHP生态 Kratos-plus v1.1.24版本更新说明 关于主题加载速度优化的一点儿小演进
关于主题加载速度优化的一点儿小演进给主题增加了Now、每日心情、年度回顾、岁月同一天、随机漫步等功能WordPress缓存插件WP Fastest Cache、WP Rocket 、FlyingPress对比关于使用AI的一些思考Kratos+ v1.1.16版本更新说明AI时代,个人技术博客的出路在哪里?
MybatisCodeHelperPro激活 服务器差点被一群垃圾爬虫搞挂了 玉楼春·尊前拟把归期说 SpringBoot启动概述(SpringBoot2.1.7) Java 灵魂拷问 13 个为什么,你都会哪些? Java简介
最近评论
obaby 发布于 11 分钟前(09月22日) 我现在是个假的wp了,哈哈哈 wp终于改了主题的命名风格了
李锋镝 发布于 2 天前(09月20日) 静态博客我之前也用过,但是感觉不是很方便,后来就一直用的WordPress
Sheep5 发布于 2 天前(09月20日) 我直接用静态博客,天然有速度优势。
不凡 发布于 2 天前(09月20日) 主要是wordpress插件丰富,需要什么功能插件,插件市场应有尽有,typecho是性能更好、更轻...
李锋镝 发布于 2 天前(09月20日) 忒极简了,而且这个布局我也搞不懂
标签聚合
架构 Claude ElasticSearch JVM SpringBoot MQ AI MySQL Spring AI编程 WordPress K8s 数据库 JAVA 日常 Redis IDEA SQL 多线程 分布式
友情链接
  • sssr7844的博客
  • 彬红茶日记
  • 九仞之行
  • 蜗牛工作室
  • 搬砖日记
  • 瓦匠个人小站
  • 若梦博客
  • 懋和道人
  • 皮皮社
  • 知向前端
  • Honesty
  • 老张博客
  • 韩情脉脉
  • 临窗旋墨
  • lijie blog
  • 韩小韩博客
  • Mr.Sun的博客
  • 哥斯拉
  • Serendipity
  • 林羽凡

COPYRIGHT © 2016-2026 lifengdi.com. ALL RIGHTS RESERVED.

lifengdi.com

Domain age badge for lifengdi.com

Theme Kratos-plus By Dylan Li

津ICP备2024022503号-3

京公网安备11011502039375号