diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index e8c206b609..b238335f2b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -26,7 +26,6 @@ import apache.rocketmq.v2.QueryAssignmentRequest; import apache.rocketmq.v2.QueryAssignmentResponse; import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.QueryRouteResponse; -import apache.rocketmq.v2.Settings; import io.grpc.Context; import java.util.ArrayList; import java.util.List; @@ -63,12 +62,12 @@ public class RouteService extends AbstractRouteService { List queueDataList = topicRouteData.getQueueDatas(); List messageQueueList = new ArrayList<>(); - Settings clientSettings = grpcClientManager.getClientSettings(ctx); - Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); + Endpoints endpoints = request.getEndpoints(); + Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, endpoints); if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) { future.complete(QueryRouteResponse.newBuilder() .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + - clientSettings.getAccessPoint() + " is invalidate")) + endpoints + " is invalidate")) .build()); return future; } @@ -110,12 +109,12 @@ public class RouteService extends AbstractRouteService { try { List assignments = new ArrayList<>(); List messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request); - Settings clientSettings = grpcClientManager.getClientSettings(ctx); - Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); + Endpoints endpoints = request.getEndpoints(); + Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, endpoints); if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) { future.complete(QueryAssignmentResponse.newBuilder() .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + - clientSettings.getAccessPoint() + " is invalidate")) + endpoints + " is invalidate")) .build()); return future; }