This is an automated email from the ASF dual-hosted git repository.
albumenj pushed a commit to branch 3.2
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.2 by this push:
new e92f85d980 Support client gracefully shutdown (#12504)
e92f85d980 is described below
commit e92f85d980a80126f547cc3ea549c126034dde73
Author: Albumen Kevin <[email protected]>
AuthorDate: Mon Jun 12 20:27:55 2023 +0800
Support client gracefully shutdown (#12504)
* Support client gracefully shutdown
* fix check
* Fix init
---
.../java/org/apache/dubbo/remoting/Constants.java | 2 ++
.../remoting/exchange/support/DefaultFuture.java | 33 ++++++++++++++++------
.../support/header/HeaderExchangeChannel.java | 10 ++++++-
.../support/header/HeaderExchangeHandler.java | 10 ++++++-
.../exchange/support/DefaultFutureTest.java | 2 +-
5 files changed, 46 insertions(+), 11 deletions(-)
diff --git
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/Constants.java
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/Constants.java
index bd2ad05220..a5f7fcebef 100644
---
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/Constants.java
+++
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/Constants.java
@@ -134,6 +134,8 @@ public interface Constants {
int DEFAULT_RECONNECT_PERIOD = 2000;
+ String CHANNEL_SHUTDOWN_TIMEOUT_KEY = "channel.shutdown.timeout";
+
String SEND_RECONNECT_KEY = "send.reconnect";
String CHECK_KEY = "check";
diff --git
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/DefaultFuture.java
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/DefaultFuture.java
index 3e38371d1e..3539723adb 100644
---
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/DefaultFuture.java
+++
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/DefaultFuture.java
@@ -19,6 +19,7 @@ package org.apache.dubbo.remoting.exchange.support;
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.resource.GlobalResourceInitializer;
+import org.apache.dubbo.common.serialize.SerializationException;
import org.apache.dubbo.common.timer.HashedWheelTimer;
import org.apache.dubbo.common.timer.Timeout;
import org.apache.dubbo.common.timer.Timer;
@@ -26,7 +27,6 @@ import org.apache.dubbo.common.timer.TimerTask;
import org.apache.dubbo.common.utils.NamedThreadFactory;
import org.apache.dubbo.remoting.Channel;
import org.apache.dubbo.remoting.RemotingException;
-import org.apache.dubbo.common.serialize.SerializationException;
import org.apache.dubbo.remoting.TimeoutException;
import org.apache.dubbo.remoting.exchange.Request;
import org.apache.dubbo.remoting.exchange.Response;
@@ -150,22 +150,39 @@ public class DefaultFuture extends
CompletableFuture<Object> {
*
* @param channel channel to close
*/
- public static void closeChannel(Channel channel) {
+ public static void closeChannel(Channel channel, long timeout) {
+ long deadline = timeout > 0 ? System.currentTimeMillis() + timeout : 0;
for (Map.Entry<Long, Channel> entry : CHANNELS.entrySet()) {
if (channel.equals(entry.getValue())) {
DefaultFuture future = getFuture(entry.getKey());
if (future != null && !future.isDone()) {
- Response disconnectResponse = new Response(future.getId());
- disconnectResponse.setStatus(Response.CHANNEL_INACTIVE);
- disconnectResponse.setErrorMessage("Channel " +
- channel + " is inactive. Directly return the
unFinished request : " +
- (logger.isDebugEnabled() ? future.getRequest() :
future.getRequest().copyWithoutData()));
- DefaultFuture.received(channel, disconnectResponse);
+ long restTime = deadline - System.currentTimeMillis();
+ if (restTime > 0) {
+ try {
+ future.get(restTime, TimeUnit.MILLISECONDS);
+ } catch (java.util.concurrent.TimeoutException ignore)
{
+ logger.warn(PROTOCOL_TIMEOUT_SERVER, "", "",
+ "Trying to close channel " + channel + ", but
response is not received in "
+ + timeout + "ms, and the request id is " +
future.id);
+ } catch (Throwable ignore) {}
+ }
+ if (!future.isDone()) {
+ respInactive(channel, future);
+ }
}
}
}
}
+ private static void respInactive(Channel channel, DefaultFuture future) {
+ Response disconnectResponse = new Response(future.getId());
+ disconnectResponse.setStatus(Response.CHANNEL_INACTIVE);
+ disconnectResponse.setErrorMessage("Channel " +
+ channel + " is inactive. Directly return the unFinished request :
" +
+ (logger.isDebugEnabled() ? future.getRequest() :
future.getRequest().copyWithoutData()));
+ DefaultFuture.received(channel, disconnectResponse);
+ }
+
public static void received(Channel channel, Response response) {
received(channel, response, false);
}
diff --git
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeChannel.java
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeChannel.java
index 6c4e7fd09c..5bb08fa8e2 100644
---
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeChannel.java
+++
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeChannel.java
@@ -18,6 +18,7 @@ package org.apache.dubbo.remoting.exchange.support.header;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.Version;
+import org.apache.dubbo.common.config.ConfigurationUtils;
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.remoting.Channel;
@@ -30,6 +31,7 @@ import org.apache.dubbo.remoting.exchange.Response;
import org.apache.dubbo.remoting.exchange.support.DefaultFuture;
import java.net.InetSocketAddress;
+import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
@@ -48,6 +50,8 @@ final class HeaderExchangeChannel implements ExchangeChannel {
private final Channel channel;
+ private final long shutdownTimeout;
+
private volatile boolean closed = false;
HeaderExchangeChannel(Channel channel) {
@@ -55,6 +59,10 @@ final class HeaderExchangeChannel implements ExchangeChannel
{
throw new IllegalArgumentException("channel == null");
}
this.channel = channel;
+ this.shutdownTimeout = Optional.ofNullable(channel.getUrl())
+ .map(URL::getOrDefaultApplicationModel)
+ .map(ConfigurationUtils::getServerShutdownTimeout)
+ .orElse(DEFAULT_TIMEOUT);
}
static HeaderExchangeChannel getOrAddChannel(Channel ch) {
@@ -159,7 +167,7 @@ final class HeaderExchangeChannel implements
ExchangeChannel {
closed = true;
try {
// graceful close
- DefaultFuture.closeChannel(channel);
+ DefaultFuture.closeChannel(channel, shutdownTimeout);
} catch (Exception e) {
logger.warn(TRANSPORT_FAILED_CLOSE, "", "", e.getMessage(), e);
}
diff --git
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeHandler.java
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeHandler.java
index fd9d660d05..2fe6253140 100644
---
a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeHandler.java
+++
b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/support/header/HeaderExchangeHandler.java
@@ -17,6 +17,7 @@
package org.apache.dubbo.remoting.exchange.support.header;
import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.config.ConfigurationUtils;
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.utils.NetUtils;
@@ -132,6 +133,8 @@ public class HeaderExchangeHandler implements
ChannelHandlerDelegate {
public void connected(Channel channel) throws RemotingException {
ExchangeChannel exchangeChannel =
HeaderExchangeChannel.getOrAddChannel(channel);
handler.connected(exchangeChannel);
+ channel.setAttribute(Constants.CHANNEL_SHUTDOWN_TIMEOUT_KEY,
+
ConfigurationUtils.getServerShutdownTimeout(channel.getUrl().getOrDefaultApplicationModel()));
}
@Override
@@ -140,7 +143,12 @@ public class HeaderExchangeHandler implements
ChannelHandlerDelegate {
try {
handler.disconnected(exchangeChannel);
} finally {
- DefaultFuture.closeChannel(channel);
+ int shutdownTimeout = 0;
+ Object timeoutObj =
channel.getAttribute(Constants.CHANNEL_SHUTDOWN_TIMEOUT_KEY);
+ if (timeoutObj instanceof Integer) {
+ shutdownTimeout = (Integer) timeoutObj;
+ }
+ DefaultFuture.closeChannel(channel, shutdownTimeout);
HeaderExchangeChannel.removeChannel(channel);
}
}
diff --git
a/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/exchange/support/DefaultFutureTest.java
b/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/exchange/support/DefaultFutureTest.java
index 9aec7713e7..1e28237370 100644
---
a/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/exchange/support/DefaultFutureTest.java
+++
b/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/exchange/support/DefaultFutureTest.java
@@ -167,7 +167,7 @@ class DefaultFutureTest {
ExecutorService executor =
ExtensionLoader.getExtensionLoader(ExecutorRepository.class)
.getDefaultExtension().createExecutorIfAbsent(URL.valueOf("dubbo://127.0.0.1:23456"));
DefaultFuture.newFuture(channel, request, 1000, executor);
- DefaultFuture.closeChannel(channel);
+ DefaultFuture.closeChannel(channel, 0);
Assertions.assertFalse(executor.isTerminated());
}