diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java index 3b751aec64..6c6349bdd1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java @@ -17,7 +17,6 @@ package org.apache.rocketmq.proxy.channel; -import com.google.common.base.Strings; import io.grpc.Context; import java.util.ArrayList; import java.util.Collections; @@ -28,9 +27,10 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.function.Supplier; +import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.constant.LoggerName; -import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.common.Cleaner; +import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.slf4j.Logger; @@ -54,14 +54,12 @@ public class ChannelManager { } public T createChannel(String clientId, Supplier creator, Class clazz) { - if (Strings.isNullOrEmpty(clientId)) { + if (StringUtils.isBlank(clientId)) { log.warn("ClientId is unexpected null or empty"); return creator.get(); } - if (!clientIdChannelMap.containsKey(clientId)) { - clientIdChannelMap.putIfAbsent(clientId, creator.get()); - } + clientIdChannelMap.computeIfAbsent(clientId, key -> creator.get()); T channel = clazz.cast(clientIdChannelMap.get(clientId)); channel.updateLastAccessTime(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java index 8b5265a049..46c8d3d6ae 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java @@ -678,10 +678,10 @@ public class GrpcConverter { return consumeMessageDirectlyResult; } - public static Resource buildResource(String resourceNameWithNamespace) { + public static Resource buildResource(String resourceStr) { return Resource.newBuilder() - .setResourceNamespace(NamespaceUtil.getNamespaceFromResource(resourceNameWithNamespace)) - .setName(NamespaceUtil.withoutNamespace(resourceNameWithNamespace)) + .setResourceNamespace(NamespaceUtil.getNamespaceFromResource(resourceStr)) + .setName(NamespaceUtil.withoutNamespace(resourceStr)) .build(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java index bdf2ef601f..1260b0b04b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java @@ -17,8 +17,8 @@ package org.apache.rocketmq.proxy.grpc.adapter; +import java.time.Duration; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; public class InvocationContext { private final R request; @@ -31,7 +31,7 @@ public class InvocationContext { } public boolean expired(long expiredTimeSec) { - return System.currentTimeMillis() - timestamp >= TimeUnit.SECONDS.toMillis(expiredTimeSec); + return System.currentTimeMillis() - timestamp >= Duration.ofSeconds(expiredTimeSec).toMillis(); } public R getRequest() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java index 44860e1919..88d123758e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java @@ -44,27 +44,18 @@ public class ResponseWriter { } } - public static void writeException(StreamObserver observer, final Throwable e) { + public static void writeException(StreamObserver observer, final Throwable e) { if (observer instanceof ServerCallStreamObserver) { - final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver) observer; if (null == e) { return; } + final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver) observer; if (serverCallStreamObserver.isCancelled()) { log.warn("Client has cancelled the request. Exception to write", e); return; } -// if (e instanceof CompletionException) { -// if (e.getCause() instanceof ProxyException) { -// ProxyException proxyException = (ProxyException) e.getCause(); -// serverCallStreamObserver.onNext(ResponseBuilder.buildCommon(proxyException.getCode(), proxyException.getMessage())); -// serverCallStreamObserver.onCompleted(); -// return; -// } -// } - log.debug("Start to write error response", e); serverCallStreamObserver.onError(e); serverCallStreamObserver.onCompleted(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java index e389b02713..bcdc46173f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java @@ -53,7 +53,7 @@ public class GrpcClientChannel extends SimpleChannel { this.manager = manager; } - public void addClientObserver(CompletableFuture future) { + public void setClientObserver(CompletableFuture future) { this.pollCommandResponseFutureRef.set(future); } @@ -131,8 +131,7 @@ public class GrpcClientChannel extends SimpleChannel { break; } case RequestCode.GET_CONSUMER_RUNNING_INFO: { - final GetConsumerRunningInfoRequestHeader requestHeader = - (GetConsumerRunningInfoRequestHeader) command.decodeCommandCustomHeader(GetConsumerRunningInfoRequestHeader.class); + final GetConsumerRunningInfoRequestHeader requestHeader = command.decodeCommandCustomHeader(GetConsumerRunningInfoRequestHeader.class); if (!requestHeader.isJstackEnable()) { break; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index 88b1853022..b562c15132 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -420,7 +420,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo .build()); break; } - producerChannel.addClientObserver(future); + producerChannel.setClientObserver(future); break; case CONSUMER_GROUP: Resource consumerGroup = request.getConsumerGroup(); @@ -432,7 +432,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo .build()); break; } - consumerChannel.addClientObserver(future); + consumerChannel.setClientObserver(future); break; default: break; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java index ca6f1e79b4..52f79f6e8c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java @@ -145,7 +145,7 @@ public class ForwardClientService extends BaseService { if (producerChannel == null) { future.complete(noopCommandResponse); } else { - producerChannel.addClientObserver(future); + producerChannel.setClientObserver(future); } break; case CONSUMER_GROUP: @@ -155,7 +155,7 @@ public class ForwardClientService extends BaseService { if (consumerChannel == null) { future.complete(noopCommandResponse); } else { - consumerChannel.addClientObserver(future); + consumerChannel.setClientObserver(future); } break; default: