cherry pick pr7109 to 3.0 (#7950)
* cherrypick * remove unused import * modify comment, remove issue-id
This commit is contained in:
parent
d5b46c3e65
commit
259dff2d04
|
|
@ -16,6 +16,15 @@
|
|||
*/
|
||||
package org.apache.dubbo.common.threadpool.manager;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.extension.ExtensionLoader;
|
||||
import org.apache.dubbo.common.logger.Logger;
|
||||
import org.apache.dubbo.common.logger.LoggerFactory;
|
||||
import org.apache.dubbo.common.threadlocal.NamedInternalThreadFactory;
|
||||
import org.apache.dubbo.common.threadpool.ThreadPool;
|
||||
import org.apache.dubbo.common.utils.ExecutorUtil;
|
||||
import org.apache.dubbo.common.utils.NamedThreadFactory;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
|
|
@ -26,16 +35,9 @@ import java.util.concurrent.ScheduledExecutorService;
|
|||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.apache.dubbo.common.URL;
|
||||
import org.apache.dubbo.common.extension.ExtensionLoader;
|
||||
import org.apache.dubbo.common.logger.Logger;
|
||||
import org.apache.dubbo.common.logger.LoggerFactory;
|
||||
import org.apache.dubbo.common.threadlocal.NamedInternalThreadFactory;
|
||||
import org.apache.dubbo.common.threadpool.ThreadPool;
|
||||
import org.apache.dubbo.common.utils.NamedThreadFactory;
|
||||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.CONSUMER_SIDE;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.EXECUTOR_SERVICE_COMPONENT_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.SIDE_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.THREADS_KEY;
|
||||
|
||||
/**
|
||||
|
|
@ -94,12 +96,9 @@ public class DefaultExecutorRepository implements ExecutorRepository {
|
|||
* @return
|
||||
*/
|
||||
public synchronized ExecutorService createExecutorIfAbsent(URL url) {
|
||||
String componentKey = EXECUTOR_SERVICE_COMPONENT_KEY;
|
||||
if (CONSUMER_SIDE.equalsIgnoreCase(url.getSide())) {
|
||||
componentKey = CONSUMER_SIDE;
|
||||
}
|
||||
Map<Integer, ExecutorService> executors = data.computeIfAbsent(componentKey, k -> new ConcurrentHashMap<>());
|
||||
Integer portKey = url.getPort();
|
||||
Map<Integer, ExecutorService> executors = data.computeIfAbsent(EXECUTOR_SERVICE_COMPONENT_KEY, k -> new ConcurrentHashMap<>());
|
||||
// Consumer's executor is sharing globally, key=Integer.MAX_VALUE. Provider's executor is sharing by protocol.
|
||||
Integer portKey = CONSUMER_SIDE.equalsIgnoreCase(url.getParameter(SIDE_KEY)) ? Integer.MAX_VALUE : url.getPort();
|
||||
ExecutorService executor = executors.computeIfAbsent(portKey, k -> createExecutor(url));
|
||||
// If executor has been shut down, create a new one
|
||||
if (executor.isShutdown() || executor.isTerminated()) {
|
||||
|
|
@ -111,11 +110,7 @@ public class DefaultExecutorRepository implements ExecutorRepository {
|
|||
}
|
||||
|
||||
public ExecutorService getExecutor(URL url) {
|
||||
String componentKey = EXECUTOR_SERVICE_COMPONENT_KEY;
|
||||
if (CONSUMER_SIDE.equalsIgnoreCase(url.getSide())) {
|
||||
componentKey = CONSUMER_SIDE;
|
||||
}
|
||||
Map<Integer, ExecutorService> executors = data.get(componentKey);
|
||||
Map<Integer, ExecutorService> executors = data.get(EXECUTOR_SERVICE_COMPONENT_KEY);
|
||||
|
||||
/**
|
||||
* It's guaranteed that this method is called after {@link #createExecutorIfAbsent(URL)}, so data should already
|
||||
|
|
@ -127,16 +122,20 @@ public class DefaultExecutorRepository implements ExecutorRepository {
|
|||
return null;
|
||||
}
|
||||
|
||||
Integer portKey = url.getPort();
|
||||
// Consumer's executor is sharing globally, key=Integer.MAX_VALUE. Provider's executor is sharing by protocol.
|
||||
Integer portKey = CONSUMER_SIDE.equalsIgnoreCase(url.getParameter(SIDE_KEY)) ? Integer.MAX_VALUE : url.getPort();
|
||||
ExecutorService executor = executors.get(portKey);
|
||||
if (executor != null) {
|
||||
if (executor.isShutdown() || executor.isTerminated()) {
|
||||
executors.remove(portKey);
|
||||
executor = createExecutor(url);
|
||||
executors.put(portKey, executor);
|
||||
}
|
||||
if (executor != null && (executor.isShutdown() || executor.isTerminated())) {
|
||||
executors.remove(portKey);
|
||||
// Does not re-create a shutdown executor, use SHARED_EXECUTOR for downgrade.
|
||||
executor = null;
|
||||
logger.info("Executor for " + url + " is shutdown.");
|
||||
}
|
||||
if (executor == null) {
|
||||
return SHARED_EXECUTOR;
|
||||
} else {
|
||||
return executor;
|
||||
}
|
||||
return executor;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -209,4 +208,17 @@ public class DefaultExecutorRepository implements ExecutorRepository {
|
|||
public ExecutorService getPoolRouterExecutor() {
|
||||
return poolRouterExecutor;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroyAll() {
|
||||
data.values().forEach(executors -> {
|
||||
if (executors != null) {
|
||||
executors.values().forEach(executor -> {
|
||||
if (executor != null && !executor.isShutdown()) {
|
||||
ExecutorUtil.shutdownNow(executor, 100);
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -79,4 +79,8 @@ public interface ExecutorRepository {
|
|||
|
||||
ExecutorService getPoolRouterExecutor();
|
||||
|
||||
/**
|
||||
* Destroy all executors that are not in shutdown state
|
||||
*/
|
||||
void destroyAll();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -39,8 +39,6 @@ public class ExecutorRepositoryTest {
|
|||
}
|
||||
|
||||
private void testGet(URL url) {
|
||||
Assertions.assertNull(executorRepository.getExecutor(url));
|
||||
|
||||
ExecutorService executorService = executorRepository.createExecutorIfAbsent(url);
|
||||
executorService.shutdown();
|
||||
executorService = executorRepository.createExecutorIfAbsent(url);
|
||||
|
|
|
|||
|
|
@ -1225,7 +1225,7 @@ public class DubboBootstrap {
|
|||
destroyRegistries();
|
||||
|
||||
destroyServiceDiscoveries();
|
||||
|
||||
destroyExecutorRepository();
|
||||
clear();
|
||||
shutdown();
|
||||
release();
|
||||
|
|
@ -1238,6 +1238,10 @@ public class DubboBootstrap {
|
|||
}
|
||||
}
|
||||
|
||||
private void destroyExecutorRepository() {
|
||||
ExtensionLoader.getExtensionLoader(ExecutorRepository.class).getDefaultExtension().destroyAll();
|
||||
}
|
||||
|
||||
private void destroyRegistries() {
|
||||
AbstractRegistryFactory.destroyAll();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ import org.apache.dubbo.common.extension.ExtensionLoader;
|
|||
import org.apache.dubbo.common.logger.Logger;
|
||||
import org.apache.dubbo.common.logger.LoggerFactory;
|
||||
import org.apache.dubbo.common.threadpool.manager.ExecutorRepository;
|
||||
import org.apache.dubbo.common.utils.ExecutorUtil;
|
||||
import org.apache.dubbo.common.utils.NetUtils;
|
||||
import org.apache.dubbo.remoting.Channel;
|
||||
import org.apache.dubbo.remoting.ChannelHandler;
|
||||
|
|
@ -38,6 +37,7 @@ import java.util.concurrent.locks.ReentrantLock;
|
|||
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_CLIENT_THREADPOOL;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.THREADPOOL_KEY;
|
||||
import static org.apache.dubbo.common.constants.CommonConstants.THREAD_NAME_KEY;
|
||||
|
||||
/**
|
||||
* AbstractClient
|
||||
|
|
@ -90,7 +90,8 @@ public abstract class AbstractClient extends AbstractEndpoint implements Client
|
|||
}
|
||||
|
||||
private void initExecutor(URL url) {
|
||||
url = ExecutorUtil.setThreadName(url, CLIENT_THREAD_POOL_NAME);
|
||||
// Consumer's executor is sharing globally, thread name not require provider ip.
|
||||
url = url.addParameter(THREAD_NAME_KEY, CLIENT_THREAD_POOL_NAME);
|
||||
url = url.addParameterIfAbsent(THREADPOOL_KEY, DEFAULT_CLIENT_THREADPOOL);
|
||||
executor = executorRepository.createExecutorIfAbsent(url);
|
||||
}
|
||||
|
|
@ -277,14 +278,6 @@ public abstract class AbstractClient extends AbstractEndpoint implements Client
|
|||
logger.warn(e.getMessage(), e);
|
||||
}
|
||||
|
||||
try {
|
||||
if (executor != null) {
|
||||
ExecutorUtil.shutdownNow(executor, 100);
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
logger.warn(e.getMessage(), e);
|
||||
}
|
||||
|
||||
try {
|
||||
disconnect();
|
||||
} catch (Throwable e) {
|
||||
|
|
@ -304,7 +297,6 @@ public abstract class AbstractClient extends AbstractEndpoint implements Client
|
|||
|
||||
@Override
|
||||
public void close(int timeout) {
|
||||
ExecutorUtil.gracefulShutdown(executor, timeout);
|
||||
close();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -43,7 +43,7 @@ public class ThreadNameTest {
|
|||
private ThreadNameVerifyHandler clientHandler;
|
||||
|
||||
private static String serverRegex = "DubboServerHandler\\-localhost:(\\d+)\\-thread\\-(\\d+)";
|
||||
private static String clientRegex = "DubboClientHandler\\-localhost:(\\d+)\\-thread\\-(\\d+)";
|
||||
private static String clientRegex = "DubboClientHandler\\-thread\\-(\\d+)";
|
||||
|
||||
private final CountDownLatch serverLatch = new CountDownLatch(1);
|
||||
private final CountDownLatch clientLatch = new CountDownLatch(1);
|
||||
|
|
|
|||
Loading…
Reference in New Issue