diff --git a/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java b/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java
index 5df5cc8fa1a..9aef301f50d 100644
--- a/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java
+++ b/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java
@@ -47,6 +47,7 @@
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
+import java.util.concurrent.ExecutorService;
/**
* In most scenarios, this is the mostly recommended class to consume messages.
@@ -160,6 +161,8 @@ public class DefaultMQPushConsumer extends ClientConfig implements MQPushConsume
*/
private int consumeThreadMin = 20;
+ private ExecutorService consumeExecutor;
+
/**
* Max consumer thread number
*/
@@ -558,6 +561,37 @@ public void setConsumerGroup(String consumerGroup) {
this.consumerGroup = consumerGroup;
}
+ /**
+ * Returns the externally managed consumption executor, or null for a dedicated pool.
+ */
+ public ExecutorService getConsumeExecutor() {
+ return consumeExecutor;
+ }
+
+ /**
+ * Sets an externally managed executor before starting this consumer.
+ *
+ *
This is an advanced API intended for controlled integrations such as Proxy. Ordinary
+ * applications should use the default consumption pool instead of injecting an executor.
+ * The executor may be shared with other consumers. Virtual-thread executors are also supported
+ * when supplied by applications running on a compatible JDK.
+ *
+ *
While consumers are running, the external executor must avoid capacity-based rejection
+ * and must not discard or cancel pending consumption tasks. The client does not guarantee
+ * automatic recovery from rejected tasks. Discarding or cancelling tasks can retain cached
+ * messages and pin consumption offsets, eventually stalling consumption.
+ *
+ *
The caller controls concurrency and owns the executor's lifecycle. This consumer never
+ * shuts down or resizes an external executor. Consumer shutdown does not await or cancel tasks
+ * submitted to it; the caller must stop all consumers using the executor before shutting it
+ * down and awaiting its termination.
+ *
+ * @param consumeExecutor external executor, or null to use the default dedicated pool
+ */
+ public void setConsumeExecutor(ExecutorService consumeExecutor) {
+ this.consumeExecutor = consumeExecutor;
+ }
+
public int getConsumeThreadMax() {
return consumeThreadMax;
}
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/AbstractConsumeMessageService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/AbstractConsumeMessageService.java
new file mode 100644
index 00000000000..27887989bcb
--- /dev/null
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/AbstractConsumeMessageService.java
@@ -0,0 +1,81 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.client.impl.consumer;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
+import org.apache.rocketmq.common.utils.ThreadUtils;
+
+public abstract class AbstractConsumeMessageService implements ConsumeMessageService {
+ protected final DefaultMQPushConsumer defaultMQPushConsumer;
+ protected final ExecutorService consumeExecutor;
+ private final boolean ownsConsumeExecutor;
+
+ protected AbstractConsumeMessageService(DefaultMQPushConsumer defaultMQPushConsumer, ThreadFactory threadFactory) {
+ this.defaultMQPushConsumer = defaultMQPushConsumer;
+ ExecutorService externalExecutor = defaultMQPushConsumer.getConsumeExecutor();
+ this.ownsConsumeExecutor = externalExecutor == null;
+ if (this.ownsConsumeExecutor) {
+ this.consumeExecutor = new ThreadPoolExecutor(
+ defaultMQPushConsumer.getConsumeThreadMin(),
+ defaultMQPushConsumer.getConsumeThreadMax(),
+ 1000 * 60,
+ TimeUnit.MILLISECONDS,
+ new LinkedBlockingQueue<>(),
+ threadFactory);
+ } else {
+ this.consumeExecutor = externalExecutor;
+ }
+ }
+
+ protected static String getConsumerGroupTag(String consumerGroup) {
+ return (consumerGroup.length() > 100 ? consumerGroup.substring(0, 100) : consumerGroup) + "_";
+ }
+
+ protected void shutdownConsumeExecutor(long awaitTerminateMillis) {
+ if (this.ownsConsumeExecutor) {
+ ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
+ }
+ }
+
+ @Override
+ public void updateCorePoolSize(int corePoolSize) {
+ if (this.ownsConsumeExecutor
+ && corePoolSize > 0
+ && corePoolSize <= Short.MAX_VALUE
+ && corePoolSize < this.defaultMQPushConsumer.getConsumeThreadMax()) {
+ ((ThreadPoolExecutor) this.consumeExecutor).setCorePoolSize(corePoolSize);
+ }
+ }
+
+ @Override
+ public void incCorePoolSize() {
+ }
+
+ @Override
+ public void decCorePoolSize() {
+ }
+
+ @Override
+ public int getCorePoolSize() {
+ return this.ownsConsumeExecutor ? ((ThreadPoolExecutor) this.consumeExecutor).getCorePoolSize() : -1;
+ }
+}
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
index b151fefbbb3..361da434ff2 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
@@ -22,14 +22,10 @@
import java.util.Iterator;
import java.util.List;
import java.util.Map;
-import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
-import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.ConsumeReturnType;
@@ -42,19 +38,15 @@
import org.apache.rocketmq.common.message.MessageAccessor;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.common.utils.ThreadUtils;
import org.apache.rocketmq.remoting.protocol.body.CMResult;
import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
-public class ConsumeMessageConcurrentlyService implements ConsumeMessageService {
+public class ConsumeMessageConcurrentlyService extends AbstractConsumeMessageService {
private static final Logger log = LoggerFactory.getLogger(ConsumeMessageConcurrentlyService.class);
private final DefaultMQPushConsumerImpl defaultMQPushConsumerImpl;
- private final DefaultMQPushConsumer defaultMQPushConsumer;
private final MessageListenerConcurrently messageListener;
- private final BlockingQueue consumeRequestQueue;
- private final ThreadPoolExecutor consumeExecutor;
private final String consumerGroup;
private final ScheduledExecutorService scheduledExecutorService;
@@ -62,22 +54,14 @@ public class ConsumeMessageConcurrentlyService implements ConsumeMessageService
public ConsumeMessageConcurrentlyService(DefaultMQPushConsumerImpl defaultMQPushConsumerImpl,
MessageListenerConcurrently messageListener) {
+ super(defaultMQPushConsumerImpl.getDefaultMQPushConsumer(), new ThreadFactoryImpl("ConsumeMessageThread_"
+ + getConsumerGroupTag(defaultMQPushConsumerImpl.getDefaultMQPushConsumer().getConsumerGroup())));
this.defaultMQPushConsumerImpl = defaultMQPushConsumerImpl;
this.messageListener = messageListener;
- this.defaultMQPushConsumer = this.defaultMQPushConsumerImpl.getDefaultMQPushConsumer();
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
- this.consumeRequestQueue = new LinkedBlockingQueue<>();
-
- 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("ConsumeMessageThread_" + consumerGroupTag));
+ String consumerGroupTag = getConsumerGroupTag(consumerGroup);
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_" + consumerGroupTag));
this.cleanExpireMsgExecutors = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("CleanExpireMsgScheduledThread_" + consumerGroupTag));
}
@@ -99,34 +83,10 @@ public void run() {
public void shutdown(long awaitTerminateMillis) {
this.scheduledExecutorService.shutdown();
- ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
+ shutdownConsumeExecutor(awaitTerminateMillis);
this.cleanExpireMsgExecutors.shutdown();
}
- @Override
- public void updateCorePoolSize(int corePoolSize) {
- if (corePoolSize > 0
- && corePoolSize <= Short.MAX_VALUE
- && corePoolSize < this.defaultMQPushConsumer.getConsumeThreadMax()) {
- this.consumeExecutor.setCorePoolSize(corePoolSize);
- }
- }
-
- @Override
- public void incCorePoolSize() {
-
- }
-
- @Override
- public void decCorePoolSize() {
-
- }
-
- @Override
- public int getCorePoolSize() {
- return this.consumeExecutor.getCorePoolSize();
- }
-
@Override
public ConsumeMessageDirectlyResult consumeMessageDirectly(MessageExt msg, String brokerName) {
ConsumeMessageDirectlyResult result = new ConsumeMessageDirectlyResult();
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyService.java
index 3ca465da70d..776eb129c75 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyService.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyService.java
@@ -20,14 +20,10 @@
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
-import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.commons.lang3.StringUtils;
-import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.ConsumeReturnType;
@@ -42,7 +38,6 @@
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.common.utils.ThreadUtils;
import org.apache.rocketmq.remoting.protocol.NamespaceUtil;
import org.apache.rocketmq.remoting.protocol.body.CMResult;
import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
@@ -50,15 +45,12 @@
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
-public class ConsumeMessageOrderlyService implements ConsumeMessageService {
+public class ConsumeMessageOrderlyService extends AbstractConsumeMessageService {
private static final Logger log = LoggerFactory.getLogger(ConsumeMessageOrderlyService.class);
private final static long MAX_TIME_CONSUME_CONTINUOUSLY =
Long.parseLong(System.getProperty("rocketmq.client.maxTimeConsumeContinuously", "60000"));
private final DefaultMQPushConsumerImpl defaultMQPushConsumerImpl;
- private final DefaultMQPushConsumer defaultMQPushConsumer;
private final MessageListenerOrderly messageListener;
- private final BlockingQueue consumeRequestQueue;
- private final ThreadPoolExecutor consumeExecutor;
private final String consumerGroup;
private final MessageQueueLock messageQueueLock = new MessageQueueLock();
private final ScheduledExecutorService scheduledExecutorService;
@@ -66,22 +58,14 @@ public class ConsumeMessageOrderlyService implements ConsumeMessageService {
public ConsumeMessageOrderlyService(DefaultMQPushConsumerImpl defaultMQPushConsumerImpl,
MessageListenerOrderly messageListener) {
+ super(defaultMQPushConsumerImpl.getDefaultMQPushConsumer(), new ThreadFactoryImpl("ConsumeMessageThread_"
+ + getConsumerGroupTag(defaultMQPushConsumerImpl.getDefaultMQPushConsumer().getConsumerGroup())));
this.defaultMQPushConsumerImpl = defaultMQPushConsumerImpl;
this.messageListener = messageListener;
- this.defaultMQPushConsumer = this.defaultMQPushConsumerImpl.getDefaultMQPushConsumer();
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
- this.consumeRequestQueue = new LinkedBlockingQueue<>();
-
- 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("ConsumeMessageThread_" + consumerGroupTag));
+ String consumerGroupTag = getConsumerGroupTag(consumerGroup);
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_" + consumerGroupTag));
}
@@ -105,7 +89,7 @@ public void run() {
public void shutdown(long awaitTerminateMillis) {
this.stopped = true;
this.scheduledExecutorService.shutdown();
- ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
+ shutdownConsumeExecutor(awaitTerminateMillis);
if (MessageModel.CLUSTERING.equals(this.defaultMQPushConsumerImpl.messageModel())) {
this.unlockAllMQ();
}
@@ -115,28 +99,6 @@ public synchronized void unlockAllMQ() {
this.defaultMQPushConsumerImpl.getRebalanceImpl().unlockAll(false);
}
- @Override
- public void updateCorePoolSize(int corePoolSize) {
- if (corePoolSize > 0
- && corePoolSize <= Short.MAX_VALUE
- && corePoolSize < this.defaultMQPushConsumer.getConsumeThreadMax()) {
- this.consumeExecutor.setCorePoolSize(corePoolSize);
- }
- }
-
- @Override
- public void incCorePoolSize() {
- }
-
- @Override
- public void decCorePoolSize() {
- }
-
- @Override
- public int getCorePoolSize() {
- return this.consumeExecutor.getCorePoolSize();
- }
-
@Override
public ConsumeMessageDirectlyResult consumeMessageDirectly(MessageExt msg, String brokerName) {
ConsumeMessageDirectlyResult result = new ConsumeMessageDirectlyResult();
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopConcurrentlyService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopConcurrentlyService.java
index d5191871106..9d70400903c 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopConcurrentlyService.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopConcurrentlyService.java
@@ -20,16 +20,12 @@
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
-import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.client.consumer.AckCallback;
import org.apache.rocketmq.client.consumer.AckResult;
-import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.ConsumeReturnType;
@@ -43,40 +39,27 @@
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.common.utils.ThreadUtils;
import org.apache.rocketmq.remoting.protocol.body.CMResult;
import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
import org.apache.rocketmq.remoting.protocol.header.ExtraInfoUtil;
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
-public class ConsumeMessagePopConcurrentlyService implements ConsumeMessageService {
+public class ConsumeMessagePopConcurrentlyService extends AbstractConsumeMessageService {
private static final Logger log = LoggerFactory.getLogger(ConsumeMessagePopConcurrentlyService.class);
private final DefaultMQPushConsumerImpl defaultMQPushConsumerImpl;
- private final DefaultMQPushConsumer defaultMQPushConsumer;
private final MessageListenerConcurrently messageListener;
- private final BlockingQueue consumeRequestQueue;
- private final ThreadPoolExecutor consumeExecutor;
private final String consumerGroup;
private final ScheduledExecutorService scheduledExecutorService;
public ConsumeMessagePopConcurrentlyService(DefaultMQPushConsumerImpl defaultMQPushConsumerImpl,
MessageListenerConcurrently messageListener) {
+ super(defaultMQPushConsumerImpl.getDefaultMQPushConsumer(), new ThreadFactoryImpl("ConsumeMessageThread_"));
this.defaultMQPushConsumerImpl = defaultMQPushConsumerImpl;
this.messageListener = messageListener;
- this.defaultMQPushConsumer = this.defaultMQPushConsumerImpl.getDefaultMQPushConsumer();
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
- this.consumeRequestQueue = new LinkedBlockingQueue<>();
-
- this.consumeExecutor = new ThreadPoolExecutor(
- this.defaultMQPushConsumer.getConsumeThreadMin(),
- this.defaultMQPushConsumer.getConsumeThreadMax(),
- 1000 * 60,
- TimeUnit.MILLISECONDS,
- this.consumeRequestQueue,
- new ThreadFactoryImpl("ConsumeMessageThread_"));
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_"));
}
@@ -86,32 +69,9 @@ public void start() {
public void shutdown(long awaitTerminateMillis) {
this.scheduledExecutorService.shutdown();
- ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
- }
-
- @Override
- public void updateCorePoolSize(int corePoolSize) {
- if (corePoolSize > 0
- && corePoolSize <= Short.MAX_VALUE
- && corePoolSize < this.defaultMQPushConsumer.getConsumeThreadMax()) {
- this.consumeExecutor.setCorePoolSize(corePoolSize);
- }
- }
-
- @Override
- public void incCorePoolSize() {
+ shutdownConsumeExecutor(awaitTerminateMillis);
}
- @Override
- public void decCorePoolSize() {
- }
-
- @Override
- public int getCorePoolSize() {
- return this.consumeExecutor.getCorePoolSize();
- }
-
-
@Override
public ConsumeMessageDirectlyResult consumeMessageDirectly(MessageExt msg, String brokerName) {
ConsumeMessageDirectlyResult result = new ConsumeMessageDirectlyResult();
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopOrderlyService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopOrderlyService.java
index 4eab1ccf664..8f6a5ee6dfd 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopOrderlyService.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessagePopOrderlyService.java
@@ -19,14 +19,10 @@
import io.netty.util.internal.ConcurrentSet;
import java.util.ArrayList;
import java.util.List;
-import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.commons.lang3.StringUtils;
-import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
@@ -39,7 +35,6 @@
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.common.utils.ThreadUtils;
import org.apache.rocketmq.remoting.protocol.NamespaceUtil;
import org.apache.rocketmq.remoting.protocol.body.CMResult;
import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
@@ -47,14 +42,11 @@
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
-public class ConsumeMessagePopOrderlyService implements ConsumeMessageService {
+public class ConsumeMessagePopOrderlyService extends AbstractConsumeMessageService {
private static final Logger log = LoggerFactory.getLogger(ConsumeMessagePopOrderlyService.class);
private final DefaultMQPushConsumerImpl defaultMQPushConsumerImpl;
- private final DefaultMQPushConsumer defaultMQPushConsumer;
private final MessageListenerOrderly messageListener;
- private final BlockingQueue consumeRequestQueue;
private final ConcurrentSet consumeRequestSet = new ConcurrentSet<>();
- private final ThreadPoolExecutor consumeExecutor;
private final String consumerGroup;
private final MessageQueueLock messageQueueLock = new MessageQueueLock();
private final MessageQueueLock consumeRequestLock = new MessageQueueLock();
@@ -63,20 +55,11 @@ public class ConsumeMessagePopOrderlyService implements ConsumeMessageService {
public ConsumeMessagePopOrderlyService(DefaultMQPushConsumerImpl defaultMQPushConsumerImpl,
MessageListenerOrderly messageListener) {
+ super(defaultMQPushConsumerImpl.getDefaultMQPushConsumer(), new ThreadFactoryImpl("ConsumeMessageThread_"));
this.defaultMQPushConsumerImpl = defaultMQPushConsumerImpl;
this.messageListener = messageListener;
- this.defaultMQPushConsumer = this.defaultMQPushConsumerImpl.getDefaultMQPushConsumer();
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
- this.consumeRequestQueue = new LinkedBlockingQueue<>();
-
- this.consumeExecutor = new ThreadPoolExecutor(
- this.defaultMQPushConsumer.getConsumeThreadMin(),
- this.defaultMQPushConsumer.getConsumeThreadMax(),
- 1000 * 60,
- TimeUnit.MILLISECONDS,
- this.consumeRequestQueue,
- new ThreadFactoryImpl("ConsumeMessageThread_"));
this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("ConsumeMessageScheduledThread_"));
}
@@ -97,7 +80,7 @@ public void run() {
public void shutdown(long awaitTerminateMillis) {
this.stopped = true;
this.scheduledExecutorService.shutdown();
- ThreadUtils.shutdownGracefully(this.consumeExecutor, awaitTerminateMillis, TimeUnit.MILLISECONDS);
+ shutdownConsumeExecutor(awaitTerminateMillis);
if (MessageModel.CLUSTERING.equals(this.defaultMQPushConsumerImpl.messageModel())) {
this.unlockAllMessageQueues();
}
@@ -107,28 +90,6 @@ public synchronized void unlockAllMessageQueues() {
this.defaultMQPushConsumerImpl.getRebalanceImpl().unlockAll(false);
}
- @Override
- public void updateCorePoolSize(int corePoolSize) {
- if (corePoolSize > 0
- && corePoolSize <= Short.MAX_VALUE
- && corePoolSize < this.defaultMQPushConsumer.getConsumeThreadMax()) {
- this.consumeExecutor.setCorePoolSize(corePoolSize);
- }
- }
-
- @Override
- public void incCorePoolSize() {
- }
-
- @Override
- public void decCorePoolSize() {
- }
-
- @Override
- public int getCorePoolSize() {
- return this.consumeExecutor.getCorePoolSize();
- }
-
@Override
public ConsumeMessageDirectlyResult consumeMessageDirectly(MessageExt msg, String brokerName) {
ConsumeMessageDirectlyResult result = new ConsumeMessageDirectlyResult();
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageService.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageService.java
index ee684730aed..1842b5eac15 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageService.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageService.java
@@ -32,6 +32,7 @@ public interface ConsumeMessageService {
void decCorePoolSize();
+ /** Returns the owned pool core size, or -1 when execution is managed externally. */
int getCorePoolSize();
ConsumeMessageDirectlyResult consumeMessageDirectly(final MessageExt msg, final String brokerName);
diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageExecutorInjectionTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageExecutorInjectionTest.java
new file mode 100644
index 00000000000..8dfffdf8ba6
--- /dev/null
+++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageExecutorInjectionTest.java
@@ -0,0 +1,203 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.client.impl.consumer;
+
+import java.lang.reflect.Method;
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.Future;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
+import org.apache.rocketmq.client.consumer.store.OffsetStore;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.junit.Assume;
+import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
+import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
+import org.apache.rocketmq.remoting.protocol.heartbeat.MessageModel;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+@RunWith(Parameterized.class)
+public class ConsumeMessageExecutorInjectionTest {
+ @Parameterized.Parameters(name = "{0}")
+ public static Collection