mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
This commit is contained in:
@@ -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<Runnable> 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<RequestTask> 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) {
|
||||
|
||||
@@ -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<Runnable> 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;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user