mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do refactor some code for readability.
This commit is contained in:
@@ -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 extends SimpleChannel> T createChannel(String clientId, Supplier<T> creator, Class<T> 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();
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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<R, W> {
|
||||
private final R request;
|
||||
@@ -31,7 +31,7 @@ public class InvocationContext<R, W> {
|
||||
}
|
||||
|
||||
public boolean expired(long expiredTimeSec) {
|
||||
return System.currentTimeMillis() - timestamp >= TimeUnit.SECONDS.toMillis(expiredTimeSec);
|
||||
return System.currentTimeMillis() - timestamp >= Duration.ofSeconds(expiredTimeSec).toMillis();
|
||||
}
|
||||
|
||||
public R getRequest() {
|
||||
|
||||
@@ -44,27 +44,18 @@ public class ResponseWriter {
|
||||
}
|
||||
}
|
||||
|
||||
public static <T> void writeException(StreamObserver<?> observer, final Throwable e) {
|
||||
public static <T> void writeException(StreamObserver<T> observer, final Throwable e) {
|
||||
if (observer instanceof ServerCallStreamObserver) {
|
||||
final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver<T>) observer;
|
||||
if (null == e) {
|
||||
return;
|
||||
}
|
||||
|
||||
final ServerCallStreamObserver<T> serverCallStreamObserver = (ServerCallStreamObserver<T>) 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();
|
||||
|
||||
+2
-3
@@ -53,7 +53,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
this.manager = manager;
|
||||
}
|
||||
|
||||
public void addClientObserver(CompletableFuture<PollCommandResponse> future) {
|
||||
public void setClientObserver(CompletableFuture<PollCommandResponse> 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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
+2
-2
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user