解決SpringBoot整合RocketMQ遇到的坑
應(yīng)用場景
在實現(xiàn)RocketMQ消費時,一般會用到@RocketMQMessageListener注解定義Group、Topic以及selectorExpression(數(shù)據(jù)過濾、選擇的規(guī)則)為了能支持動態(tài)篩選數(shù)據(jù),一般都會使用表達(dá)式,然后通過apollo或者cloud config進(jìn)行動態(tài)切換。
引入依賴
<!-- RocketMq Spring Boot Starter--> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.0.4</version> </dependency>
消費者代碼
@RocketMQMessageListener(consumerGroup = "${rocketmq.group}",topic ="${rocketmq.topic}",selectorExpression = "${rocketmq.selectorExpression}")
public class Consumer implements RocketMQListener<String> {
@Override
public void onMessage(String s) {
System.out.println("消費到的數(shù)據(jù)為:"+s);
}
}
問題排查
RocketMQMessageListener整個注解默認(rèn)selectorExpression為*,表示接收當(dāng)前Topic下的所有數(shù)據(jù),如果我們想對tags進(jìn)行動態(tài)配置,在使用${rocketmq.selectorExpression}表達(dá)式時會發(fā)現(xiàn)所有數(shù)據(jù)全被過濾了,跟蹤源碼(ListenerContainerConfiguration.java)發(fā)現(xiàn)在創(chuàng)建listener時selectorExpression的數(shù)據(jù)在通environment環(huán)境變量中獲取對應(yīng)的數(shù)據(jù)后又被覆蓋了,導(dǎo)致整個過濾條件被變更為表達(dá)式。
@Override
public void afterSingletonsInstantiated() {
// 獲取所有所有使用了RocketMQMessageListener注解的bean
Map<String, Object> beans = this.applicationContext.getBeansWithAnnotation(RocketMQMessageListener.class);
if (Objects.nonNull(beans)) {
// 循環(huán)注冊容器
beans.forEach(this::registerContainer);
}
}
private void registerContainer(String beanName, Object bean) {
Class<?> clazz = AopProxyUtils.ultimateTargetClass(bean);
// 校驗當(dāng)前bean是否實現(xiàn)了RocketMQListener接口
if (!RocketMQListener.class.isAssignableFrom(bean.getClass())) {
throw new IllegalStateException(clazz + " is not instance of " + RocketMQListener.class.getName());
}
// 獲取bean上的annotation
RocketMQMessageListener annotation = clazz.getAnnotation(RocketMQMessageListener.class);
// 解析group及topic,可支持表達(dá)式
String consumerGroup = this.environment.resolvePlaceholders(annotation.consumerGroup());
String topic = this.environment.resolvePlaceholders(annotation.topic());
boolean listenerEnabled =
(boolean)rocketMQProperties.getConsumer().getListeners().getOrDefault(consumerGroup, Collections.EMPTY_MAP)
.getOrDefault(topic, true);
if (!listenerEnabled) {
log.debug(
"Consumer Listener (group:{},topic:{}) is not enabled by configuration, will ignore initialization.",
consumerGroup, topic);
return;
}
validate(annotation);
String containerBeanName = String.format("%s_%s", DefaultRocketMQListenerContainer.class.getName(),
counter.incrementAndGet());
GenericApplicationContext genericApplicationContext = (GenericApplicationContext)applicationContext;
// 注冊bean的,調(diào)用createRocketMQListenerContainer
genericApplicationContext.registerBean(containerBeanName, DefaultRocketMQListenerContainer.class,
() -> createRocketMQListenerContainer(containerBeanName, bean, annotation));
DefaultRocketMQListenerContainer container = genericApplicationContext.getBean(containerBeanName,
DefaultRocketMQListenerContainer.class);
if (!container.isRunning()) {
try {
container.start();
} catch (Exception e) {
log.error("Started container failed. {}", container, e);
throw new RuntimeException(e);
}
}
log.info("Register the listener to container, listenerBeanName:{}, containerBeanName:{}", beanName, containerBeanName);
}
private DefaultRocketMQListenerContainer createRocketMQListenerContainer(String name, Object bean,
RocketMQMessageListener annotation) {
DefaultRocketMQListenerContainer container = new DefaultRocketMQListenerContainer();
container.setRocketMQMessageListener(annotation);
String nameServer = environment.resolvePlaceholders(annotation.nameServer());
nameServer = StringUtils.isEmpty(nameServer) ? rocketMQProperties.getNameServer() : nameServer;
String accessChannel = environment.resolvePlaceholders(annotation.accessChannel());
container.setNameServer(nameServer);
if (!StringUtils.isEmpty(accessChannel)) {
container.setAccessChannel(AccessChannel.valueOf(accessChannel));
}
container.setTopic(environment.resolvePlaceholders(annotation.topic()));
// 此處已經(jīng)根據(jù)表達(dá)式將數(shù)據(jù)取出
String tags = environment.resolvePlaceholders(annotation.selectorExpression());
if (!StringUtils.isEmpty(tags)) {
container.setSelectorExpression(tags);
}
container.setConsumerGroup(environment.resolvePlaceholders(annotation.consumerGroup()));
// 此處將SelectorExpression的數(shù)據(jù)覆蓋成了表達(dá)式
container.setRocketMQMessageListener(annotation);
container.setRocketMQListener((RocketMQListener)bean);
container.setObjectMapper(objectMapper);
container.setMessageConverter(rocketMQMessageConverter.getMessageConverter());
container.setName(name); // REVIEW ME, use the same clientId or multiple?
return container;
}
問題解決
因為ListenerContainerConfiguration類是實現(xiàn)了SmartInitializingSingleton接口的afterSingletonsInstantiated方法,我們可以通過反射對selectorExpression的數(shù)據(jù)在ListenerContainerConfiguration進(jìn)行初始化前進(jìn)行解析并賦值回去。
/**
* 在springboot初始化后,RocketMQ容器初始化前利用反射動態(tài)改變數(shù)據(jù)
**/
@Configuration
public class ChangeSelectorExpressionBeforeMQInit implements InitializingBean {
@Autowired
private ApplicationContext applicationContext;
@Autowired
private StandardEnvironment environment;
@Override
public void afterPropertiesSet() throws Exception {
Map<String,Object> beans =applicationContext.getBeansWithAnnotation(RocketMQMessageListener.class);
for (Object bean : beans.values()){
Class<?> clazz = AopProxyUtils.ultimateTargetClass(bean);
if (!RocketMQListener.class.isAssignableFrom(bean.getClass())) {
continue;
}
RocketMQMessageListener annotation = clazz.getAnnotation(RocketMQMessageListener.class);
InvocationHandler invocationHandler = Proxy.getInvocationHandler(annotation);
Field field = invocationHandler.getClass().getDeclaredField("memberValues");
field.setAccessible(true);
Map<String, Object> memberValues = (Map<String, Object>) field.get(invocationHandler);
for (Map.Entry<String,Object> entry: memberValues.entrySet()) {
if(Objects.nonNull(entry)){
memberValues.put(entry.getKey(),environment.resolvePlaceholders(String.valueOf(entry.getValue())));
}
}
}
}
}
初次之外,在2.1.0版本的依賴包中已經(jīng)修復(fù)了此Bug,在不造成依賴沖突的前提下,建議使用2.1.0以上的版本包。
以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
- 淺談Springboot整合RocketMQ使用心得
- springBoot整合RocketMQ及坑的示例代碼
- Springboot RocketMq實現(xiàn)過程詳解
- Springboot詳解RocketMQ實現(xiàn)廣播消息流程
- SpringBoot整合RocketMQ實現(xiàn)發(fā)送同步消息
- 解決springboot集成rocketmq關(guān)于tag的坑
- SpringBoot集成RocketMQ的使用示例
- 在SpringBoot中利用RocketMQ實現(xiàn)批量消息消費功能
- SpringBoot項目嵌入RocketMQ的實現(xiàn)示例
- RocketMQ在Spring Boot上的基礎(chǔ)使用
相關(guān)文章
Java 輕松實現(xiàn)二維數(shù)組與稀疏數(shù)組互轉(zhuǎn)
在某些應(yīng)用場景中需要大量的二維數(shù)組來進(jìn)行數(shù)據(jù)存儲,但是二維數(shù)組中卻有著大量的無用的位置占據(jù)著內(nèi)存空間,稀疏數(shù)組就是為了優(yōu)化二維數(shù)組,節(jié)省內(nèi)存空間2022-04-04
SpringBoot Validation入?yún)⑿r瀲H化的項目實踐
在Spring Boot中,可以使用Validation和國際化來實現(xiàn)對入?yún)⒌男r?本文就來介紹一下SpringBoot Validation入?yún)⑿r瀲H化,具有一定的參考價值,感興趣的可以了解一下2023-10-10
SpringBoot項目jar發(fā)布后如何獲取jar包所在目錄路徑
這篇文章主要介紹了SpringBoot項目jar發(fā)布后如何獲取jar包所在目錄路徑,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-11-11
Object.wait()與Object.notify()的用法詳細(xì)解析
以下是對java中Object.wait()與Object.notify()的用法進(jìn)行了詳細(xì)的分析介紹,需要的朋友可以過來參考下2013-09-09

