mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
Fix bug that the broker will hang after merge the pr that fix the headWaitTimeMills of sendThreadPoolQueue (#3631)
This commit is contained in:
@@ -33,8 +33,6 @@ 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;
|
||||
@@ -134,6 +132,7 @@ public class BrokerController {
|
||||
"BrokerControllerScheduledThread"));
|
||||
private final SlaveSynchronize slaveSynchronize;
|
||||
private final BlockingQueue<Runnable> sendThreadPoolQueue;
|
||||
private final BlockingQueue<Runnable> putThreadPoolQueue;
|
||||
private final BlockingQueue<Runnable> pullThreadPoolQueue;
|
||||
private final BlockingQueue<Runnable> replyThreadPoolQueue;
|
||||
private final BlockingQueue<Runnable> queryThreadPoolQueue;
|
||||
@@ -150,6 +149,7 @@ public class BrokerController {
|
||||
private RemotingServer fastRemotingServer;
|
||||
private TopicConfigManager topicConfigManager;
|
||||
private ExecutorService sendMessageExecutor;
|
||||
private ExecutorService putMessageFutureExecutor;
|
||||
private ExecutorService pullMessageExecutor;
|
||||
private ExecutorService replyMessageExecutor;
|
||||
private ExecutorService queryMessageExecutor;
|
||||
@@ -198,6 +198,7 @@ public class BrokerController {
|
||||
this.slaveSynchronize = new SlaveSynchronize(this);
|
||||
|
||||
this.sendThreadPoolQueue = new LinkedBlockingQueue<Runnable>(this.brokerConfig.getSendThreadPoolQueueCapacity());
|
||||
this.putThreadPoolQueue = new LinkedBlockingQueue<Runnable>(this.brokerConfig.getPutThreadPoolQueueCapacity());
|
||||
this.pullThreadPoolQueue = new LinkedBlockingQueue<Runnable>(this.brokerConfig.getPullThreadPoolQueueCapacity());
|
||||
this.replyThreadPoolQueue = new LinkedBlockingQueue<Runnable>(this.brokerConfig.getReplyThreadPoolQueueCapacity());
|
||||
this.queryThreadPoolQueue = new LinkedBlockingQueue<Runnable>(this.brokerConfig.getQueryThreadPoolQueueCapacity());
|
||||
@@ -275,6 +276,14 @@ public class BrokerController {
|
||||
this.sendThreadPoolQueue,
|
||||
new ThreadFactoryImpl("SendMessageThread_"));
|
||||
|
||||
this.putMessageFutureExecutor = new BrokerFixedThreadPoolExecutor(
|
||||
this.brokerConfig.getPutMessageFutureThreadPoolNums(),
|
||||
this.brokerConfig.getPutMessageFutureThreadPoolNums(),
|
||||
1000 * 60,
|
||||
TimeUnit.MILLISECONDS,
|
||||
this.putThreadPoolQueue,
|
||||
new ThreadFactoryImpl("PutMessageThread_"));
|
||||
|
||||
this.pullMessageExecutor = new BrokerFixedThreadPoolExecutor(
|
||||
this.brokerConfig.getPullMessageThreadPoolNums(),
|
||||
this.brokerConfig.getPullMessageThreadPoolNums(),
|
||||
@@ -652,13 +661,10 @@ public class BrokerController {
|
||||
|
||||
public long headSlowTimeMills(BlockingQueue<Runnable> q) {
|
||||
long slowTimeMills = 0;
|
||||
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();
|
||||
final Runnable peek = q.peek();
|
||||
if (peek != null) {
|
||||
RequestTask rt = BrokerFastFailure.castRunnable(peek);
|
||||
slowTimeMills = rt == null ? 0 : this.messageStore.now() - rt.getCreateTimestamp();
|
||||
}
|
||||
|
||||
if (slowTimeMills < 0) {
|
||||
@@ -787,6 +793,10 @@ public class BrokerController {
|
||||
this.sendMessageExecutor.shutdown();
|
||||
}
|
||||
|
||||
if (this.putMessageFutureExecutor != null) {
|
||||
this.putMessageFutureExecutor.shutdown();
|
||||
}
|
||||
|
||||
if (this.pullMessageExecutor != null) {
|
||||
this.pullMessageExecutor.shutdown();
|
||||
}
|
||||
@@ -1245,7 +1255,7 @@ public class BrokerController {
|
||||
}
|
||||
}
|
||||
|
||||
public ExecutorService getSendMessageExecutor() {
|
||||
return sendMessageExecutor;
|
||||
public ExecutorService getPutMessageFutureExecutor() {
|
||||
return putMessageFutureExecutor;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,7 +84,7 @@ public class SendMessageProcessor extends AbstractSendMessageProcessor implement
|
||||
|
||||
@Override
|
||||
public void asyncProcessRequest(ChannelHandlerContext ctx, RemotingCommand request, RemotingResponseCallback responseCallback) throws Exception {
|
||||
asyncProcessRequest(ctx, request).thenAcceptAsync(responseCallback::callback, this.brokerController.getSendMessageExecutor());
|
||||
asyncProcessRequest(ctx, request).thenAcceptAsync(responseCallback::callback, this.brokerController.getPutMessageFutureExecutor());
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> asyncProcessRequest(ChannelHandlerContext ctx,
|
||||
|
||||
@@ -71,7 +71,6 @@ public class BrokerControllerTest {
|
||||
|
||||
}
|
||||
};
|
||||
queue.add(runnable);
|
||||
|
||||
RequestTask requestTask = new RequestTask(runnable, null, null);
|
||||
// the requestTask is not the head of queue;
|
||||
@@ -80,6 +79,5 @@ public class BrokerControllerTest {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,6 +60,7 @@ public class BrokerConfig {
|
||||
* thread numbers for send message thread pool.
|
||||
*/
|
||||
private int sendMessageThreadPoolNums = Math.min(Runtime.getRuntime().availableProcessors(), 4);
|
||||
private int putMessageFutureThreadPoolNums = Math.min(Runtime.getRuntime().availableProcessors(), 4);
|
||||
private int pullMessageThreadPoolNums = 16 + Runtime.getRuntime().availableProcessors() * 2;
|
||||
private int processReplyMessageThreadPoolNums = 16 + Runtime.getRuntime().availableProcessors() * 2;
|
||||
private int queryMessageThreadPoolNums = 8 + Runtime.getRuntime().availableProcessors();
|
||||
@@ -84,6 +85,7 @@ public class BrokerConfig {
|
||||
@ImportantField
|
||||
private boolean fetchNamesrvAddrByAddressServer = false;
|
||||
private int sendThreadPoolQueueCapacity = 10000;
|
||||
private int putThreadPoolQueueCapacity = 10000;
|
||||
private int pullThreadPoolQueueCapacity = 100000;
|
||||
private int replyThreadPoolQueueCapacity = 10000;
|
||||
private int queryThreadPoolQueueCapacity = 20000;
|
||||
@@ -375,6 +377,14 @@ public class BrokerConfig {
|
||||
this.sendMessageThreadPoolNums = sendMessageThreadPoolNums;
|
||||
}
|
||||
|
||||
public int getPutMessageFutureThreadPoolNums() {
|
||||
return putMessageFutureThreadPoolNums;
|
||||
}
|
||||
|
||||
public void setPutMessageFutureThreadPoolNums(int putMessageFutureThreadPoolNums) {
|
||||
this.putMessageFutureThreadPoolNums = putMessageFutureThreadPoolNums;
|
||||
}
|
||||
|
||||
public int getPullMessageThreadPoolNums() {
|
||||
return pullMessageThreadPoolNums;
|
||||
}
|
||||
@@ -479,6 +489,14 @@ public class BrokerConfig {
|
||||
this.sendThreadPoolQueueCapacity = sendThreadPoolQueueCapacity;
|
||||
}
|
||||
|
||||
public int getPutThreadPoolQueueCapacity() {
|
||||
return putThreadPoolQueueCapacity;
|
||||
}
|
||||
|
||||
public void setPutThreadPoolQueueCapacity(int putThreadPoolQueueCapacity) {
|
||||
this.putThreadPoolQueueCapacity = putThreadPoolQueueCapacity;
|
||||
}
|
||||
|
||||
public int getPullThreadPoolQueueCapacity() {
|
||||
return pullThreadPoolQueueCapacity;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user