記一次RocketMQConsumer 服務關閉出現InterruptException例外
背景提要
出現問題主要還是版本升級
-
老版本核心rocketmq依賴
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>spring-boot-starter-rocketmq</artifactId> <version>${vesion}</version> </dependency> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.3.2</version> </dependency> -
新版本核心rocketmq依賴
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.9.2</version> </dependency> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.2.1</version> </dependency>
java.lang.InterruptedException
簡單列舉一個InterruptedException
java.sql.SQLException: interrupt
at com.alibaba.druid.pool.DruidDataSource.getConnectionInternal(DruidDataSource.java:1430) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.pool.DruidDataSource.getConnectionDirect(DruidDataSource.java:1272) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.filter.FilterChainImpl.dataSource_connect(FilterChainImpl.java:5007) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.filter.FilterAdapter.dataSource_getConnection(FilterAdapter.java:2745) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.filter.FilterChainImpl.dataSource_connect(FilterChainImpl.java:5003) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.filter.stat.StatFilter.dataSource_getConnection(StatFilter.java:680) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.filter.FilterChainImpl.dataSource_connect(FilterChainImpl.java:5003) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.pool.DruidDataSource.getConnection(DruidDataSource.java:1250) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.pool.DruidDataSource.getConnection(DruidDataSource.java:1242) ~[druid-1.1.12.jar!/:1.1.12]
at com.alibaba.druid.pool.DruidDataSource.getConnection(DruidDataSource.java:89) ~[druid-1.1.12.jar!/:1.1.12]
// 省略部分堆疊資訊
at org.apache.rocketmq.spring.support.DefaultRocketMQListenerContainer.handleMessage(DefaultRocketMQListenerContainer.java:399) [rocketmq-spring-boot-2.2.1.jar!/:2.2.1]
at org.apache.rocketmq.spring.support.DefaultRocketMQListenerContainer.access$100(DefaultRocketMQListenerContainer.java:71) [rocketmq-spring-boot-2.2.1.jar!/:2.2.1]
at org.apache.rocketmq.spring.support.DefaultRocketMQListenerContainer$DefaultMessageListenerConcurrently.consumeMessage(DefaultRocketMQListenerContainer.java:359) [rocketmq-spring-boot-2.2.1.jar!/:2.2.1]
at cn.techwolf.trace.rocketmq.spring.TracingMessageListenerConcurrently.consumeMessage(TracingMessageListenerConcurrently.java:37) [instrument-rocketmq-spring-1.101.jar!/:1.101]
at org.apache.rocketmq.client.impl.consumer.ConsumeMessageConcurrentlyService$ConsumeRequest.run(ConsumeMessageConcurrentlyService.java:392) [rocketmq-client-4.9.2.jar!/:4.9.2]
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) [?:1.8.0_202]
at java.util.concurrent.FutureTask.run(FutureTask.java:266) [?:1.8.0_202]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [?:1.8.0_202]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [?:1.8.0_202]
at java.lang.Thread.run(Thread.java:748) [?:1.8.0_202]
Caused by: java.lang.InterruptedException
注意
以下分析僅為個人觀點,并不一定正確,如有杠精,請勿繼續觀看,歡迎留言討論哈
提示:本文涉及的一些類,在 記一次RocketMQ服務啟動時 NullPointerException問題 本文不做一些詳細解釋
spring關閉&rocketmq關閉時機(穿插 并不重要)
-
spring

