mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
+6
-1
@@ -91,13 +91,18 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker());
|
||||
this.consumerService = new ConsumerService(connectorManager, grpcClientManager);
|
||||
this.producerService = new ProducerService(connectorManager);
|
||||
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager, grpcClientManager);
|
||||
this.routeService = new RouteService(connectorManager, grpcClientManager);
|
||||
this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService,
|
||||
channelManager, grpcClientManager, pollCommandResponseManager);
|
||||
this.transactionService = new TransactionService(connectorManager, channelManager);
|
||||
|
||||
this.appendStartAndShutdown(new ClusterGrpcServiceStartAndShutdown());
|
||||
this.appendStartAndShutdown(this.connectorManager);
|
||||
this.appendStartAndShutdown(this.consumerService);
|
||||
this.appendStartAndShutdown(this.producerService);
|
||||
this.appendStartAndShutdown(this.routeService);
|
||||
this.appendStartAndShutdown(this.clientService);
|
||||
this.appendStartAndShutdown(this.transactionService);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+20
-4
@@ -22,11 +22,13 @@ import apache.rocketmq.v2.Resource;
|
||||
import io.grpc.Context;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteCache;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException;
|
||||
|
||||
public class BaseService {
|
||||
public abstract class BaseService implements StartAndShutdown {
|
||||
|
||||
protected final ConnectorManager connectorManager;
|
||||
|
||||
@@ -34,7 +36,7 @@ public class BaseService {
|
||||
this.connectorManager = connectorManager;
|
||||
}
|
||||
|
||||
protected ReceiptHandle resolveReceiptHandle(Context ctx, String receiptHandleStr) {
|
||||
public static ReceiptHandle resolveReceiptHandle(Context ctx, String receiptHandleStr) {
|
||||
ReceiptHandle receiptHandle = ReceiptHandle.decode(receiptHandleStr);
|
||||
if (receiptHandle.isExpired()) {
|
||||
throw new ProxyException(Code.RECEIPT_HANDLE_EXPIRED, "handle has expired");
|
||||
@@ -42,20 +44,34 @@ public class BaseService {
|
||||
return receiptHandle;
|
||||
}
|
||||
|
||||
protected String getBrokerAddr(Context ctx, String brokerName) throws Exception {
|
||||
public static String getBrokerAddr(Context ctx, TopicRouteCache topicRouteCache, String brokerName) throws Exception {
|
||||
if (StringUtils.isBlank(brokerName)) {
|
||||
throw new ProxyException(Code.UNRECOGNIZED, "broker name is empty");
|
||||
}
|
||||
String addr = this.connectorManager.getTopicRouteCache().getBrokerAddr(brokerName);
|
||||
String addr = topicRouteCache.getBrokerAddr(brokerName);
|
||||
if (StringUtils.isBlank(addr)) {
|
||||
throw new ProxyException(Code.UNRECOGNIZED, brokerName + " not exist");
|
||||
}
|
||||
return addr;
|
||||
}
|
||||
|
||||
protected String getBrokerAddr(Context ctx, String brokerName) throws Exception {
|
||||
return getBrokerAddr(ctx, this.connectorManager.getTopicRouteCache(), brokerName);
|
||||
}
|
||||
|
||||
protected void checkSubscriptionData(Resource topic, FilterExpression filterExpression) {
|
||||
// for checking filterExpression.
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcConverter.buildSubscriptionData(topicName, filterExpression);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() throws Exception {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
+26
-137
@@ -22,20 +22,17 @@ import apache.rocketmq.v2.AckMessageResponse;
|
||||
import apache.rocketmq.v2.AckMessageResultEntry;
|
||||
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
|
||||
import apache.rocketmq.v2.ChangeInvisibleDurationResponse;
|
||||
import apache.rocketmq.v2.ClientType;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.Message;
|
||||
import apache.rocketmq.v2.NackMessageRequest;
|
||||
import apache.rocketmq.v2.NackMessageResponse;
|
||||
import apache.rocketmq.v2.ReceiveMessageRequest;
|
||||
import apache.rocketmq.v2.ReceiveMessageResponse;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.RetryPolicy;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import io.grpc.Context;
|
||||
import io.grpc.stub.StreamObserver;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.client.consumer.AckResult;
|
||||
@@ -43,14 +40,10 @@ import org.apache.rocketmq.client.consumer.AckStatus;
|
||||
import org.apache.rocketmq.client.consumer.PopResult;
|
||||
import org.apache.rocketmq.client.consumer.PopStatus;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.proxy.common.utils.FilterUtils;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardProducer;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
|
||||
@@ -74,9 +67,8 @@ public class ConsumerService extends BaseService {
|
||||
protected final GrpcClientManager grpcClientManager;
|
||||
|
||||
private volatile ReadQueueSelector readQueueSelector;
|
||||
private volatile ReceiveMessageResultFilter receiveMessageResultFilter;
|
||||
private volatile ResponseHook<ReceiveMessageRequest, List<ReceiveMessageResponse>> receiveMessageHook;
|
||||
private volatile ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook;
|
||||
private volatile ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> forwardToDLQInRecvMessageHook;
|
||||
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook;
|
||||
private volatile ResponseHook<ChangeInvisibleDurationRequest, ChangeInvisibleDurationResponse> changeInvisibleDurationHook;
|
||||
@@ -86,12 +78,16 @@ public class ConsumerService extends BaseService {
|
||||
this.readConsumer = connectorManager.getForwardReadConsumer();
|
||||
this.writeConsumer = connectorManager.getForwardWriteConsumer();
|
||||
this.producer = connectorManager.getForwardProducer();
|
||||
|
||||
this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache());
|
||||
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache());
|
||||
this.receiveMessageResultFilter = new DefaultReceiveMessageResultFilter(
|
||||
producer, writeConsumer, grpcClientManager, connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
public void receiveMessage(Context ctx, ReceiveMessageRequest request,
|
||||
StreamObserver<ReceiveMessageResponse> responseObserver) {
|
||||
this.receiveMessage(ctx, request)
|
||||
@@ -146,10 +142,10 @@ public class ConsumerService extends BaseService {
|
||||
PopStatus status = result.getPopStatus();
|
||||
switch (status) {
|
||||
case FOUND:
|
||||
List<Message> messageList = filterMessage(ctx, request, result.getMsgFoundList());
|
||||
List<Message> messageList = this.receiveMessageResultFilter.filterMessage(ctx, request, result.getMsgFoundList());
|
||||
if (messageList.isEmpty()) {
|
||||
responseList.add(ReceiveMessageResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, "no new message"))
|
||||
.build());
|
||||
} else {
|
||||
for (Message message : messageList) {
|
||||
@@ -176,97 +172,6 @@ public class ConsumerService extends BaseService {
|
||||
return responseList;
|
||||
}
|
||||
|
||||
protected List<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> messageExtList) {
|
||||
if (messageExtList == null || messageExtList.isEmpty()) {
|
||||
return Collections.emptyList();
|
||||
}
|
||||
Settings settings = grpcClientManager.getClientSettings(ctx);
|
||||
ClientType clientType = settings.getClientType();
|
||||
int maxAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts();
|
||||
Resource topic = request.getMessageQueue().getTopic();
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression());
|
||||
|
||||
List<Message> resMessageList = new ArrayList<>();
|
||||
for (MessageExt messageExt : messageExtList) {
|
||||
if (ClientType.SIMPLE_CONSUMER.equals(clientType) && messageExt.getReconsumeTimes() >= maxAttempts) {
|
||||
forwardMessageToDLQ(ctx, request, messageExt, maxAttempts);
|
||||
continue;
|
||||
}
|
||||
if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) {
|
||||
this.ackNoMatchedMessage(ctx, request, messageExt);
|
||||
continue;
|
||||
}
|
||||
resMessageList.add(GrpcConverter.buildMessage(messageExt));
|
||||
}
|
||||
return resMessageList;
|
||||
}
|
||||
|
||||
protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt,
|
||||
int maxReconsumeTimes) {
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
|
||||
|
||||
try {
|
||||
ReceiptHandle handle = ReceiptHandle.create(messageExt);
|
||||
if (handle == null) {
|
||||
return;
|
||||
}
|
||||
String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName());
|
||||
Resource topic = request.getMessageQueue().getTopic();
|
||||
Resource group = request.getGroup();
|
||||
ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader(
|
||||
topic,
|
||||
group,
|
||||
handle,
|
||||
messageExt.getMsgId(),
|
||||
maxReconsumeTimes);
|
||||
AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(
|
||||
topic,
|
||||
group,
|
||||
handle);
|
||||
|
||||
future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader);
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
|
||||
future.whenComplete((result, throwable) -> {
|
||||
if (forwardToDLQInRecvMessageHook != null) {
|
||||
forwardToDLQInRecvMessageHook.beforeResponse(ctx, consumerSendMsgBackRequestHeader, result, throwable);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) {
|
||||
CompletableFuture<AckResult> future = new CompletableFuture<>();
|
||||
|
||||
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
|
||||
try {
|
||||
ReceiptHandle handle = ReceiptHandle.create(messageExt);
|
||||
if (handle == null) {
|
||||
return;
|
||||
}
|
||||
String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName());
|
||||
ackMessageRequestHeader.setConsumerGroup(GrpcConverter.wrapResourceWithNamespace(request.getGroup()));
|
||||
ackMessageRequestHeader.setTopic(messageExt.getTopic());
|
||||
ackMessageRequestHeader.setQueueId(handle.getQueueId());
|
||||
ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle());
|
||||
ackMessageRequestHeader.setOffset(handle.getOffset());
|
||||
|
||||
future = this.writeConsumer.ackMessage(ctx, brokerAddr, ackMessageRequestHeader);
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
|
||||
future.whenComplete((ackResult, throwable) -> {
|
||||
if (ackNoMatchedMessageHook != null) {
|
||||
ackNoMatchedMessageHook.beforeResponse(ctx, ackMessageRequestHeader, ackResult, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
}
|
||||
|
||||
public CompletableFuture<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request) {
|
||||
CompletableFuture<AckMessageResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
@@ -309,7 +214,7 @@ public class ConsumerService extends BaseService {
|
||||
.setReceiptHandle(ackMessageEntry.getReceiptHandle());
|
||||
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, ackMessageEntry.getReceiptHandle());
|
||||
ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, ackMessageEntry.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
|
||||
AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request, receiptHandle);
|
||||
@@ -350,25 +255,18 @@ public class ConsumerService extends BaseService {
|
||||
public CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
|
||||
CompletableFuture<NackMessageResponse> future = new CompletableFuture<>();
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
|
||||
Settings settings = grpcClientManager.getClientSettings(ctx);
|
||||
int maxDeliveryAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts();
|
||||
if (request.getDeliveryAttempt() >= maxDeliveryAttempts) {
|
||||
future = this.producer.sendMessageBack(
|
||||
future = this.producer.sendMessageBackThenAckOrg(
|
||||
ctx,
|
||||
brokerAddr,
|
||||
this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts)
|
||||
).thenApply(result -> {
|
||||
if (result.getCode() == ResponseCode.SUCCESS) {
|
||||
writeConsumer.ackMessage(
|
||||
ctx,
|
||||
brokerAddr,
|
||||
this.buildAckMessageRequestHeader(ctx, request));
|
||||
}
|
||||
return convertToNackMessageResponse(ctx, request, result);
|
||||
});
|
||||
this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts),
|
||||
this.buildAckMessageRequestHeader(ctx, request)
|
||||
).thenApply(result -> convertToNackMessageResponse(ctx, request, result));
|
||||
} else {
|
||||
ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request);
|
||||
future = this.writeConsumer.changeInvisibleTimeAsync(ctx, brokerAddr, receiptHandle.getBrokerName(), requestHeader)
|
||||
@@ -425,7 +323,7 @@ public class ConsumerService extends BaseService {
|
||||
CompletableFuture<ChangeInvisibleDurationResponse> future = new CompletableFuture<>();
|
||||
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
|
||||
ChangeInvisibleTimeRequestHeader requestHeader = convertToChangeInvisibleTimeRequestHeader(ctx, request);
|
||||
@@ -468,6 +366,15 @@ public class ConsumerService extends BaseService {
|
||||
this.readQueueSelector = readQueueSelector;
|
||||
}
|
||||
|
||||
public ReceiveMessageResultFilter getReceiveMessageResultFilter() {
|
||||
return receiveMessageResultFilter;
|
||||
}
|
||||
|
||||
public void setReceiveMessageResultFilter(
|
||||
ReceiveMessageResultFilter receiveMessageResultFilter) {
|
||||
this.receiveMessageResultFilter = receiveMessageResultFilter;
|
||||
}
|
||||
|
||||
public ResponseHook<ReceiveMessageRequest, List<ReceiveMessageResponse>> getReceiveMessageHook() {
|
||||
return receiveMessageHook;
|
||||
}
|
||||
@@ -477,24 +384,6 @@ public class ConsumerService extends BaseService {
|
||||
this.receiveMessageHook = receiveMessageHook;
|
||||
}
|
||||
|
||||
public ResponseHook<AckMessageRequestHeader, AckResult> getAckNoMatchedMessageHook() {
|
||||
return ackNoMatchedMessageHook;
|
||||
}
|
||||
|
||||
public void setAckNoMatchedMessageHook(
|
||||
ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook) {
|
||||
this.ackNoMatchedMessageHook = ackNoMatchedMessageHook;
|
||||
}
|
||||
|
||||
public ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> getForwardToDLQInRecvMessageHook() {
|
||||
return forwardToDLQInRecvMessageHook;
|
||||
}
|
||||
|
||||
public void setForwardToDLQInRecvMessageHook(
|
||||
ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> forwardToDLQInRecvMessageHook) {
|
||||
this.forwardToDLQInRecvMessageHook = forwardToDLQInRecvMessageHook;
|
||||
}
|
||||
|
||||
public ResponseHook<AckMessageRequest, AckMessageResponse> getAckMessageHook() {
|
||||
return ackMessageHook;
|
||||
}
|
||||
|
||||
+173
@@ -0,0 +1,173 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
|
||||
import apache.rocketmq.v2.ClientType;
|
||||
import apache.rocketmq.v2.Message;
|
||||
import apache.rocketmq.v2.ReceiveMessageRequest;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.client.consumer.AckResult;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.proxy.common.utils.FilterUtils;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardProducer;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteCache;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
import static org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseService.getBrokerAddr;
|
||||
|
||||
public class DefaultReceiveMessageResultFilter implements ReceiveMessageResultFilter {
|
||||
|
||||
protected final ForwardProducer producer;
|
||||
protected final ForwardWriteConsumer writeConsumer;
|
||||
protected final GrpcClientManager grpcClientManager;
|
||||
protected final TopicRouteCache topicRouteCache;
|
||||
|
||||
private volatile ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook;
|
||||
private volatile ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> forwardToDLQInRecvMessageHook;
|
||||
|
||||
public DefaultReceiveMessageResultFilter(ForwardProducer producer, ForwardWriteConsumer writeConsumer,
|
||||
GrpcClientManager grpcClientManager, TopicRouteCache topicRouteCache) {
|
||||
this.producer = producer;
|
||||
this.writeConsumer = writeConsumer;
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
this.topicRouteCache = topicRouteCache;
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> messageExtList) {
|
||||
if (messageExtList == null || messageExtList.isEmpty()) {
|
||||
return Collections.emptyList();
|
||||
}
|
||||
Settings settings = grpcClientManager.getClientSettings(ctx);
|
||||
ClientType clientType = settings.getClientType();
|
||||
int maxAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts();
|
||||
Resource topic = request.getMessageQueue().getTopic();
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression());
|
||||
|
||||
List<Message> resMessageList = new ArrayList<>();
|
||||
for (MessageExt messageExt : messageExtList) {
|
||||
if (ClientType.SIMPLE_CONSUMER.equals(clientType) && messageExt.getReconsumeTimes() >= maxAttempts) {
|
||||
forwardMessageToDLQ(ctx, request, messageExt, maxAttempts);
|
||||
continue;
|
||||
}
|
||||
if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) {
|
||||
this.ackNoMatchedMessage(ctx, request, messageExt);
|
||||
continue;
|
||||
}
|
||||
resMessageList.add(GrpcConverter.buildMessage(messageExt));
|
||||
}
|
||||
return resMessageList;
|
||||
}
|
||||
|
||||
protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt,
|
||||
int maxReconsumeTimes) {
|
||||
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
|
||||
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
|
||||
|
||||
try {
|
||||
ReceiptHandle handle = ReceiptHandle.create(messageExt);
|
||||
if (handle == null) {
|
||||
return;
|
||||
}
|
||||
String brokerAddr = getBrokerAddr(ctx, topicRouteCache, handle.getBrokerName());
|
||||
Resource topic = request.getMessageQueue().getTopic();
|
||||
Resource group = request.getGroup();
|
||||
ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader(
|
||||
topic,
|
||||
group,
|
||||
handle,
|
||||
messageExt.getMsgId(),
|
||||
maxReconsumeTimes);
|
||||
AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(
|
||||
topic,
|
||||
group,
|
||||
handle);
|
||||
|
||||
future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader);
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
|
||||
future.whenComplete((result, throwable) -> {
|
||||
if (forwardToDLQInRecvMessageHook != null) {
|
||||
forwardToDLQInRecvMessageHook.beforeResponse(ctx, consumerSendMsgBackRequestHeader, result, throwable);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) {
|
||||
CompletableFuture<AckResult> future = new CompletableFuture<>();
|
||||
|
||||
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
|
||||
try {
|
||||
ReceiptHandle handle = ReceiptHandle.create(messageExt);
|
||||
if (handle == null) {
|
||||
return;
|
||||
}
|
||||
String brokerAddr = getBrokerAddr(ctx, topicRouteCache, handle.getBrokerName());
|
||||
ackMessageRequestHeader.setConsumerGroup(GrpcConverter.wrapResourceWithNamespace(request.getGroup()));
|
||||
ackMessageRequestHeader.setTopic(messageExt.getTopic());
|
||||
ackMessageRequestHeader.setQueueId(handle.getQueueId());
|
||||
ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle());
|
||||
ackMessageRequestHeader.setOffset(handle.getOffset());
|
||||
|
||||
future = this.writeConsumer.ackMessage(ctx, brokerAddr, ackMessageRequestHeader);
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
|
||||
future.whenComplete((ackResult, throwable) -> {
|
||||
if (ackNoMatchedMessageHook != null) {
|
||||
ackNoMatchedMessageHook.beforeResponse(ctx, ackMessageRequestHeader, ackResult, throwable);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public ResponseHook<AckMessageRequestHeader, AckResult> getAckNoMatchedMessageHook() {
|
||||
return ackNoMatchedMessageHook;
|
||||
}
|
||||
|
||||
public void setAckNoMatchedMessageHook(
|
||||
ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook) {
|
||||
this.ackNoMatchedMessageHook = ackNoMatchedMessageHook;
|
||||
}
|
||||
|
||||
public ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> getForwardToDLQInRecvMessageHook() {
|
||||
return forwardToDLQInRecvMessageHook;
|
||||
}
|
||||
|
||||
public void setForwardToDLQInRecvMessageHook(
|
||||
ResponseHook<ConsumerSendMsgBackRequestHeader, RemotingCommand> forwardToDLQInRecvMessageHook) {
|
||||
this.forwardToDLQInRecvMessageHook = forwardToDLQInRecvMessageHook;
|
||||
}
|
||||
}
|
||||
+3
-1
@@ -82,9 +82,11 @@ public class ForwardClientService extends BaseService {
|
||||
this.channelManager = channelManager;
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
this.telemetryCommandManager = telemetryCommandManager;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.clientSettingsService = new ClientSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager);
|
||||
|
||||
this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListenerImpl());
|
||||
this.producerManager = new ProducerManager();
|
||||
this.producerManager.appendProducerChangeListener(new ProducerChangeListenerImpl());
|
||||
|
||||
+6
-2
@@ -53,7 +53,11 @@ public class ProducerService extends BaseService {
|
||||
public ProducerService(ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
this.producer = connectorManager.getForwardProducer();
|
||||
writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
public CompletableFuture<SendMessageResponse> sendMessage(Context ctx, SendMessageRequest request) {
|
||||
@@ -123,7 +127,7 @@ public class ProducerService extends BaseService {
|
||||
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> future = new CompletableFuture<>();
|
||||
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = this.buildConsumerSendMsgBackRequestHeader(ctx, request);
|
||||
AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(
|
||||
|
||||
+29
@@ -0,0 +1,29 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
|
||||
import apache.rocketmq.v2.Message;
|
||||
import apache.rocketmq.v2.ReceiveMessageRequest;
|
||||
import io.grpc.Context;
|
||||
import java.util.List;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
|
||||
public interface ReceiveMessageResultFilter {
|
||||
|
||||
List<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> messageExtList);
|
||||
}
|
||||
+44
-124
@@ -16,8 +16,6 @@
|
||||
*/
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
|
||||
import apache.rocketmq.v2.Address;
|
||||
import apache.rocketmq.v2.AddressScheme;
|
||||
import apache.rocketmq.v2.Assignment;
|
||||
import apache.rocketmq.v2.Broker;
|
||||
import apache.rocketmq.v2.Code;
|
||||
@@ -29,32 +27,23 @@ import apache.rocketmq.v2.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import com.google.common.base.Preconditions;
|
||||
import com.google.common.net.HostAndPort;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
|
||||
public class RouteService extends BaseService {
|
||||
private final ProxyMode mode;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook;
|
||||
|
||||
@@ -64,14 +53,16 @@ public class RouteService extends BaseService {
|
||||
|
||||
protected final GrpcClientManager grpcClientManager;
|
||||
|
||||
public RouteService(ProxyMode mode, ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager);
|
||||
Preconditions.checkArgument(ProxyMode.isClusterMode(mode) || ProxyMode.isLocalMode(mode));
|
||||
this.mode = mode;
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.queryRouteEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
}
|
||||
|
||||
public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
@@ -87,44 +78,26 @@ public class RouteService extends BaseService {
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName);
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
List<QueueData> queueDataList = topicRouteData.getQueueDatas();
|
||||
List<BrokerData> brokerDataList = topicRouteData.getBrokerDatas();
|
||||
|
||||
List<MessageQueue> messageQueueList = new ArrayList<>();
|
||||
if (ProxyMode.isClusterMode(mode.name())) {
|
||||
Settings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " +
|
||||
clientSettings.getAccessPoint() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
for (QueueData queueData : queueDataList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(queueData.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
|
||||
messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
Settings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " +
|
||||
clientSettings.getAccessPoint() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
if (ProxyMode.isLocalMode(mode.name())) {
|
||||
Map<String, Map<Long, Broker>> brokerMap = buildBrokerMap(brokerDataList);
|
||||
for (QueueData queueData : queueDataList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(queueData.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
|
||||
for (QueueData queueData : queueDataList) {
|
||||
String brokerName = queueData.getBrokerName();
|
||||
Map<Long, Broker> brokerIdMap = brokerMap.get(brokerName);
|
||||
if (brokerIdMap == null) {
|
||||
break;
|
||||
}
|
||||
for (Broker broker : brokerIdMap.values()) {
|
||||
messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
}
|
||||
messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
|
||||
QueryRouteResponse response = QueryRouteResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
|
||||
.addAllMessageQueues(messageQueueList)
|
||||
@@ -153,57 +126,32 @@ public class RouteService extends BaseService {
|
||||
try {
|
||||
List<Assignment> assignments = new ArrayList<>();
|
||||
List<SelectableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request);
|
||||
if (ProxyMode.isLocalMode(mode)) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName);
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
Map<String, Map<Long, Broker>> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas());
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Map<Long, Broker> brokerIdMap = brokerMap.get(messageQueue.getBrokerName());
|
||||
if (brokerIdMap != null) {
|
||||
Broker broker = brokerIdMap.get(0L);
|
||||
|
||||
MessageQueue defaultMessageQueue = MessageQueue.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setMessageQueue(defaultMessageQueue)
|
||||
.build());
|
||||
}
|
||||
}
|
||||
Settings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " +
|
||||
clientSettings.getAccessPoint() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
if (ProxyMode.isClusterMode(mode)) {
|
||||
Settings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " +
|
||||
clientSettings.getAccessPoint() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(messageQueue.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(messageQueue.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
|
||||
MessageQueue defaultMessageQueue = MessageQueue.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
MessageQueue defaultMessageQueue = MessageQueue.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setMessageQueue(defaultMessageQueue)
|
||||
.build());
|
||||
}
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setMessageQueue(defaultMessageQueue)
|
||||
.build());
|
||||
}
|
||||
|
||||
QueryAssignmentResponse response = QueryAssignmentResponse.newBuilder()
|
||||
@@ -217,34 +165,6 @@ public class RouteService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
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<>();
|
||||
String brokerName = brokerData.getBrokerName();
|
||||
for (Map.Entry<Long, String> entry : brokerData.getBrokerAddrs().entrySet()) {
|
||||
Long brokerId = entry.getKey();
|
||||
HostAndPort hostAndPort = HostAndPort.fromString(entry.getValue());
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(brokerName)
|
||||
.setId(Math.toIntExact(brokerId))
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.setScheme(AddressScheme.IPv4)
|
||||
.addAddresses(
|
||||
Address.newBuilder()
|
||||
.setPort(ConfigurationManager.getProxyConfig().getGrpcServerPort())
|
||||
.setHost(hostAndPort.getHost())
|
||||
)
|
||||
.build())
|
||||
.build();
|
||||
|
||||
brokerIdMap.put(brokerId, broker);
|
||||
}
|
||||
brokerMap.put(brokerName, brokerIdMap);
|
||||
}
|
||||
return brokerMap;
|
||||
}
|
||||
|
||||
public ParameterConverter<Endpoints, Endpoints> getQueryRouteEndpointConverter() {
|
||||
return queryRouteEndpointConverter;
|
||||
}
|
||||
|
||||
+7
-2
@@ -51,10 +51,15 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
private ReadQueueSelector readQueueSelector;
|
||||
|
||||
private ConsumerService consumerService;
|
||||
private DefaultReceiveMessageResultFilter receiveMessageResultFilter;
|
||||
|
||||
@Override
|
||||
public void beforeEach() throws Throwable {
|
||||
consumerService = new ConsumerService(this.connectorManager, this.grpcClientManager);
|
||||
consumerService.start();
|
||||
|
||||
receiveMessageResultFilter = new DefaultReceiveMessageResultFilter(producerClient, writeConsumerClient, grpcClientManager, topicRouteCache);
|
||||
consumerService.setReceiveMessageResultFilter(receiveMessageResultFilter);
|
||||
consumerService.setReadQueueSelector(readQueueSelector);
|
||||
}
|
||||
|
||||
@@ -84,7 +89,7 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
|
||||
Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor());
|
||||
AtomicReference<String> ackHandler = new AtomicReference<>();
|
||||
consumerService.setAckNoMatchedMessageHook((ctx1, request, response, t) -> ackHandler.set(request.getExtraInfo()));
|
||||
receiveMessageResultFilter.setAckNoMatchedMessageHook((ctx1, request, response, t) -> ackHandler.set(request.getExtraInfo()));
|
||||
List<ReceiveMessageResponse> responseList = consumerService.receiveMessage(ctx,
|
||||
ReceiveMessageRequest.newBuilder()
|
||||
.setMessageQueue(apache.rocketmq.v2.MessageQueue.newBuilder()
|
||||
@@ -189,7 +194,7 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
doAnswer(mock -> {
|
||||
headerRef.set(mock.getArgument(2));
|
||||
return CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, ""));
|
||||
}).when(producerClient).sendMessageBack(any(), anyString(), any());
|
||||
}).when(producerClient).sendMessageBackThenAckOrg(any(), anyString(), any(), any());
|
||||
when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr");
|
||||
|
||||
Settings clientSettings = createClientSettings(3);
|
||||
|
||||
+1
@@ -47,6 +47,7 @@ public class ForwardClientServiceTest extends BaseServiceTest {
|
||||
this.channelManager,
|
||||
this.grpcClientManager,
|
||||
this.telemetryCommandManager);
|
||||
clientService.start();
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+4
-5
@@ -46,6 +46,8 @@ import static org.mockito.Mockito.when;
|
||||
|
||||
public class ProducerServiceTest extends BaseServiceTest {
|
||||
|
||||
private ProducerService producerService;
|
||||
|
||||
private static final SendMessageRequest REQUEST = SendMessageRequest.newBuilder()
|
||||
.addMessages(Message.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -61,6 +63,8 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
|
||||
@Override
|
||||
public void beforeEach() throws Throwable {
|
||||
producerService = new ProducerService(this.connectorManager);
|
||||
producerService.start();
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -71,7 +75,6 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
sendResultFuture.complete(new SendResult(SendStatus.SEND_OK, "msgId", new MessageQueue(),
|
||||
1L, "txId", "offsetMsgId", "regionId"));
|
||||
|
||||
ProducerService producerService = new ProducerService(this.connectorManager);
|
||||
producerService.setWriteQueueSelector((ctx, request) ->
|
||||
new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr"));
|
||||
|
||||
@@ -88,8 +91,6 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
|
||||
@Test
|
||||
public void testSendMessageNoQueueSelect() {
|
||||
ProducerService producerService = new ProducerService(this.connectorManager);
|
||||
|
||||
producerService.setWriteQueueSelector((ctx, request) -> null);
|
||||
|
||||
CompletableFuture<SendMessageResponse> future = producerService.sendMessage(Context.current(), SendMessageRequest.newBuilder()
|
||||
@@ -125,7 +126,6 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
.thenReturn(sendResultFuture);
|
||||
sendResultFuture.completeExceptionally(ex);
|
||||
|
||||
ProducerService producerService = new ProducerService(this.connectorManager);
|
||||
producerService.setWriteQueueSelector((ctx, request) ->
|
||||
new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr"));
|
||||
|
||||
@@ -145,7 +145,6 @@ public class ProducerServiceTest extends BaseServiceTest {
|
||||
public void testSendMessageWithErrorThrow() {
|
||||
RuntimeException ex = new RuntimeException();
|
||||
|
||||
ProducerService producerService = new ProducerService(this.connectorManager);
|
||||
producerService.setWriteQueueSelector((ctx, request) -> {
|
||||
throw ex;
|
||||
});
|
||||
|
||||
+5
-52
@@ -30,7 +30,6 @@ import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import com.google.common.net.HostAndPort;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
@@ -44,7 +43,6 @@ import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
@@ -77,6 +75,8 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
.setAccessPoint(Endpoints.getDefaultInstance())
|
||||
.build();
|
||||
|
||||
private RouteService routeService;
|
||||
|
||||
@Override
|
||||
public void beforeEach() throws Exception {
|
||||
TopicRouteData routeData = new TopicRouteData();
|
||||
@@ -106,6 +106,9 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
when(this.topicRouteCache.getMessageQueue("topic")).thenReturn(messageQueueWrapper);
|
||||
|
||||
when(this.topicRouteCache.getMessageQueue("notExistTopic")).thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, ""));
|
||||
|
||||
routeService = new RouteService(this.connectorManager, this.grpcClientManager);
|
||||
routeService.start();
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -161,28 +164,8 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
return queueData;
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryRoute() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber());
|
||||
assertEquals(8, response.getMessageQueuesCount());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getMessageQueues(0).getBroker()
|
||||
.getEndpoints().getAddresses(0).getHost());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWithInvalidEndpoints() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS);
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -196,8 +179,6 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
|
||||
@Test
|
||||
public void testQueryRoute() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
@@ -215,8 +196,6 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWhenTopicNotExist() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
@@ -231,8 +210,6 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
|
||||
@Test
|
||||
public void testQueryAssignmentInvalidEndpoints() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS);
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(
|
||||
@@ -246,32 +223,8 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
assertEquals(Code.ILLEGAL_ACCESS_POINT.getNumber(), response.getStatus().getCode().getNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryAssignment() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber());
|
||||
assertEquals(1, response.getAssignmentsCount());
|
||||
assertEquals("brokerName", response.getAssignments(0).getMessageQueue().getBroker().getName());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getAssignments(0).getMessageQueue().getBroker().getEndpoints().getAddresses(0).getHost());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryAssignment() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
|
||||
Reference in New Issue
Block a user