springboot1.x版本兼容rocketmq5.1.0问题记录
·
项目依赖初始版本:

启动时发现问题报错:
因为业务需求,我们自定义MessagePushRocketMQConsumerListener继承RocketMQListener使用@RocketMQMessageListener注解 但是出现创建这个bean失败

看日志发现是rocketmq源码中ListenerContainerConfiguration的registerContainer方法中没有找到默认类 下载rocketMQ对应的源码
https://github.com/apache/rocketmq-spring/tags

找到错误信息ListenerContainerConfiguration中的registerContainer方法进行修改(参考博客https://blog.csdn.net/weixin_43811294/article/details/131427047)
修改对应报错调用代码
genericApplicationContext.registerBean(containerBeanName, DefaultRocketMQListenerContainer.class,
() -> createRocketMQListenerContainer(containerBeanName, bean, annotation));
这里用下面的代码替换上面的
BeanDefinition beanDefinition = buildBeanDefinition(containerBeanName, bean, annotation);
genericApplicationContext.registerBeanDefinition(containerBeanName, beanDefinition);
添加自定义方法
private BeanDefinition buildBeanDefinition(String name, Object bean,
RocketMQMessageListener annotation) {
String nameServer = environment.resolvePlaceholders(annotation.nameServer());
nameServer = StringUtils.isEmpty(nameServer) ? rocketMQProperties.getNameServer() : nameServer;
String accessChannel = environment.resolvePlaceholders(annotation.accessChannel());
String tags = environment.resolvePlaceholders(annotation.selectorExpression());
BeanDefinitionBuilder builder = BeanDefinitionBuilder
.genericBeanDefinition(DefaultRocketMQListenerContainer.class)
.addPropertyValue("rocketMQMessageListener", annotation)
.addPropertyValue("nameServer", nameServer)
.addPropertyValue("topic", environment.resolvePlaceholders(annotation.topic()))
.addPropertyValue("consumerGroup", environment.resolvePlaceholders(annotation.consumerGroup()))
.addPropertyValue("tlsEnable", environment.resolvePlaceholders(annotation.tlsEnable()))
.addPropertyValue("messageConverter", rocketMQMessageConverter.getMessageConverter())
.addPropertyValue("name", name);
if (!StringUtils.isEmpty(accessChannel)) {
builder.addPropertyValue("accessChannel", AccessChannel.valueOf(accessChannel));
}
if (!StringUtils.isEmpty(tags)) {
builder.addPropertyValue("selectorExpression", tags);
}
if (RocketMQListener.class.isAssignableFrom(bean.getClass())) {
builder.addPropertyValue("rocketMQListener", bean);
} else if (RocketMQReplyListener.class.isAssignableFrom(bean.getClass())) {
builder.addPropertyValue("rocketMQReplyListener", bean);
}
return builder.getBeanDefinition();
}

修改完后重新打包指定版本这里我改为2.2.3-hc-SNAPSHOT 再次编译发现错误发生了变化

仔细排查后发现 2.2.3的源码中引用的rocketmq版本是5.0.0的,而我们使用的是5.1.0的 修改为5.1.0,更新maven

来到上次编译报错的代码段,发现这两个类没有找到,查看set方法发现了目录结构发生了变化我们进行替换,再次编译

在微服务中重新拉取maven(因为使用的快照版本 不需要更新版本号),重启

成功!
更多推荐


所有评论(0)