diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java index 6b9ea2cf7e..61eb131ee2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java @@ -17,5 +17,6 @@ package org.apache.rocketmq.proxy.common.utils; public class ProxyUtils { + public static final int MAX_MSG_NUMS_FOR_POP_REQUEST = 32; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 1c531d92d9..6d56457a14 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -23,6 +23,8 @@ public class ProxyConfig { public final static String CONFIG_FILE_NAME = "rmq-proxy.json"; private static final int PROCESSOR_NUMBER = Runtime.getRuntime().availableProcessors(); + private String rocketMQClusterName = ""; + /** * configuration for ThreadPoolMonitor */ @@ -36,7 +38,6 @@ public class ProxyConfig { * gRPC */ private String proxyMode = ProxyMode.CLUSTER.name(); - private Boolean startGrpcServer = true; private Integer grpcServerPort = 8081; private boolean grpcTlsTestModeEnable = true; private String grpcTlsKeyPath = ConfigurationManager.getProxyHome() + "/conf/tls/rocketmq.key"; @@ -81,6 +82,8 @@ public class ProxyConfig { private int topicConfigCacheExpiredInSeconds = 20; private int topicConfigCacheMaxNum = 20000; + private int subscriptionGroupConfigCacheExpiredInSeconds = 20; + private int subscriptionGroupConfigCacheMaxNum = 20000; private int metadataThreadPoolNums = 3; private int metadataThreadPoolQueueCapacity = 1000; @@ -95,6 +98,14 @@ public class ProxyConfig { private boolean enableTopicMessageTypeCheck = true; + public String getRocketMQClusterName() { + return rocketMQClusterName; + } + + public void setRocketMQClusterName(String rocketMQClusterName) { + this.rocketMQClusterName = rocketMQClusterName; + } + public boolean isEnablePrintJstack() { return enablePrintJstack; } @@ -143,14 +154,6 @@ public class ProxyConfig { this.proxyMode = proxyMode; } - public Boolean getStartGrpcServer() { - return startGrpcServer; - } - - public void setStartGrpcServer(Boolean startGrpcServer) { - this.startGrpcServer = startGrpcServer; - } - public Integer getGrpcServerPort() { return grpcServerPort; } @@ -423,6 +426,22 @@ public class ProxyConfig { this.topicConfigCacheMaxNum = topicConfigCacheMaxNum; } + public int getSubscriptionGroupConfigCacheExpiredInSeconds() { + return subscriptionGroupConfigCacheExpiredInSeconds; + } + + public void setSubscriptionGroupConfigCacheExpiredInSeconds(int subscriptionGroupConfigCacheExpiredInSeconds) { + this.subscriptionGroupConfigCacheExpiredInSeconds = subscriptionGroupConfigCacheExpiredInSeconds; + } + + public int getSubscriptionGroupConfigCacheMaxNum() { + return subscriptionGroupConfigCacheMaxNum; + } + + public void setSubscriptionGroupConfigCacheMaxNum(int subscriptionGroupConfigCacheMaxNum) { + this.subscriptionGroupConfigCacheMaxNum = subscriptionGroupConfigCacheMaxNum; + } + public int getMetadataThreadPoolNums() { return metadataThreadPoolNums; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java index f63ee61816..50d5f180a0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java @@ -58,7 +58,7 @@ import org.apache.rocketmq.proxy.processor.MessagingProcessor; public class DefaultGrpcMessingActivity extends AbstractStartAndShutdown implements GrpcMessingActivity { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); - private GrpcClientSettingsManager grpcClientSettingsManager; + private final GrpcClientSettingsManager grpcClientSettingsManager; private final ReceiveMessageActivity receiveMessageActivity; private final AckMessageActivity ackMessageActivity; @@ -70,7 +70,7 @@ public class DefaultGrpcMessingActivity extends AbstractStartAndShutdown impleme private final ClientActivity clientActivity; protected DefaultGrpcMessingActivity(MessagingProcessor messagingProcessor) { - this.grpcClientSettingsManager = new GrpcClientSettingsManager(); + this.grpcClientSettingsManager = new GrpcClientSettingsManager(messagingProcessor); this.receiveMessageActivity = new ReceiveMessageActivity(messagingProcessor, this.grpcClientSettingsManager); this.ackMessageActivity = new AckMessageActivity(messagingProcessor, this.grpcClientSettingsManager); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java index 854507370b..5cbcc94244 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java @@ -95,7 +95,7 @@ public class ClientActivity extends AbstractMessingActivity { String clientId = context.getVal(GrpcContextConstants.CLIENT_ID); LanguageCode languageCode = context.getVal(GrpcContextConstants.LANGUAGE); - Settings clientSettings = grpcClientSettingsManager.getClientSettings(clientId); + Settings clientSettings = grpcClientSettingsManager.getClientSettings(context); switch (clientSettings.getClientType()) { case PRODUCER: { for (Resource topic : clientSettings.getPublishing().getTopicsList()) { @@ -150,7 +150,7 @@ public class ClientActivity extends AbstractMessingActivity { ProxyContext context = createContext(ctx); String clientId = context.getVal(GrpcContextConstants.CLIENT_ID); LanguageCode languageCode = context.getVal(GrpcContextConstants.LANGUAGE); - Settings clientSettings = grpcClientSettingsManager.getClientSettings(clientId); + Settings clientSettings = grpcClientSettingsManager.getClientSettings(context); switch (clientSettings.getClientType()) { case PRODUCER: @@ -225,7 +225,7 @@ public class ClientActivity extends AbstractMessingActivity { ProxyContext context = createContext(ctx); String clientId = context.getVal(GrpcContextConstants.CLIENT_ID); grpcClientSettingsManager.updateClientSettings(clientId, request.getSettings()); - Settings settings = grpcClientSettingsManager.getClientSettings(clientId); + Settings settings = grpcClientSettingsManager.getClientSettings(context); if (settings.hasPublishing()) { for (Resource topic : settings.getPublishing().getTopicsList()) { String topicName = GrpcConverter.wrapResourceWithNamespace(topic); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java index b5f8bde099..834a4d37b9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java @@ -17,15 +17,28 @@ package org.apache.rocketmq.proxy.grpc.v2.common; +import apache.rocketmq.v2.CustomizedBackoff; import apache.rocketmq.v2.ExponentialBackoff; import apache.rocketmq.v2.Publishing; import apache.rocketmq.v2.RetryPolicy; import apache.rocketmq.v2.Settings; import apache.rocketmq.v2.Subscription; +import com.google.protobuf.Duration; import com.google.protobuf.util.Durations; +import java.util.Arrays; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.stream.Collectors; +import org.apache.rocketmq.common.subscription.CustomizedRetryPolicy; +import org.apache.rocketmq.common.subscription.ExponentialRetryPolicy; +import org.apache.rocketmq.common.subscription.GroupRetryPolicy; +import org.apache.rocketmq.common.subscription.GroupRetryPolicyType; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; +import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; +import org.apache.rocketmq.proxy.grpc.v2.GrpcContextConstants; +import org.apache.rocketmq.proxy.processor.MessagingProcessor; public class GrpcClientSettingsManager { @@ -44,25 +57,79 @@ public class GrpcClientSettingsManager { .setMaxBodySize(4 * 1024 * 1024) .build()) .build(); - protected static final Settings DEFAULT_CONSUMER_SETTINGS = Settings.newBuilder() - .setBackoffPolicy(RetryPolicy.newBuilder() - .setMaxAttempts(3) - .setExponentialBackoff(ExponentialBackoff.newBuilder() - .setInitial(Durations.fromSeconds(1)) - .setMax(Durations.fromSeconds(3)) - .setMultiplier(2) - .build()) - .build()) + protected static final Settings DEFAULT_CONSUMER_SETTINGS = mergeSubscriptionData(Settings.newBuilder() .setSubscription(Subscription.newBuilder() - .setFifo(false) .setReceiveBatchSize(ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST) .setLongPollingTimeout(Durations.fromSeconds(30)) .build()) - .build(); + .build(), new SubscriptionGroupConfig()); + protected static final Map CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>(); - public Settings getClientSettings(String clientId) { - return CLIENT_SETTINGS_MAP.get(clientId); + private final MessagingProcessor messagingProcessor; + + public GrpcClientSettingsManager(MessagingProcessor messagingProcessor) { + this.messagingProcessor = messagingProcessor; + } + + public Settings getClientSettings(ProxyContext ctx) { + String clientId = ctx.getVal(GrpcContextConstants.CLIENT_ID); + Settings settings = CLIENT_SETTINGS_MAP.get(clientId); + if (settings.hasSubscription()) { + settings = mergeSubscriptionData(ctx, settings, + GrpcConverter.wrapResourceWithNamespace(settings.getSubscription().getGroup())); + } + return settings; + } + + private Settings mergeSubscriptionData(ProxyContext ctx, Settings settings, String consumerGroup) { + SubscriptionGroupConfig config = this.messagingProcessor.getSubscriptionGroupConfig(ctx, consumerGroup); + if (config == null) { + return settings; + } + + return mergeSubscriptionData(settings, config); + } + + protected static Settings mergeSubscriptionData(Settings settings, SubscriptionGroupConfig config) { + Settings.Builder resultSettingsBuilder = settings.toBuilder(); + + resultSettingsBuilder.getSubscriptionBuilder().setFifo(config.isConsumeMessageOrderly()); + + resultSettingsBuilder.getBackoffPolicyBuilder().setMaxAttempts(config.getRetryMaxTimes()); + + GroupRetryPolicy groupRetryPolicy = config.getGroupRetryPolicy(); + if (groupRetryPolicy.getType().equals(GroupRetryPolicyType.EXPONENTIAL)) { + ExponentialRetryPolicy exponentialRetryPolicy = groupRetryPolicy.getExponentialRetryPolicy(); + if (exponentialRetryPolicy == null) { + exponentialRetryPolicy = new ExponentialRetryPolicy(); + } + resultSettingsBuilder.getBackoffPolicyBuilder().setExponentialBackoff(convertToExponentialBackoff(exponentialRetryPolicy)); + } else { + CustomizedRetryPolicy customizedRetryPolicy = groupRetryPolicy.getCustomizedRetryPolicy(); + if (customizedRetryPolicy == null) { + customizedRetryPolicy = new CustomizedRetryPolicy(); + } + resultSettingsBuilder.getBackoffPolicyBuilder().setCustomizedBackoff(convertToCustomizedRetryPolicy(customizedRetryPolicy)); + } + + return resultSettingsBuilder.build(); + } + + protected static ExponentialBackoff convertToExponentialBackoff(ExponentialRetryPolicy retryPolicy) { + return ExponentialBackoff.newBuilder() + .setInitial(Durations.fromMillis(retryPolicy.getInitial())) + .setMax(Durations.fromMillis(retryPolicy.getMax())) + .setMultiplier(retryPolicy.getMultiplier()) + .build(); + } + + protected static CustomizedBackoff convertToCustomizedRetryPolicy(CustomizedRetryPolicy retryPolicy) { + List durationList = Arrays.stream(retryPolicy.getNext()) + .mapToObj(Durations::fromMillis).collect(Collectors.toList()); + return CustomizedBackoff.newBuilder() + .addAllNext(durationList) + .build(); } public void updateClientSettings(String clientId, Settings settings) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java index 23842c6fab..d411184149 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java @@ -21,7 +21,6 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.common.utils.FilterUtils; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager; -import org.apache.rocketmq.proxy.grpc.v2.GrpcContextConstants; import org.apache.rocketmq.proxy.processor.PopMessageResultFilter; public class PopMessageResultFilterImpl implements PopMessageResultFilter { @@ -34,7 +33,7 @@ public class PopMessageResultFilterImpl implements PopMessageResultFilter { @Override public FilterResult filterMessage(ProxyContext ctx, String consumerGroup, SubscriptionData subscriptionData, MessageExt messageExt) { - int maxAttempts = grpcClientSettingsManager.getClientSettings(ctx.getVal(GrpcContextConstants.CLIENT_ID)).getBackoffPolicy().getMaxAttempts(); + int maxAttempts = grpcClientSettingsManager.getClientSettings(ctx).getBackoffPolicy().getMaxAttempts(); if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { return FilterResult.NO_MATCH; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java index 673abca54d..f3965037a8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java @@ -20,6 +20,8 @@ import apache.rocketmq.v2.Code; import apache.rocketmq.v2.FilterExpression; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; +import apache.rocketmq.v2.Settings; +import apache.rocketmq.v2.Subscription; import com.google.protobuf.util.Durations; import io.grpc.Context; import io.grpc.stub.StreamObserver; @@ -31,7 +33,6 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.v2.AbstractMessingActivity; -import org.apache.rocketmq.proxy.grpc.v2.GrpcContextConstants; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter; import org.apache.rocketmq.proxy.processor.MessagingProcessor; @@ -50,63 +51,67 @@ public class ReceiveMessageActivity extends AbstractMessingActivity { public void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver) { ProxyContext proxyContext = createContext(ctx); - boolean fifo = false; - ReceiveMessageResponseStreamWriter writer = new ReceiveMessageResponseStreamWriter( this.messagingProcessor, responseObserver ); - long timeRemaining = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); - if (pollTime <= 0) { - pollTime = timeRemaining; - } - if (pollTime <= 0) { - writer.write(proxyContext, Code.MESSAGE_NOT_FOUND, "time remaining is too small"); - return; - } - - long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); - if (request.getAutoRenew()) { - invisibleTime = Durations.toMillis( - this.grpcClientSettingsManager.getClientSettings(proxyContext.getVal(GrpcContextConstants.CLIENT_ID)) - .getSubscription().getLongPollingTimeout() - ); - } - - String topic = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()); - String group = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); - FilterExpression filterExpression = request.getFilterExpression(); - SubscriptionData subscriptionData; try { - subscriptionData = FilterAPI.build(topic, filterExpression.getExpression(), - GrpcConverter.buildExpressionType(filterExpression.getType())); - } catch (Exception e) { - writer.write(proxyContext, Code.ILLEGAL_FILTER_EXPRESSION, e.getMessage()); - return; - } + Settings settings = this.grpcClientSettingsManager.getClientSettings(proxyContext); + Subscription subscription = settings.getSubscription(); + boolean fifo = subscription.getFifo(); - this.messagingProcessor.popMessage( - proxyContext, - new ReceiveMessageQueueSelector( - request.getMessageQueue().getBroker().getName() - ), - group, - topic, - request.getBatchSize(), - invisibleTime, - pollTime, - ConsumeInitMode.MAX, - subscriptionData, - fifo, - new PopMessageResultFilterImpl(grpcClientSettingsManager), - timeRemaining - ).thenAccept(popResult -> writer.write(proxyContext, request, popResult)) - .exceptionally(t -> { - writer.write(proxyContext, request, t); - return null; - }); + long timeRemaining = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS); + long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); + if (pollTime <= 0) { + pollTime = timeRemaining; + } + if (pollTime <= 0) { + writer.write(proxyContext, Code.MESSAGE_NOT_FOUND, "time remaining is too small"); + return; + } + + long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); + if (request.getAutoRenew()) { + invisibleTime = Durations.toMillis(subscription.getLongPollingTimeout() + ); + } + + String topic = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()); + String group = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + FilterExpression filterExpression = request.getFilterExpression(); + SubscriptionData subscriptionData; + try { + subscriptionData = FilterAPI.build(topic, filterExpression.getExpression(), + GrpcConverter.buildExpressionType(filterExpression.getType())); + } catch (Exception e) { + writer.write(proxyContext, Code.ILLEGAL_FILTER_EXPRESSION, e.getMessage()); + return; + } + + this.messagingProcessor.popMessage( + proxyContext, + new ReceiveMessageQueueSelector( + request.getMessageQueue().getBroker().getName() + ), + group, + topic, + request.getBatchSize(), + invisibleTime, + pollTime, + ConsumeInitMode.MAX, + subscriptionData, + fifo, + new PopMessageResultFilterImpl(grpcClientSettingsManager), + timeRemaining + ).thenAccept(popResult -> writer.write(proxyContext, request, popResult)) + .exceptionally(t -> { + writer.write(proxyContext, request, t); + return null; + }); + } catch (Throwable t) { + writer.write(proxyContext, request, t); + } } protected static class ReceiveMessageQueueSelector implements QueueSelector { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index 0269b752be..c159dec05d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -36,6 +36,7 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; import org.apache.rocketmq.common.thread.ThreadPoolMonitor; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.common.Address; @@ -112,6 +113,11 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen this.appendShutdown(this.consumerProcessorExecutor::shutdown); } + @Override + public SubscriptionGroupConfig getSubscriptionGroupConfig(ProxyContext ctx, String consumerGroupName) { + return this.serviceManager.getMetadataService().getSubscriptionGroupConfig(consumerGroupName); + } + @Override public ProxyTopicRouteData getTopicRouteDataForProxy(ProxyContext ctx, List
requestHostAndPortList, String topicName) throws Exception { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java index 10ac52c1c0..eee9f98731 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java @@ -34,6 +34,7 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; import org.apache.rocketmq.proxy.common.Address; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.common.StartAndShutdown; @@ -47,6 +48,11 @@ public interface MessagingProcessor extends StartAndShutdown { long DEFAULT_TIMEOUT_MILLS = Duration.ofSeconds(2).toMillis(); + SubscriptionGroupConfig getSubscriptionGroupConfig( + ProxyContext ctx, + String consumerGroupName + ); + ProxyTopicRouteData getTopicRouteDataForProxy( ProxyContext ctx, List
requestHostAndPortList, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java deleted file mode 100644 index 03742f200e..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java +++ /dev/null @@ -1,69 +0,0 @@ -/* - * 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.service.metadata; - -import java.util.concurrent.ThreadPoolExecutor; -import java.util.concurrent.TimeUnit; -import org.apache.rocketmq.common.attribute.TopicMessageType; -import org.apache.rocketmq.common.constant.LoggerName; -import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; -import org.apache.rocketmq.common.thread.ThreadPoolMonitor; -import org.apache.rocketmq.logging.InternalLogger; -import org.apache.rocketmq.logging.InternalLoggerFactory; -import org.apache.rocketmq.proxy.common.AbstractCacheLoader; -import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; -import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.config.ProxyConfig; - -public abstract class AbstractMetadataService extends AbstractStartAndShutdown implements MetadataService { - protected static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); - protected final ThreadPoolExecutor cacheRefreshExecutor; - - public AbstractMetadataService() { - ProxyConfig config = ConfigurationManager.getProxyConfig(); - this.cacheRefreshExecutor = ThreadPoolMonitor.createAndMonitor( - config.getMetadataThreadPoolNums(), - config.getMetadataThreadPoolNums(), - 1000 * 60, - TimeUnit.MILLISECONDS, - "MetadataCacheRefresh", - config.getMetadataThreadPoolQueueCapacity() - ); - } - - public abstract TopicMessageType getTopicMessageType(String topic); - - protected abstract class AbstractTopicConfigCacheLoader extends AbstractCacheLoader { - - public AbstractTopicConfigCacheLoader() { - super(cacheRefreshExecutor); - } - - protected abstract TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception; - - @Override - public TopicConfigAndQueueMapping getDirectly(String topic) throws Exception { - return loadTopicConfig(topic); - } - - @Override - protected void onErr(String key, Exception e) { - log.error("load topic config failed. topic:{}", key, e); - } - } -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java index edf011608d..3be7b051b9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java @@ -20,32 +20,64 @@ package org.apache.rocketmq.proxy.service.metadata; import com.google.common.cache.CacheBuilder; import com.google.common.cache.LoadingCache; import java.util.Optional; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.common.attribute.TopicMessageType; +import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; +import org.apache.rocketmq.common.thread.ThreadPoolMonitor; +import org.apache.rocketmq.logging.InternalLogger; +import org.apache.rocketmq.logging.InternalLoggerFactory; +import org.apache.rocketmq.proxy.common.AbstractCacheLoader; +import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.route.TopicRouteHelper; import org.apache.rocketmq.proxy.service.route.TopicRouteService; -public class ClusterMetadataService extends AbstractMetadataService { +public class ClusterMetadataService extends AbstractStartAndShutdown implements MetadataService { + protected static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); + private static final long DEFAULT_TIMEOUT = 3000; + + private final ThreadPoolExecutor cacheRefreshExecutor; + private final TopicRouteService topicRouteService; + private final MQClientAPIFactory mqClientAPIFactory; + private final LoadingCache topicCache; - private TopicRouteService topicRouteService; private final static TopicConfigAndQueueMapping EMPTY_TOPIC_CONFIG = new TopicConfigAndQueueMapping(); + private final LoadingCache subscriptionGroupConfigCache; + private final static SubscriptionGroupConfig EMPTY_SUBSCRIPTION_GROUP_CONFIG = new SubscriptionGroupConfig(); + public ClusterMetadataService(TopicRouteService topicRouteService, MQClientAPIFactory mqClientAPIFactory) { this.topicRouteService = topicRouteService; + this.mqClientAPIFactory = mqClientAPIFactory; + ProxyConfig config = ConfigurationManager.getProxyConfig(); + this.cacheRefreshExecutor = ThreadPoolMonitor.createAndMonitor( + config.getMetadataThreadPoolNums(), + config.getMetadataThreadPoolNums(), + 1000 * 60, + TimeUnit.MILLISECONDS, + "MetadataCacheRefresh", + config.getMetadataThreadPoolQueueCapacity() + ); this.topicCache = CacheBuilder.newBuilder() .maximumSize(config.getTopicConfigCacheMaxNum()) .refreshAfterWrite(config.getTopicConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) - .build(new ClusterTopicConfigCacheLoader(mqClientAPIFactory)); + .build(new ClusterTopicConfigCacheLoader()); + this.subscriptionGroupConfigCache = CacheBuilder.newBuilder() + .maximumSize(config.getSubscriptionGroupConfigCacheMaxNum()) + .refreshAfterWrite(config.getSubscriptionGroupConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) + .build(new ClusterSubscriptionGroupConfigCacheLoader()); } - @Override public TopicMessageType getTopicMessageType(String topic) { + @Override + public TopicMessageType getTopicMessageType(String topic) { TopicConfigAndQueueMapping topicConfigAndQueueMapping; try { topicConfigAndQueueMapping = topicCache.get(topic); @@ -58,28 +90,74 @@ public class ClusterMetadataService extends AbstractMetadataService { return topicConfigAndQueueMapping.getTopicMessageType(); } - protected class ClusterTopicConfigCacheLoader extends AbstractTopicConfigCacheLoader { - private final MQClientAPIFactory mqClientAPIFactory; + @Override + public SubscriptionGroupConfig getSubscriptionGroupConfig(String group) { + SubscriptionGroupConfig config; + try { + config = this.subscriptionGroupConfigCache.get(group); + } catch (Exception e) { + return null; + } + if (config == EMPTY_SUBSCRIPTION_GROUP_CONFIG) { + return null; + } + return config; + } - public ClusterTopicConfigCacheLoader(MQClientAPIFactory mqClientAPIFactory) { - this.mqClientAPIFactory = mqClientAPIFactory; + protected class ClusterSubscriptionGroupConfigCacheLoader extends AbstractCacheLoader { + + public ClusterSubscriptionGroupConfigCacheLoader() { + super(cacheRefreshExecutor); } - @Override protected TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception { - try { - Optional brokerDataOptional = topicRouteService.getAllMessageQueueView(topic).getTopicRouteData().getBrokerDatas().stream().findAny(); - if (brokerDataOptional.isPresent()) { - String brokerAddress = topicRouteService.getBrokerAddr(brokerDataOptional.get().getBrokerName()); - return this.mqClientAPIFactory.getClient().getTopicConfig(brokerAddress, topic, 1000L); - } - return EMPTY_TOPIC_CONFIG; - } catch (MQClientException e) { - if (TopicRouteHelper.isTopicNotExistError(e)) { - log.warn("topic is not exist", e); - return EMPTY_TOPIC_CONFIG; - } - throw e; + @Override + protected SubscriptionGroupConfig getDirectly(String consumerGroup) throws Exception { + ProxyConfig config = ConfigurationManager.getProxyConfig(); + String clusterName = config.getRocketMQClusterName(); + Optional brokerDataOptional = findOneBroker(clusterName); + if (brokerDataOptional.isPresent()) { + String brokerAddress = brokerDataOptional.get().selectBrokerAddr(); + return mqClientAPIFactory.getClient().getSubscriptionGroupConfig(brokerAddress, consumerGroup, DEFAULT_TIMEOUT); } + return EMPTY_SUBSCRIPTION_GROUP_CONFIG; + } + + @Override + protected void onErr(String consumerGroup, Exception e) { + log.error("load subscription config failed. consumerGroup:{}", consumerGroup, e); + } + } + + protected class ClusterTopicConfigCacheLoader extends AbstractCacheLoader { + + public ClusterTopicConfigCacheLoader() { + super(cacheRefreshExecutor); + } + + @Override + protected TopicConfigAndQueueMapping getDirectly(String topic) throws Exception { + Optional brokerDataOptional = findOneBroker(topic); + if (brokerDataOptional.isPresent()) { + String brokerAddress = brokerDataOptional.get().selectBrokerAddr(); + return mqClientAPIFactory.getClient().getTopicConfig(brokerAddress, topic, DEFAULT_TIMEOUT); + } + return EMPTY_TOPIC_CONFIG; + } + + @Override + protected void onErr(String key, Exception e) { + log.error("load topic config failed. topic:{}", key, e); + } + } + + protected Optional findOneBroker(String topic) throws Exception { + try { + return topicRouteService.getAllMessageQueueView(topic).getTopicRouteData().getBrokerDatas().stream().findAny(); + } catch (MQClientException e) { + if (TopicRouteHelper.isTopicNotExistError(e)) { + return Optional.empty(); + } + throw e; } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java index 0c2689268d..6f06f84888 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java @@ -17,52 +17,29 @@ package org.apache.rocketmq.proxy.service.metadata; -import com.google.common.cache.CacheBuilder; -import com.google.common.cache.LoadingCache; -import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.common.TopicConfig; import org.apache.rocketmq.common.attribute.TopicMessageType; -import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; -import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.config.ProxyConfig; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; -public class LocalMetadataService extends AbstractMetadataService { - private final LoadingCache topicCache; +public class LocalMetadataService implements MetadataService { + private final BrokerController brokerController; public LocalMetadataService(BrokerController brokerController) { - ProxyConfig config = ConfigurationManager.getProxyConfig(); - - this.topicCache = CacheBuilder.newBuilder() - .maximumSize(config.getTopicConfigCacheMaxNum()) - .refreshAfterWrite(config.getTopicConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) - .build(new LocalTopicConfigCacheLoader(brokerController)); + this.brokerController = brokerController; } @Override public TopicMessageType getTopicMessageType(String topic) { - try { - TopicConfigAndQueueMapping topicConfigAndQueueMapping = topicCache.get(topic); - if (topicConfigAndQueueMapping == null) { - return TopicMessageType.UNSPECIFIED; - } - return topicConfigAndQueueMapping.getTopicMessageType(); - } catch (Exception e) { - log.error("getTopicMessageType error", topic, e); + TopicConfig topicConfig = brokerController.getTopicConfigManager().selectTopicConfig(topic); + if (topicConfig == null) { return TopicMessageType.UNSPECIFIED; } + return topicConfig.getTopicMessageType(); } - protected class LocalTopicConfigCacheLoader extends AbstractTopicConfigCacheLoader { - private final BrokerController brokerController; - - public LocalTopicConfigCacheLoader(BrokerController brokerController) { - this.brokerController = brokerController; - } - - @Override protected TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception { - TopicConfig topicConfig = brokerController.getTopicConfigManager().selectTopicConfig(topic); - return new TopicConfigAndQueueMapping(topicConfig, null); - } + @Override + public SubscriptionGroupConfig getSubscriptionGroupConfig(String group) { + return this.brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().get(group); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java index a87ec1484f..6951845e52 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java @@ -18,7 +18,11 @@ package org.apache.rocketmq.proxy.service.metadata; import org.apache.rocketmq.common.attribute.TopicMessageType; +import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; public interface MetadataService { + TopicMessageType getTopicMessageType(String topic); + + SubscriptionGroupConfig getSubscriptionGroupConfig(String group); } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java index c155b382e8..751c62f7c0 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java @@ -51,9 +51,11 @@ public class BaseActivityTest extends InitConfigAndLoggerTest { protected static final String LOCAL_ADDR = "127.0.0.1:8080"; protected Metadata metadata = new Metadata(); + protected static final String CLIENT_ID = "client-id" + UUID.randomUUID(); + public void before() throws Throwable { super.before(); - metadata.put(InterceptorConstants.CLIENT_ID, "client-id" + UUID.randomUUID()); + metadata.put(InterceptorConstants.CLIENT_ID, CLIENT_ID); metadata.put(InterceptorConstants.LANGUAGE, "JAVA"); metadata.put(InterceptorConstants.REMOTE_ADDRESS, REMOTE_ADDR); metadata.put(InterceptorConstants.LOCAL_ADDRESS, LOCAL_ADDR); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java index 86cb54dd59..ae5bef6eec 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java @@ -24,6 +24,7 @@ import apache.rocketmq.v2.MessageQueue; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.Settings; import io.grpc.stub.ServerCallStreamObserver; import io.grpc.stub.StreamObserver; import java.util.ArrayList; @@ -63,6 +64,8 @@ public class ReceiveMessageActivityTest extends BaseActivityTest { ArgumentCaptor responseArgumentCaptor = ArgumentCaptor.forClass(ReceiveMessageResponse.class); doNothing().when(receiveStreamObserver).onNext(responseArgumentCaptor.capture()); + when(this.grpcClientSettingsManager.getClientSettings(any())).thenReturn(Settings.newBuilder().getDefaultInstanceForType()); + this.receiveMessageActivity.receiveMessage( createContext(), ReceiveMessageRequest.newBuilder() @@ -85,6 +88,8 @@ public class ReceiveMessageActivityTest extends BaseActivityTest { ArgumentCaptor responseArgumentCaptor = ArgumentCaptor.forClass(ReceiveMessageResponse.class); doNothing().when(receiveStreamObserver).onNext(responseArgumentCaptor.capture()); + when(this.grpcClientSettingsManager.getClientSettings(any())).thenReturn(Settings.newBuilder().getDefaultInstanceForType()); + PopResult popResult = new PopResult(PopStatus.NO_NEW_MSG, new ArrayList<>()); when(this.messagingProcessor.popMessage( any(),