mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-19 02:23:24 +08:00
[ISSUE #6440] Optimize the code of consumer thread name,and support tag the name of scheduledExecutor thread. (#6441)
Co-authored-by: cyu <dingchengyu@hydee.cn>
This commit is contained in:
+4
-9
@@ -69,22 +69,17 @@ public class ConsumeMessageConcurrentlyService implements ConsumeMessageService
|
||||
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
|
||||
this.consumeRequestQueue = new LinkedBlockingQueue<>();
|
||||
|
||||
String consumeThreadPrefix = null;
|
||||
if (consumerGroup.length() > 100) {
|
||||
consumeThreadPrefix = new StringBuilder("ConsumeMessageThread_").append(consumerGroup, 0, 100).append("_").toString();
|
||||
} else {
|
||||
consumeThreadPrefix = new StringBuilder("ConsumeMessageThread_").append(consumerGroup).append("_").toString();
|
||||
}
|
||||
String consumerGroupTag = (consumerGroup.length() > 100 ? consumerGroup.substring(0, 100) : consumerGroup) + "_";
|
||||
this.consumeExecutor = new ThreadPoolExecutor(
|
||||
this.defaultMQPushConsumer.getConsumeThreadMin(),
|
||||
this.defaultMQPushConsumer.getConsumeThreadMax(),
|
||||
1000 * 60,
|
||||
TimeUnit.MILLISECONDS,
|
||||
this.consumeRequestQueue,
|
||||
new ThreadFactoryImpl(consumeThreadPrefix));
|
||||
new ThreadFactoryImpl("ConsumeMessageThread_" + consumerGroupTag));
|
||||
|
||||
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_"));
|
||||
this.cleanExpireMsgExecutors = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("CleanExpireMsgScheduledThread_"));
|
||||
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_" + consumerGroupTag));
|
||||
this.cleanExpireMsgExecutors = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("CleanExpireMsgScheduledThread_" + consumerGroupTag));
|
||||
}
|
||||
|
||||
public void start() {
|
||||
|
||||
+3
-8
@@ -73,21 +73,16 @@ public class ConsumeMessageOrderlyService implements ConsumeMessageService {
|
||||
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
|
||||
this.consumeRequestQueue = new LinkedBlockingQueue<>();
|
||||
|
||||
String consumeThreadPrefix = null;
|
||||
if (consumerGroup.length() > 100) {
|
||||
consumeThreadPrefix = new StringBuilder("ConsumeMessageThread_").append(consumerGroup.substring(0, 100)).append("_").toString();
|
||||
} else {
|
||||
consumeThreadPrefix = new StringBuilder("ConsumeMessageThread_").append(consumerGroup).append("_").toString();
|
||||
}
|
||||
String consumerGroupTag = (consumerGroup.length() > 100 ? consumerGroup.substring(0, 100) : consumerGroup) + "_";
|
||||
this.consumeExecutor = new ThreadPoolExecutor(
|
||||
this.defaultMQPushConsumer.getConsumeThreadMin(),
|
||||
this.defaultMQPushConsumer.getConsumeThreadMax(),
|
||||
1000 * 60,
|
||||
TimeUnit.MILLISECONDS,
|
||||
this.consumeRequestQueue,
|
||||
new ThreadFactoryImpl(consumeThreadPrefix));
|
||||
new ThreadFactoryImpl("ConsumeMessageThread_" + consumerGroupTag));
|
||||
|
||||
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_"));
|
||||
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_" + consumerGroupTag));
|
||||
}
|
||||
|
||||
public void start() {
|
||||
|
||||
Reference in New Issue
Block a user