李锋镝的博客

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

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

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

背景

我们使用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日
相关文章
  • SpringBoot 中内置的 49 个常用工具类2025年5月26日
  • 一篇文章帮你彻底搞清楚“I/O多路复用”和“异步I/O”的前世今生2020年5月23日
  • 几款Java开发者必备常用的工具,准点下班不在话下2021年2月19日
  • 配置Jackson使用字段而不是getter/setter来序列化和反序列化2026年3月19日
  • Java设计模式:状态模式2025年5月22日

李锋镝

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

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

文章评论

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

Your time is limited, so don't waste it living someone else's life. Don't be trapped by dogma -- which is living with the results of other people's thinking. Don't let the noise of other's opinions drown out your own inner voice. And most important, have the courage to follow your heart and intuition. They somehow already know what you truly want to become. Everything else is secondary.

听点儿音乐吧 朋友~
文章目录
最新 热点 随机
最新 热点 随机
给主题增加了游客中心功能 祝大家中秋安康 游本昌去世 Redis7+&8.X 全新进阶系列(02):Redis Functions 详解——替代Lua脚本的官方轻量化函数方案 Redis7.x&8.x 全新进阶系列(01):划时代升级总览——从6.x到7.x/8.x全版本变革全景 让WordPress静态化之Rocket‑Nginx
游本昌去世关于主题加载速度优化的一点儿小演进给主题增加了Now、每日心情、年度回顾、岁月同一天、随机漫步等功能WordPress缓存插件WP Fastest Cache、WP Rocket 、FlyingPress对比关于使用AI的一些思考WordPress下一代默认主题Ipsum预览
每日早报 · 2026年8月20日 · 早上好 优化了MYSQL大量写入问题,老板奖励了1000块给我 LangGraph 深度实战指南:从基础架构到生产级 AI Agent 工作流构建 Docker核心概念解析及使用 Java布尔运算 Spring Boot 配置加载优先级总结
最近评论
老张博客 发布于 1 小时前(09月28日) 不需要注册,而又有自己的“中心”,这个功能真的不错,极大增加了互动性。
lijie blog 发布于 2 小时前(09月28日) :74: 可以继续扩展读者墙、那年今日(参考qq空间的)
不凡 发布于 3 小时前(09月28日) 是以邮箱作为唯一“注册”凭证,是吗?
李锋镝 发布于 6 小时前(09月28日) 是的,不然就有些太单调了,这样互动性也更好一些
我是军爸 发布于 6 小时前(09月28日) 还带有徽章功能啊,有点像论坛的感觉了
标签聚合
Redis 多线程 JVM Theme ElasticSearch K8s 日常 AI编程 Claude AI MySQL 数据库 分布式 SpringBoot JAVA Spring WordPress 架构 IDEA SQL
友情链接
  • lijie blog
  • 韩情脉脉
  • Serendipity
  • 林羽凡
  • 彬红茶日记
  • 韩小韩博客
  • 蜗牛工作室
  • 九仞之行
  • 老张博客
  • 懋和道人
  • 搬砖日记
  • 临窗旋墨
  • Mr.Sun的博客
  • 若梦博客
  • 皮皮社
  • 知向前端
  • sssr7844的博客
  • 志文工作室
  • 瓦匠个人小站
  • 哥斯拉

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号