[ISSUE #3949] support config retryPolicy and consumeMessageOrderly in subscriptionGroupConfig

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:36 +08:00
committed by zhouxiang
parent 20e9bdd3e4
commit 95b8105ffa
15 changed files with 305 additions and 205 deletions
@@ -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;
}
@@ -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;
}
@@ -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);
@@ -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);
@@ -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<String, Settings> 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<Duration> durationList = Arrays.stream(retryPolicy.getNext())
.mapToObj(Durations::fromMillis).collect(Collectors.toList());
return CustomizedBackoff.newBuilder()
.addAllNext(durationList)
.build();
}
public void updateClientSettings(String clientId, Settings settings) {
@@ -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;
}
@@ -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<ReceiveMessageResponse> 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 {
@@ -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<Address> requestHostAndPortList,
String topicName) throws Exception {
@@ -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<Address> requestHostAndPortList,
@@ -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<String, TopicConfigAndQueueMapping> {
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);
}
}
}
@@ -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<String, TopicConfigAndQueueMapping> topicCache;
private TopicRouteService topicRouteService;
private final static TopicConfigAndQueueMapping EMPTY_TOPIC_CONFIG = new TopicConfigAndQueueMapping();
private final LoadingCache<String, SubscriptionGroupConfig> 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<String, SubscriptionGroupConfig> {
public ClusterSubscriptionGroupConfigCacheLoader() {
super(cacheRefreshExecutor);
}
@Override protected TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception {
try {
Optional<BrokerData> 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<BrokerData> 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<String, TopicConfigAndQueueMapping> {
public ClusterTopicConfigCacheLoader() {
super(cacheRefreshExecutor);
}
@Override
protected TopicConfigAndQueueMapping getDirectly(String topic) throws Exception {
Optional<BrokerData> 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<BrokerData> 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;
}
}
}
@@ -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<String, TopicConfigAndQueueMapping> 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);
}
}
@@ -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);
}
@@ -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);
@@ -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<ReceiveMessageResponse> 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<ReceiveMessageResponse> 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(),