mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do some refactoring for readability.
This commit is contained in:
@@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.grpc.common;
|
||||
|
||||
import io.grpc.Context;
|
||||
|
||||
@FunctionalInterface
|
||||
public interface ParameterConverter<T, R> {
|
||||
R convert(Context ctx, T parameter) throws Throwable;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
+1
-1
@@ -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<SelectableMessageQueue> getAssignment(Context ctx, QueryAssignmentRequest request) throws Exception;
|
||||
}
|
||||
+8
-8
@@ -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<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook = null;
|
||||
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook = null;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> 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<ReceiveMessageResponse> 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<PopResult> 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(
|
||||
+2
-2
@@ -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;
|
||||
}
|
||||
|
||||
+2
-2
@@ -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;
|
||||
}
|
||||
|
||||
+3
-3
@@ -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;
|
||||
}
|
||||
|
||||
+16
-17
@@ -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<SendMessageRequest, SendMessageResponse> sendMessageHook = null;
|
||||
private volatile ResponseHook<ForwardMessageToDeadLetterQueueRequest, ForwardMessageToDeadLetterQueueResponse> forwardMessageToDLQHook = null;
|
||||
|
||||
public ProducerService(ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
messageQueueSelector = new DefaultProducerQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
public void setSendMessageHook(ResponseHook<SendMessageRequest, SendMessageResponse> 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<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> 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<RemotingCommand> 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);
|
||||
}
|
||||
|
||||
+23
-24
@@ -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<PullResult> 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);
|
||||
|
||||
+1
-1
@@ -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);
|
||||
}
|
||||
+13
-13
@@ -58,7 +58,7 @@ public class RouteService extends BaseService {
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook = null;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter;
|
||||
private volatile RouteAssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile AssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> 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<Endpoints, Endpoints> 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<String, Map<Long, Broker>> buildBrokerMap(List<BrokerData> brokerDataList) {
|
||||
private Map<String/*brokerName*/, Map<Long/*brokerID*/, Broker>> buildBrokerMap(List<BrokerData> brokerDataList) {
|
||||
Map<String, Map<Long, Broker>> brokerMap = new HashMap<>();
|
||||
for (BrokerData brokerData : brokerDataList) {
|
||||
Map<Long, Broker> brokerIdMap = new HashMap<>();
|
||||
|
||||
+1
-1
@@ -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);
|
||||
+4
-4
@@ -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()));
|
||||
|
||||
+4
-4
@@ -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<SendMessageResponse> 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<SendMessageResponse> 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<SendMessageResponse> 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) -> {
|
||||
|
||||
Reference in New Issue
Block a user