diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index 3e146b754a..bce21c5209 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -33,6 +33,8 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.Optional; +import java.util.Objects; import org.apache.rocketmq.acl.AccessValidator; import org.apache.rocketmq.broker.client.ClientHousekeepingService; import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; @@ -650,10 +652,13 @@ public class BrokerController { public long headSlowTimeMills(BlockingQueue q) { long slowTimeMills = 0; - final Runnable peek = q.peek(); - if (peek != null) { - RequestTask rt = BrokerFastFailure.castRunnable(peek); - slowTimeMills = rt == null ? 0 : this.messageStore.now() - rt.getCreateTimestamp(); + Optional op = q.stream() + .map(BrokerFastFailure::castRunnable) + .filter(Objects::nonNull) + .findFirst(); + if (op.isPresent()) { + RequestTask rt = op.get(); + slowTimeMills = this.messageStore.now() - rt.getCreateTimestamp(); } if (slowTimeMills < 0) { diff --git a/broker/src/test/java/org/apache/rocketmq/broker/BrokerControllerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/BrokerControllerTest.java index dae1335540..e8442a4d4b 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/BrokerControllerTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/BrokerControllerTest.java @@ -18,10 +18,16 @@ package org.apache.rocketmq.broker; import java.io.File; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + +import org.apache.rocketmq.broker.latency.FutureTaskExt; import org.apache.rocketmq.common.BrokerConfig; import org.apache.rocketmq.common.UtilAll; import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.apache.rocketmq.remoting.netty.NettyServerConfig; +import org.apache.rocketmq.remoting.netty.RequestTask; import org.apache.rocketmq.store.config.MessageStoreConfig; import org.junit.After; import org.junit.Ignore; @@ -47,4 +53,33 @@ public class BrokerControllerTest { public void destroy() { UtilAll.deleteFile(new File(new MessageStoreConfig().getStorePathRootDir())); } + + @Test + public void testHeadSlowTimeMills() throws Exception { + BrokerController brokerController = new BrokerController( + new BrokerConfig(), + new NettyServerConfig(), + new NettyClientConfig(), + new MessageStoreConfig()); + brokerController.initialize(); + BlockingQueue queue = new LinkedBlockingQueue<>(); + + //create task is not instance of FutureTaskExt; + Runnable runnable = new Runnable() { + @Override + public void run() { + + } + }; + queue.add(runnable); + + RequestTask requestTask = new RequestTask(runnable, null, null); + // the requestTask is not the head of queue; + queue.add(new FutureTaskExt<>(requestTask, null)); + + long headSlowTimeMills = 100; + TimeUnit.MILLISECONDS.sleep(headSlowTimeMills); + assertThat(brokerController.headSlowTimeMills(queue)).isGreaterThanOrEqualTo(headSlowTimeMills); + //Attention: if we use the previous version method BrokerController#headSlowTimeMills, it will return 0; + } }