diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java index 6e9ce1643c..7561a21bc4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java @@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.grpc.common; import io.grpc.Context; +@FunctionalInterface public interface ParameterConverter { R convert(Context ctx, T parameter) throws Throwable; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 304e4598e2..4366b80c0b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -70,7 +70,7 @@ import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService; import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService; -import org.apache.rocketmq.proxy.grpc.service.cluster.ReceiveMessageService; +import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService; import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService; import org.slf4j.Logger; @@ -85,7 +85,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc private final ChannelManager channelManager; private final ConnectorManager connectorManager; private final ProducerService producerService; - private final ReceiveMessageService receiveMessageService; + private final ConsumerService receiveMessageService; private final RouteService routeService; private final ClientService clientService; private final PullMessageService pullMessageService; @@ -96,7 +96,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc this.channelManager = new ChannelManager(); this.pollCommandResponseManager = new PollCommandResponseManager(); this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker()); - this.receiveMessageService = new ReceiveMessageService(connectorManager); + this.receiveMessageService = new ConsumerService(connectorManager); this.producerService = new ProducerService(connectorManager); this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager); this.clientService = new ClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteAssignmentQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/AssignmentQueueSelector.java similarity index 95% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteAssignmentQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/AssignmentQueueSelector.java index ecdfd63539..fb2422736d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteAssignmentQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/AssignmentQueueSelector.java @@ -21,7 +21,7 @@ import io.grpc.Context; import java.util.List; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -public interface RouteAssignmentQueueSelector { +public interface AssignmentQueueSelector { List getAssignment(Context ctx, QueryAssignmentRequest request) throws Exception; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java similarity index 94% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index ffd91a14c3..505fb46e18 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -49,22 +49,22 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; -public class ReceiveMessageService extends BaseService { +public class ConsumerService extends BaseService { private final ForwardReadConsumer readConsumer; private final ForwardWriteConsumer writeConsumer; - private volatile ReceiveMessageQueueSelector receiveMessageQueueSelector; + private volatile ReadQueueSelector readQueueSelector; private volatile ResponseHook receiveMessageHook = null; private volatile ResponseHook ackMessageHook = null; private volatile ResponseHook nackMessageHook = null; - public ReceiveMessageService(ConnectorManager connectorManager) { + public ConsumerService(ConnectorManager connectorManager) { super(connectorManager); this.readConsumer = connectorManager.getForwardReadConsumer(); this.writeConsumer = connectorManager.getForwardWriteConsumer(); - this.receiveMessageQueueSelector = new DefaultReceiveMessageQueueSelector(connectorManager.getTopicRouteCache()); + this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache()); } public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { @@ -76,7 +76,7 @@ public class ReceiveMessageService extends BaseService { }); try { PopMessageRequestHeader requestHeader = this.convertToPopMessageRequestHeader(ctx, request); - SelectableMessageQueue messageQueue = this.receiveMessageQueueSelector.select(ctx, request, requestHeader); + SelectableMessageQueue messageQueue = this.readQueueSelector.select(ctx, request, requestHeader); CompletableFuture popResultFuture = this.readConsumer.popMessage( messageQueue.getBrokerAddr(), @@ -228,9 +228,9 @@ public class ReceiveMessageService extends BaseService { .build(); } - public void setReceiveMessageQueueSelector( - ReceiveMessageQueueSelector receiveMessageQueueSelector) { - this.receiveMessageQueueSelector = receiveMessageQueueSelector; + public void setReadQueueSelector( + ReadQueueSelector readQueueSelector) { + this.readQueueSelector = readQueueSelector; } public void setReceiveMessageHook( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultRouteAssignmentQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java similarity index 90% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultRouteAssignmentQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java index e112853b54..1738cb1f3e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultRouteAssignmentQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java @@ -24,11 +24,11 @@ import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; import org.apache.rocketmq.proxy.grpc.common.Converter; -public class DefaultRouteAssignmentQueueSelector implements RouteAssignmentQueueSelector { +public class DefaultAssignmentQueueSelector implements AssignmentQueueSelector { private final TopicRouteCache topicRouteCache; - public DefaultRouteAssignmentQueueSelector(TopicRouteCache topicRouteCache) { + public DefaultAssignmentQueueSelector(TopicRouteCache topicRouteCache) { this.topicRouteCache = topicRouteCache; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReceiveMessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReadQueueSelector.java similarity index 92% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReceiveMessageQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReadQueueSelector.java index df59b5e12b..284270921c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReceiveMessageQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultReadQueueSelector.java @@ -23,11 +23,11 @@ import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; -public class DefaultReceiveMessageQueueSelector implements ReceiveMessageQueueSelector { +public class DefaultReadQueueSelector implements ReadQueueSelector { private final TopicRouteCache topicRouteCache; - public DefaultReceiveMessageQueueSelector(TopicRouteCache topicRouteCache) { + public DefaultReadQueueSelector(TopicRouteCache topicRouteCache) { this.topicRouteCache = topicRouteCache; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java similarity index 94% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java index 5b7122f1c8..1cc9ba7bc6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java @@ -26,12 +26,12 @@ import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public class DefaultProducerQueueSelector implements ProducerQueueSelector { +public class DefaultWriteQueueSelector implements WriteQueueSelector { - private static final Logger log = LoggerFactory.getLogger(DefaultProducerQueueSelector.class); + private static final Logger log = LoggerFactory.getLogger(DefaultWriteQueueSelector.class); protected final TopicRouteCache topicRouteCache; - public DefaultProducerQueueSelector(TopicRouteCache topicRouteCache) { + public DefaultWriteQueueSelector(TopicRouteCache topicRouteCache) { this.topicRouteCache = topicRouteCache; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java index a640f1df6e..0435c8f794 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java @@ -23,6 +23,7 @@ import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import com.google.rpc.Code; import io.grpc.Context; +import java.util.concurrent.CompletableFuture; import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.rocketmq.client.producer.SendResult; @@ -39,25 +40,23 @@ import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -import java.util.concurrent.CompletableFuture; - public class ProducerService extends BaseService { - private volatile ProducerQueueSelector messageQueueSelector; + private volatile WriteQueueSelector writeQueueSelector; private volatile ResponseHook sendMessageHook = null; private volatile ResponseHook forwardMessageToDLQHook = null; public ProducerService(ConnectorManager connectorManager) { super(connectorManager); - messageQueueSelector = new DefaultProducerQueueSelector(this.connectorManager.getTopicRouteCache()); + writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache()); } public void setSendMessageHook(ResponseHook sendMessageHook) { this.sendMessageHook = sendMessageHook; } - public void setMessageQueueSelector(ProducerQueueSelector messageQueueSelector) { - this.messageQueueSelector = messageQueueSelector; + public void setWriteQueueSelector(WriteQueueSelector writeQueueSelector) { + this.writeQueueSelector = writeQueueSelector; } public void setForwardMessageToDLQHook( @@ -77,7 +76,7 @@ public class ProducerService extends BaseService { Pair requestPair = this.convertSendMessageRequest(ctx, request); SendMessageRequestHeader requestHeader = requestPair.getLeft(); org.apache.rocketmq.common.message.Message message = requestPair.getRight(); - SelectableMessageQueue addressableMessageQueue = messageQueueSelector.selectQueue(ctx, request, requestHeader, message); + SelectableMessageQueue addressableMessageQueue = writeQueueSelector.selectQueue(ctx, request, requestHeader, message); String topic = requestHeader.getTopic(); if (addressableMessageQueue == null) { @@ -151,17 +150,17 @@ public class ProducerService extends BaseService { CompletableFuture resultFuture = this.connectorManager.getForwardProducer() .sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); resultFuture - .thenAccept(result -> - future.complete( - ForwardMessageToDeadLetterQueueResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(result.getCode(), result.getRemark())) - .build() - ) + .thenAccept(result -> + future.complete( + ForwardMessageToDeadLetterQueueResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(result.getCode(), result.getRemark())) + .build() ) - .exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + ) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java index 559ae458a0..9c78a2ebb2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java @@ -26,6 +26,10 @@ import apache.rocketmq.v1.QueryOffsetResponse; import com.google.protobuf.util.Timestamps; import com.google.rpc.Code; import io.grpc.Context; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.consumer.PullStatus; import org.apache.rocketmq.common.message.MessageExt; @@ -41,11 +45,6 @@ import org.apache.rocketmq.proxy.grpc.common.ProxyException; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; - public class PullMessageService extends BaseService { private final DefaultForwardClient defaultForwardClient; @@ -84,14 +83,14 @@ public class PullMessageService extends BaseService { offsetFuture = this.defaultForwardClient.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); } offsetFuture.thenAccept(result -> future.complete( - QueryOffsetResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .setOffset(result) - .build())) - .exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + QueryOffsetResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .setOffset(result) + .build())) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } @@ -113,19 +112,19 @@ public class PullMessageService extends BaseService { String brokerAddr = this.getBrokerAddr(ctx, brokerName); CompletableFuture pullResultFuture = this.connectorManager.getForwardReadConsumer() - .pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + .pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); pullResultFuture - .thenAccept(pullResult -> { - try { - future.complete(convertToPullMessageResponse(ctx, request, pullResult)); - } catch (Throwable throwable) { - future.completeExceptionally(throwable); - } - }) - .exceptionally(throwable -> { + .thenAccept(pullResult -> { + try { + future.complete(convertToPullMessageResponse(ctx, request, pullResult)); + } catch (Throwable throwable) { future.completeExceptionally(throwable); - return null; - }); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReadQueueSelector.java similarity index 96% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReadQueueSelector.java index c671709cff..cd0702d1c1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReadQueueSelector.java @@ -21,7 +21,7 @@ import io.grpc.Context; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -public interface ReceiveMessageQueueSelector { +public interface ReadQueueSelector { SelectableMessageQueue select(Context ctx, ReceiveMessageRequest request, PopMessageRequestHeader requestHeader); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index c551135e9d..e7c64036b5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -58,7 +58,7 @@ public class RouteService extends BaseService { private volatile ResponseHook queryRouteHook = null; private volatile ParameterConverter queryAssignmentEndpointConverter; - private volatile RouteAssignmentQueueSelector assignmentQueueSelector; + private volatile AssignmentQueueSelector assignmentQueueSelector; private volatile ResponseHook queryAssignmentHook = null; public RouteService(ProxyMode mode, ConnectorManager connectorManager) { @@ -67,7 +67,7 @@ public class RouteService extends BaseService { this.mode = mode; queryRouteEndpointConverter = (ctx, parameter) -> parameter; queryAssignmentEndpointConverter = (ctx, parameter) -> parameter; - assignmentQueueSelector = new DefaultRouteAssignmentQueueSelector(this.connectorManager.getTopicRouteCache()); + assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache()); } public void setQueryRouteEndpointConverter(ParameterConverter queryRouteEndpointConverter) { @@ -83,7 +83,7 @@ public class RouteService extends BaseService { this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter; } - public void setAssignmentQueueSelector(RouteAssignmentQueueSelector assignmentQueueSelector) { + public void setAssignmentQueueSelector(AssignmentQueueSelector assignmentQueueSelector) { this.assignmentQueueSelector = assignmentQueueSelector; } @@ -179,25 +179,25 @@ public class RouteService extends BaseService { int queueIdIndex = 0; for (int i = 0; i < r; i++) { Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) - .setId(queueIdIndex++) - .setPermission(Permission.READ) - .build(); + .setId(queueIdIndex++) + .setPermission(Permission.READ) + .build(); partitionList.add(partition); } for (int i = 0; i < w; i++) { Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) - .setId(queueIdIndex++) - .setPermission(Permission.WRITE) - .build(); + .setId(queueIdIndex++) + .setPermission(Permission.WRITE) + .build(); partitionList.add(partition); } for (int i = 0; i < rw; i++) { Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) - .setId(queueIdIndex++) - .setPermission(Permission.READ_WRITE) - .build(); + .setId(queueIdIndex++) + .setPermission(Permission.READ_WRITE) + .build(); partitionList.add(partition); } @@ -278,7 +278,7 @@ public class RouteService extends BaseService { return future; } - private Map> buildBrokerMap(List brokerDataList) { + private Map> buildBrokerMap(List brokerDataList) { Map> brokerMap = new HashMap<>(); for (BrokerData brokerData : brokerDataList) { Map brokerIdMap = new HashMap<>(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/WriteQueueSelector.java similarity index 96% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerQueueSelector.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/WriteQueueSelector.java index ae4619db88..8605d11062 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/WriteQueueSelector.java @@ -21,7 +21,7 @@ import io.grpc.Context; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -public interface ProducerQueueSelector { +public interface WriteQueueSelector { SelectableMessageQueue selectQueue(Context ctx, SendMessageRequest request, SendMessageRequestHeader requestHeader, org.apache.rocketmq.common.message.Message message); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java index 35495513fe..21accfcf1e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java @@ -59,7 +59,7 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .setBody(ByteString.copyFrom("hello", StandardCharsets.UTF_8)) .build()) .build(); - ProducerQueueSelector queueSelector = new DefaultProducerQueueSelector(this.topicRouteCache); + WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, Converter.buildSendMessageRequestHeader(request), Converter.buildMessage(request.getMessage())); @@ -83,7 +83,7 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .setBody(ByteString.copyFrom("hello", StandardCharsets.UTF_8)) .build()) .build(); - ProducerQueueSelector queueSelector = new DefaultProducerQueueSelector(this.topicRouteCache); + WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, Converter.buildSendMessageRequestHeader(request), Converter.buildMessage(request.getMessage())); @@ -106,7 +106,7 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .setBody(ByteString.copyFrom("hello", StandardCharsets.UTF_8)) .build()) .build(); - ProducerQueueSelector queueSelector = new DefaultProducerQueueSelector(this.topicRouteCache); + WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, Converter.buildSendMessageRequestHeader(request), Converter.buildMessage(request.getMessage())); @@ -134,7 +134,7 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .build()) .build()) .build(); - ProducerQueueSelector queueSelector = new DefaultProducerQueueSelector(this.topicRouteCache); + WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, Converter.buildSendMessageRequestHeader(request), Converter.buildMessage(request.getMessage())); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java index 3d00c3dbbd..708d683d76 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java @@ -72,7 +72,7 @@ public class ProducerServiceTest extends BaseServiceTest { 1L, "txId", "offsetMsgId", "regionId")); ProducerService producerService = new ProducerService(this.clientManager); - producerService.setMessageQueueSelector((ctx, request, requestHeader, message) -> + producerService.setWriteQueueSelector((ctx, request, requestHeader, message) -> new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr")); CompletableFuture future = producerService.sendMessage(Context.current(), REQUEST); @@ -90,7 +90,7 @@ public class ProducerServiceTest extends BaseServiceTest { public void testSendMessageNoQueueSelect() { ProducerService producerService = new ProducerService(this.clientManager); - producerService.setMessageQueueSelector((ctx, request, requestHeader, message) -> null); + producerService.setWriteQueueSelector((ctx, request, requestHeader, message) -> null); CompletableFuture future = producerService.sendMessage(Context.current(), SendMessageRequest.newBuilder() .setMessage(Message.newBuilder() @@ -126,7 +126,7 @@ public class ProducerServiceTest extends BaseServiceTest { sendResultFuture.completeExceptionally(ex); ProducerService producerService = new ProducerService(this.clientManager); - producerService.setMessageQueueSelector((ctx, request, requestHeader, message) -> + producerService.setWriteQueueSelector((ctx, request, requestHeader, message) -> new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr")); CompletableFuture future = producerService.sendMessage(Context.current(), REQUEST); @@ -146,7 +146,7 @@ public class ProducerServiceTest extends BaseServiceTest { RuntimeException ex = new RuntimeException(); ProducerService producerService = new ProducerService(this.clientManager); - producerService.setMessageQueueSelector((ctx, request, requestHeader, message) -> { + producerService.setWriteQueueSelector((ctx, request, requestHeader, message) -> { throw ex; }); producerService.setSendMessageHook((request, response, t) -> {