mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do some refactor work.
This commit is contained in:
@@ -58,15 +58,17 @@ public class ProxyStartup {
|
||||
// init thread pool monitor for proxy.
|
||||
initThreadPoolMonitor();
|
||||
|
||||
// create and start grpcServer
|
||||
// create grpcServer
|
||||
GrpcServer grpcServer = createGrpcServer();
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(grpcServer);
|
||||
|
||||
// health check server
|
||||
// create health check server
|
||||
final HealthCheckServer healthCheckServer = new HealthCheckServer();
|
||||
PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(healthCheckServer);
|
||||
|
||||
// start servers one by one.
|
||||
PROXY_START_AND_SHUTDOWN.start();
|
||||
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
||||
LOGGER.info("try to shutdown server");
|
||||
try {
|
||||
|
||||
@@ -22,6 +22,7 @@ import org.apache.rocketmq.client.exception.MQClientException;
|
||||
import org.apache.rocketmq.client.impl.MQClientAPIExt;
|
||||
import org.apache.rocketmq.common.protocol.header.GetConsumerListByGroupRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingException;
|
||||
@@ -59,6 +60,10 @@ public class DefaultForwardClient extends AbstractForwardClient {
|
||||
return this.getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> getMaxOffset(String brokerAddr, String topic, int queueId) {
|
||||
return this.getMaxOffset(brokerAddr, topic, queueId, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> getMaxOffset(
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
@@ -68,6 +73,15 @@ public class DefaultForwardClient extends AbstractForwardClient {
|
||||
return this.getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> searchOffset(
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
int queueId,
|
||||
long timestamp
|
||||
) {
|
||||
return this.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> searchOffset(
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
|
||||
@@ -26,6 +26,7 @@ import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
@@ -52,18 +53,21 @@ public class ForwardProducer extends AbstractForwardClient {
|
||||
return clientFactory.getTransactionalProducer(name, threadCount);
|
||||
}
|
||||
|
||||
|
||||
public CompletableFuture<Integer> heartBeat(String heartbeatAddr, HeartbeatData heartbeatData, long timeout) throws Exception {
|
||||
return this.getClient().sendHeartbeat(heartbeatAddr, heartbeatData, timeout);
|
||||
}
|
||||
|
||||
public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader, long timeoutMillis) throws Exception {
|
||||
this.getClient().endTransactionOneway(
|
||||
brokerAddr,
|
||||
requestHeader,
|
||||
"end transaction from rmq proxy",
|
||||
timeoutMillis
|
||||
);
|
||||
this.getClient().endTransactionOneway(brokerAddr, requestHeader, "end transaction from rmq proxy", timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<SendResult> sendMessage(
|
||||
String address,
|
||||
String brokerName,
|
||||
Message msg,
|
||||
SendMessageRequestHeader requestHeader
|
||||
) {
|
||||
return this.sendMessage(address, brokerName, msg, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<SendResult> sendMessage(
|
||||
@@ -84,6 +88,10 @@ public class ForwardProducer extends AbstractForwardClient {
|
||||
});
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) {
|
||||
return this.sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
|
||||
return this.getClient().sendMessageBack(brokerAddr, requestHeader, timeoutMillis);
|
||||
}
|
||||
|
||||
@@ -22,8 +22,9 @@ import org.apache.rocketmq.client.consumer.PullResult;
|
||||
import org.apache.rocketmq.client.impl.MQClientAPIExt;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
|
||||
public class ForwardReadConsumer extends AbstractForwardClient {
|
||||
|
||||
@@ -46,13 +47,21 @@ public class ForwardReadConsumer extends AbstractForwardClient {
|
||||
return clientFactory.getMQClient(name, threadCount);
|
||||
}
|
||||
|
||||
public CompletableFuture<PopResult> popMessage(String address, String brokerName, PopMessageRequestHeader requestHeader,
|
||||
long timeoutMillis) {
|
||||
return getClient().popMessage(address, brokerName, requestHeader, timeoutMillis);
|
||||
public CompletableFuture<PopResult> popMessage(
|
||||
String address,
|
||||
String brokerName,
|
||||
PopMessageRequestHeader requestHeader,
|
||||
long timeoutMillis
|
||||
) {
|
||||
return this.getClient().popMessage(address, brokerName, requestHeader, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<PullResult> pullMessage(String address, PullMessageRequestHeader requestHeader) {
|
||||
return this.pullMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<PullResult> pullMessage(String address, PullMessageRequestHeader requestHeader,
|
||||
long timeoutMillis) {
|
||||
return getClient().pullMessage(address, requestHeader, timeoutMillis);
|
||||
return this.getClient().pullMessage(address, requestHeader, timeoutMillis);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,8 +22,9 @@ import org.apache.rocketmq.client.impl.MQClientAPIExt;
|
||||
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.UpdateConsumerOffsetRequestHeader;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingException;
|
||||
|
||||
public class ForwardWriteConsumer extends AbstractForwardClient {
|
||||
@@ -47,12 +48,24 @@ public class ForwardWriteConsumer extends AbstractForwardClient {
|
||||
return clientFactory.getMQClient(name, threadCount);
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> ackMessage(String address, AckMessageRequestHeader requestHeader) {
|
||||
return this.ackMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> ackMessage(
|
||||
String address,
|
||||
AckMessageRequestHeader requestHeader,
|
||||
long timeoutMillis
|
||||
) {
|
||||
return getClient().ackMessage(address, requestHeader, timeoutMillis);
|
||||
return this.getClient().ackMessage(address, requestHeader, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> changeInvisibleTimeAsync(
|
||||
String address,
|
||||
String brokerName,
|
||||
ChangeInvisibleTimeRequestHeader requestHeader
|
||||
) {
|
||||
return this.changeInvisibleTimeAsync(address, brokerName, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
|
||||
public CompletableFuture<AckResult> changeInvisibleTimeAsync(
|
||||
@@ -61,7 +74,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient {
|
||||
ChangeInvisibleTimeRequestHeader requestHeader,
|
||||
long timeoutMillis
|
||||
) {
|
||||
return getClient().changeInvisibleTimeAsync(address, brokerName, requestHeader, timeoutMillis);
|
||||
return this.getClient().changeInvisibleTimeAsync(address, brokerName, requestHeader, timeoutMillis);
|
||||
}
|
||||
|
||||
public void updateConsumerOffsetOneWay(
|
||||
@@ -69,6 +82,6 @@ public class ForwardWriteConsumer extends AbstractForwardClient {
|
||||
UpdateConsumerOffsetRequestHeader header,
|
||||
long timeoutMillis
|
||||
) throws RemotingException, InterruptedException {
|
||||
getClient().updateConsumerOffsetOneWay(brokerAddr, header, timeoutMillis);
|
||||
this.getClient().updateConsumerOffsetOneWay(brokerAddr, header, timeoutMillis);
|
||||
}
|
||||
}
|
||||
|
||||
+4
-4
@@ -153,10 +153,10 @@ public class MessageQueueSelector {
|
||||
}
|
||||
|
||||
public final SelectableMessageQueue selectOne(String brokerName, int queueId) {
|
||||
for (SelectableMessageQueue addressableMessageQueue : queues) {
|
||||
String queueBrokerName = addressableMessageQueue.getBrokerName();
|
||||
if (queueBrokerName.equals(brokerName) && addressableMessageQueue.getQueueId() == queueId) {
|
||||
return addressableMessageQueue;
|
||||
for (SelectableMessageQueue targetMessageQueue : queues) {
|
||||
String queueBrokerName = targetMessageQueue.getBrokerName();
|
||||
if (queueBrokerName.equals(brokerName) && targetMessageQueue.getQueueId() == queueId) {
|
||||
return targetMessageQueue;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
|
||||
+1
-1
@@ -72,7 +72,7 @@ public class SelectableMessageQueue implements Comparable<SelectableMessageQueue
|
||||
}
|
||||
|
||||
@Override public String toString() {
|
||||
return "AddressableMessageQueue{" +
|
||||
return "SelectableMessageQueue{" +
|
||||
"messageQueue=" + messageQueue +
|
||||
", brokerAddr='" + brokerAddr + '\'' +
|
||||
'}';
|
||||
|
||||
@@ -51,6 +51,7 @@ import com.google.protobuf.Timestamp;
|
||||
import com.google.protobuf.util.Durations;
|
||||
import com.google.protobuf.util.Timestamps;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.net.SocketAddress;
|
||||
import java.net.UnknownHostException;
|
||||
import java.util.Arrays;
|
||||
@@ -700,4 +701,16 @@ public class GrpcConverter {
|
||||
}
|
||||
return header;
|
||||
}
|
||||
|
||||
public static long buildPollTimeFromContext(Context ctx) {
|
||||
long timeRemaining = ctx.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis();
|
||||
if (pollTime <= 0) {
|
||||
pollTime = timeRemaining;
|
||||
}
|
||||
|
||||
return pollTime;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -88,9 +88,14 @@ import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.PollResponseFuture;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.PullMessageChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.ReceiveMessageChannel;
|
||||
@@ -98,12 +103,6 @@ import org.apache.rocketmq.proxy.grpc.adapter.channel.SendMessageChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.handler.PullMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.PollResponseFuture;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.remoting.RemotingServer;
|
||||
@@ -226,13 +225,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
long timeRemaining = Context.current()
|
||||
.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis();
|
||||
if (pollTime <= 0) {
|
||||
pollTime = timeRemaining;
|
||||
}
|
||||
long pollTime = GrpcConverter.buildPollTimeFromContext(ctx);
|
||||
PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader);
|
||||
command.makeCustomHeaderToNet();
|
||||
@@ -392,13 +385,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
|
||||
@Override
|
||||
public CompletableFuture<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request) {
|
||||
long timeRemaining = Context.current()
|
||||
.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis();
|
||||
if (pollTime <= 0) {
|
||||
pollTime = timeRemaining;
|
||||
}
|
||||
long pollTime = GrpcConverter.buildPollTimeFromContext(ctx);
|
||||
PullMessageRequestHeader requestHeader = GrpcConverter.buildPullMessageRequestHeader(request, pollTime);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, requestHeader);
|
||||
command.makeCustomHeaderToNet();
|
||||
|
||||
@@ -16,11 +16,14 @@
|
||||
*/
|
||||
package org.apache.rocketmq.proxy.grpc.service.cluster;
|
||||
|
||||
import apache.rocketmq.v1.FilterExpression;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ProxyException;
|
||||
|
||||
public class BaseService {
|
||||
@@ -49,4 +52,10 @@ public class BaseService {
|
||||
}
|
||||
return addr;
|
||||
}
|
||||
|
||||
protected void checkSubscriptionData(Resource topic, FilterExpression filterExpression) {
|
||||
// for checking filterExpression.
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcConverter.buildSubscriptionData(topicName, filterExpression);
|
||||
}
|
||||
}
|
||||
|
||||
+58
-64
@@ -25,6 +25,9 @@ import apache.rocketmq.v1.ReceiveMessageRequest;
|
||||
import apache.rocketmq.v1.ReceiveMessageResponse;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.client.consumer.AckResult;
|
||||
import org.apache.rocketmq.client.consumer.AckStatus;
|
||||
import org.apache.rocketmq.client.consumer.PopResult;
|
||||
@@ -37,38 +40,33 @@ import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHead
|
||||
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.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardProducer;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ProxyException;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class ConsumerService extends BaseService {
|
||||
|
||||
private final DelayPolicy delayPolicy;
|
||||
private final ForwardReadConsumer readConsumer;
|
||||
private final ForwardWriteConsumer writeConsumer;
|
||||
/**
|
||||
* For sending messages back to broker.
|
||||
*/
|
||||
private final ForwardProducer producer;
|
||||
|
||||
private volatile ReadQueueSelector readQueueSelector;
|
||||
private volatile ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook = null;
|
||||
private volatile ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook = null;
|
||||
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook = null;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook = null;
|
||||
|
||||
private final DelayPolicy delayPolicy;
|
||||
private volatile ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook;
|
||||
private volatile ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook;
|
||||
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook;
|
||||
|
||||
public ConsumerService(ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
@@ -82,13 +80,15 @@ public class ConsumerService extends BaseService {
|
||||
|
||||
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
CompletableFuture<ReceiveMessageResponse> future = new CompletableFuture<>();
|
||||
// register hook.
|
||||
future.whenComplete((response, throwable) -> {
|
||||
if (receiveMessageHook != null) {
|
||||
receiveMessageHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
PopMessageRequestHeader requestHeader = this.convertToPopMessageRequestHeader(ctx, request);
|
||||
PopMessageRequestHeader requestHeader = this.buildPopMessageRequestHeader(ctx, request);
|
||||
SelectableMessageQueue messageQueue = this.readQueueSelector.select(ctx, request, requestHeader);
|
||||
|
||||
if (messageQueue == null) {
|
||||
@@ -118,27 +118,19 @@ public class ConsumerService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
protected PopMessageRequestHeader convertToPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) {
|
||||
// check filterExpression is correct or not
|
||||
GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression());
|
||||
|
||||
long timeRemaining = ctx.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis();
|
||||
if (pollTime <= 0) {
|
||||
pollTime = timeRemaining;
|
||||
}
|
||||
|
||||
return GrpcConverter.buildPopMessageRequestHeader(request, pollTime);
|
||||
protected PopMessageRequestHeader buildPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) {
|
||||
checkSubscriptionData(request.getPartition().getTopic(), request.getFilterExpression());
|
||||
return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx));
|
||||
}
|
||||
|
||||
protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) {
|
||||
SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(
|
||||
GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression());
|
||||
PopStatus status = result.getPopStatus();
|
||||
switch (status) {
|
||||
case FOUND:
|
||||
break;
|
||||
return ReceiveMessageResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.addAllMessages(checkAndGetMessagesFromPopResult(ctx, request, result))
|
||||
.build();
|
||||
case POLLING_FULL:
|
||||
return ReceiveMessageResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.RESOURCE_EXHAUSTED, "polling full"))
|
||||
@@ -150,6 +142,11 @@ public class ConsumerService extends BaseService {
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, "no new message"))
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
||||
protected List<Message> checkAndGetMessagesFromPopResult(Context ctx, ReceiveMessageRequest request, PopResult result) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic());
|
||||
SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression());
|
||||
|
||||
List<Message> messages = new ArrayList<>();
|
||||
for (MessageExt messageExt : result.getMsgFoundList()) {
|
||||
@@ -160,14 +157,12 @@ public class ConsumerService extends BaseService {
|
||||
messages.add(GrpcConverter.buildMessage(messageExt));
|
||||
}
|
||||
|
||||
return ReceiveMessageResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.addAllMessages(messages)
|
||||
.build();
|
||||
return messages;
|
||||
}
|
||||
|
||||
protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) {
|
||||
CompletableFuture<AckResult> future = new CompletableFuture<>();
|
||||
|
||||
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
|
||||
try {
|
||||
ReceiptHandle handle = ReceiptHandle.create(messageExt);
|
||||
@@ -181,15 +176,17 @@ public class ConsumerService extends BaseService {
|
||||
ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle());
|
||||
ackMessageRequestHeader.setOffset(handle.getOffset());
|
||||
|
||||
future = this.writeConsumer.ackMessage(brokerAddr, ackMessageRequestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
future = this.writeConsumer.ackMessage(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) {
|
||||
@@ -199,31 +196,32 @@ public class ConsumerService extends BaseService {
|
||||
ackMessageHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
|
||||
AckMessageRequestHeader requestHeader = this.convertToAckMessageRequestHeader(ctx, request);
|
||||
CompletableFuture<AckResult> ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request);
|
||||
CompletableFuture<AckResult> ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader);
|
||||
ackResultFuture
|
||||
.thenAccept(result -> {
|
||||
try {
|
||||
future.complete(convertToAckMessageResponse(ctx, request, result));
|
||||
} catch (Throwable throwable) {
|
||||
future.completeExceptionally(throwable);
|
||||
}
|
||||
})
|
||||
.exceptionally(throwable -> {
|
||||
.thenAccept(result -> {
|
||||
try {
|
||||
future.complete(convertToAckMessageResponse(ctx, request, result));
|
||||
} catch (Throwable throwable) {
|
||||
future.completeExceptionally(throwable);
|
||||
return null;
|
||||
});
|
||||
}
|
||||
})
|
||||
.exceptionally(throwable -> {
|
||||
future.completeExceptionally(throwable);
|
||||
return null;
|
||||
});
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
protected AckMessageRequestHeader convertToAckMessageRequestHeader(Context ctx, AckMessageRequest request) {
|
||||
protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request) {
|
||||
return GrpcConverter.buildAckMessageRequestHeader(request);
|
||||
}
|
||||
|
||||
@@ -252,8 +250,9 @@ public class ConsumerService extends BaseService {
|
||||
if (request.getDeliveryAttempt() >= request.getMaxDeliveryAttempts()) {
|
||||
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBack(
|
||||
brokerAddr,
|
||||
this.convertToConsumerSendMsgBackToDLQRequestHeader(ctx, request),
|
||||
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request)
|
||||
);
|
||||
|
||||
resultFuture
|
||||
.thenAccept(result -> {
|
||||
try {
|
||||
@@ -267,9 +266,8 @@ public class ConsumerService extends BaseService {
|
||||
return null;
|
||||
});
|
||||
} else {
|
||||
ChangeInvisibleTimeRequestHeader requestHeader = this.convertToChangeInvisibleTimeRequestHeader(ctx, request);
|
||||
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader,
|
||||
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request);
|
||||
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader);
|
||||
resultFuture
|
||||
.thenAccept(result -> {
|
||||
try {
|
||||
@@ -289,11 +287,11 @@ public class ConsumerService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
protected ChangeInvisibleTimeRequestHeader convertToChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) {
|
||||
return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy);
|
||||
protected ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) {
|
||||
return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, this.delayPolicy);
|
||||
}
|
||||
|
||||
protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) {
|
||||
protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) {
|
||||
return GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request);
|
||||
}
|
||||
|
||||
@@ -318,23 +316,19 @@ public class ConsumerService extends BaseService {
|
||||
this.readQueueSelector = readQueueSelector;
|
||||
}
|
||||
|
||||
public void setReceiveMessageHook(
|
||||
ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook) {
|
||||
public void setReceiveMessageHook(ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook) {
|
||||
this.receiveMessageHook = receiveMessageHook;
|
||||
}
|
||||
|
||||
public void setAckNoMatchedMessageHook(
|
||||
ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook) {
|
||||
public void setAckNoMatchedMessageHook(ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook) {
|
||||
this.ackNoMatchedMessageHook = ackNoMatchedMessageHook;
|
||||
}
|
||||
|
||||
public void setAckMessageHook(
|
||||
ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook) {
|
||||
public void setAckMessageHook(ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook) {
|
||||
this.ackMessageHook = ackMessageHook;
|
||||
}
|
||||
|
||||
public void setNackMessageHook(
|
||||
ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook) {
|
||||
public void setNackMessageHook(ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook) {
|
||||
this.nackMessageHook = nackMessageHook;
|
||||
}
|
||||
}
|
||||
|
||||
+14
-12
@@ -27,8 +27,8 @@ import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class DefaultWriteQueueSelector implements WriteQueueSelector {
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(DefaultWriteQueueSelector.class);
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(DefaultWriteQueueSelector.class);
|
||||
protected final TopicRouteCache topicRouteCache;
|
||||
|
||||
public DefaultWriteQueueSelector(TopicRouteCache topicRouteCache) {
|
||||
@@ -36,9 +36,12 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector {
|
||||
}
|
||||
|
||||
@Override
|
||||
public SelectableMessageQueue selectQueue(Context ctx, SendMessageRequest request,
|
||||
public SelectableMessageQueue selectQueue(
|
||||
Context ctx,
|
||||
SendMessageRequest request,
|
||||
SendMessageRequestHeader requestHeader,
|
||||
org.apache.rocketmq.common.message.Message message) {
|
||||
org.apache.rocketmq.common.message.Message message
|
||||
) {
|
||||
try {
|
||||
String topic = requestHeader.getTopic();
|
||||
String brokerName = "";
|
||||
@@ -47,19 +50,19 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector {
|
||||
}
|
||||
Integer queueId = requestHeader.getQueueId();
|
||||
String shardingKey = message.getProperty(MessageConst.PROPERTY_SHARDING_KEY);
|
||||
SelectableMessageQueue addressableMessageQueue;
|
||||
if (!StringUtils.isBlank(brokerName) && queueId != null) {
|
||||
SelectableMessageQueue targetMessageQueue;
|
||||
if (StringUtils.isNotBlank(brokerName) && queueId != null) {
|
||||
// Grpc client sendSelect situation
|
||||
addressableMessageQueue = selectTargetQueue(topic, brokerName, queueId);
|
||||
targetMessageQueue = selectTargetQueue(topic, brokerName, queueId);
|
||||
} else if (shardingKey != null) {
|
||||
// With shardingKey
|
||||
addressableMessageQueue = selectOrderQueue(topic, shardingKey);
|
||||
targetMessageQueue = selectOrderQueue(topic, shardingKey);
|
||||
} else {
|
||||
addressableMessageQueue = selectNormalQueue(topic);
|
||||
targetMessageQueue = selectNormalQueue(topic);
|
||||
}
|
||||
return addressableMessageQueue;
|
||||
return targetMessageQueue;
|
||||
} catch (Exception e) {
|
||||
log.error("error when select queue in DefaultMessageQueueSelector. request: {}", request, e);
|
||||
LOGGER.error("error when select queue in DefaultMessageQueueSelector. request: {}", request, e);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
@@ -68,8 +71,7 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector {
|
||||
return this.topicRouteCache.selectOneWriteQueue(topic, null);
|
||||
}
|
||||
|
||||
protected SelectableMessageQueue selectTargetQueue(String topic, String brokerName,
|
||||
int queueId) throws Exception {
|
||||
protected SelectableMessageQueue selectTargetQueue(String topic, String brokerName, int queueId) throws Exception {
|
||||
return this.topicRouteCache.selectOneWriteQueue(topic, brokerName, queueId);
|
||||
}
|
||||
|
||||
|
||||
+30
-25
@@ -31,8 +31,8 @@ import org.apache.rocketmq.client.producer.SendStatus;
|
||||
import org.apache.rocketmq.common.consumer.ReceiptHandle;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardProducer;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ProxyException;
|
||||
@@ -42,12 +42,15 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class ProducerService extends BaseService {
|
||||
|
||||
private final ForwardProducer producer;
|
||||
private volatile WriteQueueSelector writeQueueSelector;
|
||||
private volatile ResponseHook<SendMessageRequest, SendMessageResponse> sendMessageHook = null;
|
||||
private volatile ResponseHook<ForwardMessageToDeadLetterQueueRequest, ForwardMessageToDeadLetterQueueResponse> forwardMessageToDLQHook = null;
|
||||
private volatile ResponseHook<SendMessageRequest, SendMessageResponse> sendMessageHook;
|
||||
private volatile ResponseHook<ForwardMessageToDeadLetterQueueRequest, ForwardMessageToDeadLetterQueueResponse> forwardMessageToDLQHook;
|
||||
|
||||
|
||||
public ProducerService(ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
this.producer = connectorManager.getForwardProducer();
|
||||
writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
@@ -73,23 +76,24 @@ public class ProducerService extends BaseService {
|
||||
});
|
||||
|
||||
try {
|
||||
Pair<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> requestPair = this.convertSendMessageRequest(ctx, request);
|
||||
Pair<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> requestPair = this.buildSendMessageRequest(ctx, request);
|
||||
SendMessageRequestHeader requestHeader = requestPair.getLeft();
|
||||
org.apache.rocketmq.common.message.Message message = requestPair.getRight();
|
||||
SelectableMessageQueue addressableMessageQueue = writeQueueSelector.selectQueue(ctx, request, requestHeader, message);
|
||||
SelectableMessageQueue selectableMessageQueue = writeQueueSelector.selectQueue(ctx, request, requestHeader, message);
|
||||
|
||||
String topic = requestHeader.getTopic();
|
||||
if (addressableMessageQueue == null) {
|
||||
throw new ProxyException(Code.NOT_FOUND, "no writeable topic route for topic " + topic);
|
||||
if (selectableMessageQueue == null) {
|
||||
throw new ProxyException(Code.NOT_FOUND, "no writeable topic route for topic: " + topic);
|
||||
}
|
||||
|
||||
CompletableFuture<SendResult> sendResultCompletableFuture = this.connectorManager.getForwardProducer().sendMessage(
|
||||
addressableMessageQueue.getBrokerAddr(),
|
||||
addressableMessageQueue.getBrokerName(),
|
||||
// send message to broker.
|
||||
CompletableFuture<SendResult> sendResultCompletableFuture = this.producer.sendMessage(
|
||||
selectableMessageQueue.getBrokerAddr(),
|
||||
selectableMessageQueue.getBrokerName(),
|
||||
message,
|
||||
requestHeader,
|
||||
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT
|
||||
requestHeader
|
||||
);
|
||||
|
||||
sendResultCompletableFuture
|
||||
.thenAccept(result -> {
|
||||
try {
|
||||
@@ -108,20 +112,21 @@ public class ProducerService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
protected Pair<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> convertSendMessageRequest(
|
||||
protected Pair<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> buildSendMessageRequest(
|
||||
Context ctx, SendMessageRequest request) {
|
||||
return Pair.of(GrpcConverter.buildSendMessageRequestHeader(request), GrpcConverter.buildMessage(request.getMessage()));
|
||||
SendMessageRequestHeader requestHeader = GrpcConverter.buildSendMessageRequestHeader(request);
|
||||
org.apache.rocketmq.common.message.Message message = GrpcConverter.buildMessage(request.getMessage());
|
||||
return Pair.of(requestHeader, message);
|
||||
}
|
||||
|
||||
protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request,
|
||||
SendResult sendResult) {
|
||||
if (sendResult.getSendStatus() != SendStatus.SEND_OK) {
|
||||
protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request, SendResult result) {
|
||||
if (result.getSendStatus() != SendStatus.SEND_OK) {
|
||||
return SendMessageResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "send message failed, sendStatus=" + sendResult.getSendStatus()))
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "send message failed, sendStatus=" + result.getSendStatus()))
|
||||
.build();
|
||||
}
|
||||
|
||||
if (StringUtils.isNotBlank(sendResult.getTransactionId())) {
|
||||
if (StringUtils.isNotBlank(result.getTransactionId())) {
|
||||
Message message = request.getMessage();
|
||||
String group = GrpcConverter.wrapResourceWithNamespace(message.getSystemAttribute().getProducerGroup());
|
||||
String topic = GrpcConverter.wrapResourceWithNamespace(message.getTopic());
|
||||
@@ -130,8 +135,8 @@ public class ProducerService extends BaseService {
|
||||
|
||||
return SendMessageResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.setMessageId(StringUtils.defaultString(sendResult.getMsgId()))
|
||||
.setTransactionId(StringUtils.defaultString(sendResult.getTransactionId()))
|
||||
.setMessageId(StringUtils.defaultString(result.getMsgId()))
|
||||
.setTransactionId(StringUtils.defaultString(result.getTransactionId())) // use "" if transactionID is null.
|
||||
.build();
|
||||
}
|
||||
|
||||
@@ -143,12 +148,12 @@ public class ProducerService extends BaseService {
|
||||
forwardMessageToDLQHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
ConsumerSendMsgBackRequestHeader requestHeader = this.convertToConsumerSendMsgBackRequestHeader(ctx, request);
|
||||
CompletableFuture<RemotingCommand> resultFuture = this.connectorManager.getForwardProducer()
|
||||
.sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
ConsumerSendMsgBackRequestHeader requestHeader = this.buildConsumerSendMsgBackRequestHeader(ctx, request);
|
||||
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBack(brokerAddr, requestHeader);
|
||||
resultFuture
|
||||
.thenAccept(result ->
|
||||
future.complete(
|
||||
@@ -167,7 +172,7 @@ public class ProducerService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackRequestHeader(Context ctx,
|
||||
protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader(Context ctx,
|
||||
ForwardMessageToDeadLetterQueueRequest request) {
|
||||
return GrpcConverter.buildConsumerSendMsgBackRequestHeader(request);
|
||||
}
|
||||
|
||||
+27
-34
@@ -20,7 +20,6 @@ import apache.rocketmq.v1.Message;
|
||||
import apache.rocketmq.v1.Partition;
|
||||
import apache.rocketmq.v1.PullMessageRequest;
|
||||
import apache.rocketmq.v1.PullMessageResponse;
|
||||
import apache.rocketmq.v1.QueryOffsetPolicy;
|
||||
import apache.rocketmq.v1.QueryOffsetRequest;
|
||||
import apache.rocketmq.v1.QueryOffsetResponse;
|
||||
import com.google.protobuf.util.Timestamps;
|
||||
@@ -28,32 +27,31 @@ import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.stream.Collectors;
|
||||
import org.apache.rocketmq.client.consumer.PullResult;
|
||||
import org.apache.rocketmq.client.consumer.PullStatus;
|
||||
import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.proxy.common.utils.FilterUtils;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.DefaultForwardClient;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook;
|
||||
|
||||
public class PullMessageService extends BaseService {
|
||||
|
||||
private final DefaultForwardClient defaultForwardClient;
|
||||
private final DefaultForwardClient forwardClient;
|
||||
private final ForwardReadConsumer readConsumer;
|
||||
|
||||
private volatile ResponseHook<QueryOffsetRequest, QueryOffsetResponse> queryOffsetHook = null;
|
||||
|
||||
private volatile ResponseHook<PullMessageRequest, PullMessageResponse> pullMessageHook = null;
|
||||
private volatile ResponseHook<QueryOffsetRequest, QueryOffsetResponse> queryOffsetHook;
|
||||
private volatile ResponseHook<PullMessageRequest, PullMessageResponse> pullMessageHook;
|
||||
|
||||
public PullMessageService(ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
this.defaultForwardClient = connectorManager.getDefaultForwardClient();
|
||||
this.forwardClient = connectorManager.getDefaultForwardClient();
|
||||
this.readConsumer = connectorManager.getForwardReadConsumer();
|
||||
}
|
||||
|
||||
public CompletableFuture<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request) {
|
||||
@@ -63,24 +61,28 @@ public class PullMessageService extends BaseService {
|
||||
queryOffsetHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
Partition partition = request.getPartition();
|
||||
String topic = GrpcConverter.wrapResourceWithNamespace(partition.getTopic());
|
||||
|
||||
String brokerName = partition.getBroker().getName();
|
||||
int queueId = partition.getId();
|
||||
|
||||
CompletableFuture<Long> offsetFuture;
|
||||
if (request.getPolicy() == QueryOffsetPolicy.BEGINNING) {
|
||||
offsetFuture = CompletableFuture.completedFuture(0L);
|
||||
} else if (request.getPolicy() == QueryOffsetPolicy.END) {
|
||||
String brokerAddr = this.getBrokerAddr(ctx, brokerName);
|
||||
offsetFuture = this.defaultForwardClient.getMaxOffset(brokerAddr, topic, queueId, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
} else {
|
||||
long timestamp = Timestamps.toMillis(request.getTimePoint());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, brokerName);
|
||||
offsetFuture = this.defaultForwardClient.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
switch (request.getPolicy()) {
|
||||
case BEGINNING:
|
||||
offsetFuture = CompletableFuture.completedFuture(0L);
|
||||
break;
|
||||
case END:
|
||||
offsetFuture = this.forwardClient.getMaxOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId);
|
||||
break;
|
||||
default:
|
||||
long timestamp = Timestamps.toMillis(request.getTimePoint());
|
||||
offsetFuture = this.forwardClient.searchOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId, timestamp);
|
||||
}
|
||||
offsetFuture.thenAccept(result -> future.complete(
|
||||
|
||||
offsetFuture
|
||||
.thenAccept(result -> future.complete(
|
||||
QueryOffsetResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.setOffset(result)
|
||||
@@ -104,13 +106,12 @@ public class PullMessageService extends BaseService {
|
||||
});
|
||||
|
||||
try {
|
||||
PullMessageRequestHeader requestHeader = this.convertToPullMessageRequestHeader(ctx, request);
|
||||
PullMessageRequestHeader requestHeader = this.buildPullMessageRequestHeader(ctx, request);
|
||||
|
||||
String brokerName = request.getPartition().getBroker().getName();
|
||||
String brokerAddr = this.getBrokerAddr(ctx, brokerName);
|
||||
|
||||
CompletableFuture<PullResult> pullResultFuture = this.connectorManager.getForwardReadConsumer()
|
||||
.pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
CompletableFuture<PullResult> pullResultFuture = this.readConsumer.pullMessage(brokerAddr, requestHeader);
|
||||
pullResultFuture
|
||||
.thenAccept(pullResult -> {
|
||||
try {
|
||||
@@ -130,17 +131,9 @@ public class PullMessageService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
protected PullMessageRequestHeader convertToPullMessageRequestHeader(Context ctx, PullMessageRequest request) {
|
||||
// check filterExpression is correct or not
|
||||
GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression());
|
||||
|
||||
long timeRemaining = ctx.getDeadline()
|
||||
.timeRemaining(TimeUnit.MILLISECONDS);
|
||||
long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis();
|
||||
if (pollTime <= 0) {
|
||||
pollTime = timeRemaining;
|
||||
}
|
||||
return GrpcConverter.buildPullMessageRequestHeader(request, pollTime);
|
||||
protected PullMessageRequestHeader buildPullMessageRequestHeader(Context ctx, PullMessageRequest request) {
|
||||
checkSubscriptionData(request.getPartition().getTopic(), request.getFilterExpression());
|
||||
return GrpcConverter.buildPullMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx));
|
||||
}
|
||||
|
||||
protected PullMessageResponse convertToPullMessageResponse(Context ctx, PullMessageRequest request, PullResult result) {
|
||||
|
||||
+8
-10
@@ -56,11 +56,11 @@ public class RouteService extends BaseService {
|
||||
private final ProxyMode mode;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook = null;
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter;
|
||||
private volatile AssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook = null;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook;
|
||||
|
||||
public RouteService(ProxyMode mode, ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
@@ -79,8 +79,7 @@ public class RouteService extends BaseService {
|
||||
this.queryRouteHook = queryRouteHook;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentEndpointConverter(
|
||||
ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter) {
|
||||
public void setQueryAssignmentEndpointConverter(ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter) {
|
||||
this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter;
|
||||
}
|
||||
|
||||
@@ -88,8 +87,7 @@ public class RouteService extends BaseService {
|
||||
this.assignmentQueueSelector = assignmentQueueSelector;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentHook(
|
||||
ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook) {
|
||||
public void setQueryAssignmentHook(ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook) {
|
||||
this.queryAssignmentHook = queryAssignmentHook;
|
||||
}
|
||||
|
||||
@@ -102,8 +100,8 @@ public class RouteService extends BaseService {
|
||||
});
|
||||
|
||||
try {
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache()
|
||||
.getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic()));
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName);
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
List<QueueData> queueDataList = topicRouteData.getQueueDatas();
|
||||
List<BrokerData> brokerDataList = topicRouteData.getBrokerDatas();
|
||||
@@ -217,8 +215,8 @@ public class RouteService extends BaseService {
|
||||
List<Assignment> assignments = new ArrayList<>();
|
||||
List<SelectableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request);
|
||||
if (ProxyMode.isLocalMode(mode)) {
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache()
|
||||
.getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic()));
|
||||
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) {
|
||||
|
||||
+7
-7
@@ -46,8 +46,8 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
private final ChannelManager channelManager;
|
||||
private final ForwardProducer forwardProducer;
|
||||
|
||||
private volatile ResponseHook<TransactionStateCheckRequest, PollCommandResponse> checkTransactionStateHook = null;
|
||||
private volatile ResponseHook<EndTransactionRequest, EndTransactionResponse> endTransactionHook = null;
|
||||
private volatile ResponseHook<TransactionStateCheckRequest, PollCommandResponse> checkTransactionStateHook;
|
||||
private volatile ResponseHook<EndTransactionRequest, EndTransactionResponse> endTransactionHook;
|
||||
|
||||
public TransactionService(ConnectorManager connectorManager, ChannelManager channelManager) {
|
||||
super(connectorManager);
|
||||
@@ -63,11 +63,11 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
if (CollectionUtils.isEmpty(clientIdList)) {
|
||||
return;
|
||||
}
|
||||
|
||||
String clientId = clientIdList.get(ThreadLocalRandom.current().nextInt(clientIdList.size()));
|
||||
|
||||
GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId);
|
||||
String transactionId = checkData.getTransactionId().getProxyTransactionId();
|
||||
|
||||
String transactionId = checkData.getTransactionId().getProxyTransactionId();
|
||||
Message message = GrpcConverter.buildMessage(checkData.getMessageExt());
|
||||
PollCommandResponse response = PollCommandResponse.newBuilder()
|
||||
.setRecoverOrphanedTransactionCommand(
|
||||
@@ -76,6 +76,7 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
.setTransactionId(transactionId)
|
||||
.build()
|
||||
).build();
|
||||
|
||||
channel.writeAndFlush(response);
|
||||
if (this.checkTransactionStateHook != null) {
|
||||
this.checkTransactionStateHook.beforeResponse(ctx, checkData, response, null);
|
||||
@@ -97,10 +98,9 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
|
||||
try {
|
||||
TransactionId handle = TransactionId.decode(request.getTransactionId());
|
||||
String brokerAddr = RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr());
|
||||
EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request);
|
||||
this.forwardProducer.endTransaction(
|
||||
RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr()),
|
||||
requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
this.forwardProducer.endTransaction(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
future.complete(EndTransactionResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build());
|
||||
|
||||
Reference in New Issue
Block a user