-
rocketmq 繼承了SmartCycle 會在容器關閉的時候 回呼其stop方法
rocketmq shutdown 分析
首先確認我們的RocketMQConsumer 的實作 DefaultRocketMQListenerContainer:
// 類全路徑org.apache.rocketmq.spring.support.DefaultRocketMQListenerContainer
public class DefaultRocketMQListenerContainer implements InitializingBean, RocketMQListenerContainer, SmartLifecycle, ApplicationContextAware {
private DefaultMQPushConsumer consumer;
@Override
public void destroy() {// DisposableBean 回呼 銷毀bean的時候 呼叫
this.setRunning(false);
if (Objects.nonNull(consumer)) {
consumer.shutdown();
}
log.info("container destroyed, {}", this.toString());
}
@Override
public void stop() { // SmartLifecycle 回呼 關閉容器前回呼
if (this.isRunning()) {
if (Objects.nonNull(consumer)) {
consumer.shutdown();
}
setRunning(false);
}
}
}
可以看到關閉的話 是先呼叫到 DefaultRocketMQListenerContainer.stop方法, 接下來就是看看consumer.shutdown() 方法了:
// 類全路徑 org.apache.rocketmq.client.consumer.DefaultMQPushConsumer
public class DefaultMQPushConsumer extends ClientConfig implements MQPushConsumer {
protected final transient DefaultMQPushConsumerImpl defaultMQPushConsumerImpl;
/**
* Maximum time to await message consuming when shutdown consumer, 0 indicates no await.
*/
private long awaitTerminationMillisWhenShutdown = 0;
@Override
public void shutdown() {
this.defaultMQPushConsumerImpl.shutdown(awaitTerminationMillisWhenShutdown);
if (null != traceDispatcher) {
traceDispatcher.shutdown();
}
}
}
// 類全路徑 org.apache.rocketmq.client.impl.consumer.DefaultMQPushConsumerImpl
public class DefaultMQPushConsumerImpl implements MQConsumerInner {
private ConsumeMessageService consumeMessageService;
public synchronized void shutdown(long awaitTerminateMillis) {
switch (this.serviceState) {
case CREATE_JUST:
break;
case RUNNING:
this.consumeMessageService.shutdown(awaitTerminateMillis);
this.persistConsumerOffset();
this.mQClientFactory.unregisterConsumer(this.defaultMQPushConsumer.getConsumerGroup());
this.mQClientFactory.shutdown();
log.info("the consumer [{}] shutdown OK", this.defaultMQPushConsumer.getConsumerGroup());
this.rebalanceImpl.destroy();
this.serviceState = ServiceState.SHUTDOWN_ALREADY;
break;
case SHUTDOWN_ALREADY:
break;
default:
break;
}
}
}
// ConsumeMessageService 有兩個實作類 分別是 ConsumeMessageConcurrentlyService 和 ConsumeMessageOrderlyService
// 我們使用的 consumeMessageService 的具體實作類是 ConsumeMessageConcurrentlyService 具體問題具體分析哈
// 類全路徑 org.apache.rocketmq.client.impl.consumer.ConsumeMessageConcurrentlyService
public class ConsumeMessageConcurrentlyService implements ConsumeMessageService {
private final ThreadPoolExecutor consumeExecutor;
public void shutdown(long awaitTerminateMillis) {
this.scheduledExecutorService.shutdown();
ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
this.cleanExpireMsgExecutors.shutdown();
}
}
// 類全路徑org.apache.rocketmq.common.utils.ThreadUtils
public final class ThreadUtils {
public static void shutdownGracefully(ExecutorService executor, long timeout, TimeUnit timeUnit) {
// Disable new tasks from being submitted.
executor.shutdown();
try {
// Wait a while for existing tasks to terminate.
if (!executor.awaitTermination(timeout, timeUnit)) { // 注意這里
executor.shutdownNow();
// Wait a while for tasks to respond to being cancelled.
if (!executor.awaitTermination(timeout, timeUnit)) {
log.warn(String.format("%s didn't terminate!", executor));
}
}
} catch (InterruptedException ie) {
// (Re-)Cancel if current thread also interrupted.
executor.shutdownNow();
// Preserve interrupt status.
Thread.currentThread().interrupt();
}
}
}
從上面我們可以清楚的看到 執行到 ConsumeMessageConcurrentlyService.shutdown的時候,awaitTerminateMillis默認值是0, 執行到 ThreadUtils.shutdownGracefully時,會直接呼叫shutdonwNow,并沒有等待
shutdonw 和 shutdownNow 的區別 百度搜一搜就行了
- shutdown => 平緩關閉,等待所有已添加到執行緒池中的任務執行完在關閉
- shutdownNow => 立刻關閉,停止正在執行的任務,并回傳佇列中未執行的任務
對比以前代碼 rocketmq-client 4.3.2
public class ConsumeMessageConcurrentlyService implements ConsumeMessageService { private final ThreadPoolExecutor consumeExecutor; public void shutdown(long awaitTerminateMillis) { this.scheduledExecutorService.shutdown(); this.consumeExecutor.shutdown(); // 這里 this.cleanExpireMsgExecutors.shutdown(); } }
故 我認為是,是因為直接呼叫了 shutdonwNow 導致服務關閉的時候出現中斷 InterruptException 例外
解決方案
從上述分析可以知道 是因為直接呼叫shutdownNow 導致的,我們應該可以調整 awaitTerminateMillis引數,也就是DefaultMQPushConsumer.awaitTerminationMillisWhenShutdown引數
而目前我似乎沒看到有哪種方式支持全域配置該引數的方式(官方方式) 這里提供倆思路
RocketMQPushConsumerLifecycleListener or RocketMQPushConsumerLifecycleListener
// org.apache.rocketmq.spring.support.DefaultRocketMQListenerContainer
public class DefaultRocketMQListenerContainer implements InitializingBean, RocketMQListenerContainer, SmartLifecycle, ApplicationContextAware {
@Override
public void afterPropertiesSet() throws Exception {
initRocketMQPushConsumer();
this.messageType = getMessageType();
this.methodParameter = getMethodParameter();
log.debug("RocketMQ messageType: {}", messageType);
}
private void initRocketMQPushConsumer() throws MQClientException {
// 省略部分代碼
if (Objects.nonNull(rpcHook)) {
consumer = new DefaultMQPushConsumer(consumerGroup, rpcHook, new AllocateMessageQueueAveragely(),
enableMsgTrace, this.applicationContext.getEnvironment().
resolveRequiredPlaceholders(this.rocketMQMessageListener.customizedTraceTopic()));
consumer.setVipChannelEnabled(false);
} else {
log.debug("Access-key or secret-key not configure in " + this + ".");
consumer = new DefaultMQPushConsumer(consumerGroup, enableMsgTrace,
this.applicationContext.getEnvironment().
resolveRequiredPlaceholders(this.rocketMQMessageListener.customizedTraceTopic()));
}
// 省略部分consumer 引數配置代碼
// 可以看到這 他會回呼 prepareStart 方法
if (rocketMQListener instanceof RocketMQPushConsumerLifecycleListener) {
((RocketMQPushConsumerLifecycleListener) rocketMQListener).prepareStart(consumer);
} else if (rocketMQReplyListener instanceof RocketMQPushConsumerLifecycleListener) {
((RocketMQPushConsumerLifecycleListener) rocketMQReplyListener).prepareStart(consumer);
}
}
}
可以看到 在DefaultRocketMQListenerContainer初始化完后會回呼 RocketMQPushConsumerLifecycleListener、RocketMQPushConsumerLifecycleListener的prepareStart方法
那就很簡單了
@Service
@RocketMQMessageListener(nameServer = "${spring.rocketmq.nameServer}",
topic = "${topic}",
consumerGroup = "${group}")
public class TestConsumer implements RocketMQListener<String>, RocketMQPushConsumerLifecycleListener {
@Override
public void onMessage(String msg) {
// do something
}
@Override
public void prepareStart(DefaultMQPushConsumer consumer) {
consumer.setAwaitTerminationMillisWhenShutdown(1000); // 設定
}
}
可以抽象成一個通用類 consumer 繼承該類就行了
DefaultRocketMQListenerContainer.getConsumer
DefaultRocketMQListenerContainer會被注冊成bean 具體的實作在 org.apache.rocketmq.spring.autoconfigure.ListenerContainerConfiguration中,這里就不分析了,我們可以嘗試獲取所有的DefaultRocketMQListenerContainer然后呼叫其getConsumer方法 示例代碼 具體怎么觸發這個代碼 就得自己思考和完善了哈 本文不做解釋了
public void set() {
List<DefaultRocketMQListenerContainer> containers = getAllBeans();
for (DefaultRocketMQListenerContainer c : containers) {
c.getConsumer().setAwaitTerminationMillisWhenShutdown(1000);
}
}
最后,以上僅為本人分析,并不一定正確,用第一種方式確實沒出現了InterruptException了,不知道是否是偶然;
歡迎留言討論
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436407.html
標籤:其他
上一篇:HBase 過濾器
下一篇:hive-SQL學習筆記11
