From 7c322c8ab8dfbdd55cb63e250e547b84b751aa3c Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Tue, 10 May 2022 15:26:39 +0800 Subject: [PATCH] [ISSUE #3949] support topic message type --- .../attribute/AbstractRangeAttribute.java | 47 ++++++++++ .../common/attribute/IntRangeAttribute.java | 29 ++++++ .../common/attribute/LongRangeAttribute.java | 23 +---- .../common/constant/TopicMessageTypeName.java | 52 ++++++++++ .../rocketmq/proxy/config/ProxyConfig.java | 55 +++++++++++ .../proxy/connector/TopicConfigCache.java | 94 +++++++++++++++++++ 6 files changed, 281 insertions(+), 19 deletions(-) create mode 100644 common/src/main/java/org/apache/rocketmq/common/attribute/AbstractRangeAttribute.java create mode 100644 common/src/main/java/org/apache/rocketmq/common/attribute/IntRangeAttribute.java create mode 100644 common/src/main/java/org/apache/rocketmq/common/constant/TopicMessageTypeName.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java diff --git a/common/src/main/java/org/apache/rocketmq/common/attribute/AbstractRangeAttribute.java b/common/src/main/java/org/apache/rocketmq/common/attribute/AbstractRangeAttribute.java new file mode 100644 index 0000000000..7b773884e0 --- /dev/null +++ b/common/src/main/java/org/apache/rocketmq/common/attribute/AbstractRangeAttribute.java @@ -0,0 +1,47 @@ +/* + * 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.common.attribute; + +import static java.lang.String.format; + +public abstract class AbstractRangeAttribute> extends Attribute { + + protected final T min; + protected final T max; + protected final T defaultValue; + + public AbstractRangeAttribute(String name, boolean changeable, T min, T max, T defaultValue) { + super(name, changeable); + this.min = min; + this.max = max; + this.defaultValue = defaultValue; + } + + protected abstract T parse(String value); + + @Override + public void verify(String value) { + T l = parse(value); + if (l.compareTo(min) < 0 || l.compareTo(max) > 0) { + throw new RuntimeException(format("value is not in range(%s, %s)", min, max)); + } + } + + public T getDefaultValue() { + return defaultValue; + } +} diff --git a/common/src/main/java/org/apache/rocketmq/common/attribute/IntRangeAttribute.java b/common/src/main/java/org/apache/rocketmq/common/attribute/IntRangeAttribute.java new file mode 100644 index 0000000000..d55a3124ff --- /dev/null +++ b/common/src/main/java/org/apache/rocketmq/common/attribute/IntRangeAttribute.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.common.attribute; + +public class IntRangeAttribute extends AbstractRangeAttribute { + + public IntRangeAttribute(String name, boolean changeable, int min, int max, int defaultValue) { + super(name, changeable, min, max, defaultValue); + } + + @Override + protected Integer parse(String value) { + return Integer.parseInt(value); + } +} diff --git a/common/src/main/java/org/apache/rocketmq/common/attribute/LongRangeAttribute.java b/common/src/main/java/org/apache/rocketmq/common/attribute/LongRangeAttribute.java index eeeda72153..f4ccbc561e 100644 --- a/common/src/main/java/org/apache/rocketmq/common/attribute/LongRangeAttribute.java +++ b/common/src/main/java/org/apache/rocketmq/common/attribute/LongRangeAttribute.java @@ -16,29 +16,14 @@ */ package org.apache.rocketmq.common.attribute; -import static java.lang.String.format; - -public class LongRangeAttribute extends Attribute { - private final long min; - private final long max; - private final long defaultValue; +public class LongRangeAttribute extends AbstractRangeAttribute { public LongRangeAttribute(String name, boolean changeable, long min, long max, long defaultValue) { - super(name, changeable); - this.min = min; - this.max = max; - this.defaultValue = defaultValue; + super(name, changeable, min, max, defaultValue); } @Override - public void verify(String value) { - long l = Long.parseLong(value); - if (l < min || l > max) { - throw new RuntimeException(format("value is not in range(%d, %d)", min, max)); - } - } - - public long getDefaultValue() { - return defaultValue; + protected Long parse(String value) { + return Long.parseLong(value); } } diff --git a/common/src/main/java/org/apache/rocketmq/common/constant/TopicMessageTypeName.java b/common/src/main/java/org/apache/rocketmq/common/constant/TopicMessageTypeName.java new file mode 100644 index 0000000000..0050acee0f --- /dev/null +++ b/common/src/main/java/org/apache/rocketmq/common/constant/TopicMessageTypeName.java @@ -0,0 +1,52 @@ +/* + * 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.common.constant; + +public class TopicMessageTypeName { + public static final int INDEX_TRANSACTION = 4; + public static final int INDEX_DELAY = 3; + public static final int INDEX_FIFO = 2; + public static final int INDEX_NORMAL = 1; + + public static final int TRANSACTION = 0x1 << INDEX_TRANSACTION; + public static final int DELAY = 0x1 << INDEX_DELAY; + public static final int FIFO = 0x1 << INDEX_FIFO; + public static final int NORMAL = 0x1 << INDEX_NORMAL; + public static final int UNSPECIFIED = 0; + + public static final int ALL = NORMAL | FIFO | DELAY | TRANSACTION; + + public static boolean isUnspecified(final int type) { + return type == UNSPECIFIED; + } + + public static boolean isNormal(final int type) { + return (type & NORMAL) == NORMAL; + } + + public static boolean isFifo(final int type) { + return (type & FIFO) == FIFO; + } + + public static boolean isDelay(final int type) { + return (type & DELAY) == DELAY; + } + + public static boolean isTransaction(final int type) { + return (type & TRANSACTION) == TRANSACTION; + } +} 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 4ccc8bc57f..1eb5ab7d66 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 @@ -79,6 +79,13 @@ public class ProxyConfig { private int topicRouteServiceThreadPoolNums = PROCESSOR_NUMBER; private int topicRouteServiceThreadPoolQueueCapacity = 5000; + private int topicConfigCacheExpiredInSeconds = 20; + private int topicConfigCacheExecutorThreadNum = 3; + private int topicConfigCacheExecutorQueueCapacity = 1000; + private int topicConfigCacheMaxNum = 20000; + private int topicConfigThreadPoolNums = 36; + private int topicConfigThreadPoolQueueCapacity = 50000; + private int transactionHeartbeatThreadPoolNums = 20; private int transactionHeartbeatThreadPoolQueueCapacity = 200; private int transactionHeartbeatPeriodSecond = 20; @@ -400,6 +407,54 @@ public class ProxyConfig { this.topicRouteServiceThreadPoolQueueCapacity = topicRouteServiceThreadPoolQueueCapacity; } + public int getTopicConfigCacheExpiredInSeconds() { + return topicConfigCacheExpiredInSeconds; + } + + public void setTopicConfigCacheExpiredInSeconds(int topicConfigCacheExpiredInSeconds) { + this.topicConfigCacheExpiredInSeconds = topicConfigCacheExpiredInSeconds; + } + + public int getTopicConfigCacheExecutorThreadNum() { + return topicConfigCacheExecutorThreadNum; + } + + public void setTopicConfigCacheExecutorThreadNum(int topicConfigCacheExecutorThreadNum) { + this.topicConfigCacheExecutorThreadNum = topicConfigCacheExecutorThreadNum; + } + + public int getTopicConfigCacheExecutorQueueCapacity() { + return topicConfigCacheExecutorQueueCapacity; + } + + public void setTopicConfigCacheExecutorQueueCapacity(int topicConfigCacheExecutorQueueCapacity) { + this.topicConfigCacheExecutorQueueCapacity = topicConfigCacheExecutorQueueCapacity; + } + + public int getTopicConfigCacheMaxNum() { + return topicConfigCacheMaxNum; + } + + public void setTopicConfigCacheMaxNum(int topicConfigCacheMaxNum) { + this.topicConfigCacheMaxNum = topicConfigCacheMaxNum; + } + + public int getTopicConfigThreadPoolNums() { + return topicConfigThreadPoolNums; + } + + public void setTopicConfigThreadPoolNums(int topicConfigThreadPoolNums) { + this.topicConfigThreadPoolNums = topicConfigThreadPoolNums; + } + + public int getTopicConfigThreadPoolQueueCapacity() { + return topicConfigThreadPoolQueueCapacity; + } + + public void setTopicConfigThreadPoolQueueCapacity(int topicConfigThreadPoolQueueCapacity) { + this.topicConfigThreadPoolQueueCapacity = topicConfigThreadPoolQueueCapacity; + } + public int getTransactionHeartbeatThreadPoolNums() { return transactionHeartbeatThreadPoolNums; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java new file mode 100644 index 0000000000..c4ad518e43 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java @@ -0,0 +1,94 @@ +/* + * 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.connector; + +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.constant.LoggerName; +import org.apache.rocketmq.common.protocol.ResponseCode; +import org.apache.rocketmq.common.protocol.route.BrokerData; +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.config.ConfigurationManager; +import org.apache.rocketmq.proxy.config.ProxyConfig; +import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; +import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; + +public class TopicConfigCache { + private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); + + private final TopicRouteCache topicRouteCache; + private final ThreadPoolExecutor cacheRefreshExecutor; + private final LoadingCache topicConfigCache; + + private final DefaultForwardClient defaultClient; + + public TopicConfigCache(TopicRouteCache topicRouteCache, DefaultForwardClient client) { + this.topicRouteCache = topicRouteCache; + this.defaultClient = client; + + ProxyConfig config = ConfigurationManager.getProxyConfig(); + this.cacheRefreshExecutor = ThreadPoolMonitor.createAndMonitor( + config.getTopicConfigThreadPoolNums(), + config.getTopicConfigThreadPoolNums(), + 1000 * 60, + TimeUnit.MILLISECONDS, + "TopicConfigCacheRefresh", + config.getTopicConfigThreadPoolQueueCapacity() + ); + this.topicConfigCache = CacheBuilder.newBuilder() + .maximumSize(config.getTopicConfigCacheMaxNum()) + .refreshAfterWrite(config.getTopicConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) + .build(new TopicConfigCacheLoader()); + } + + public TopicConfigAndQueueMapping getTopicConfigAndQueueMapping(String topic) throws Exception { + return topicConfigCache.get(topic); + } + + protected class TopicConfigCacheLoader extends AbstractCacheLoader { + + public TopicConfigCacheLoader() { + super(cacheRefreshExecutor); + } + + @Override + protected TopicConfigAndQueueMapping getDirectly(String topic) throws Exception { + MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topic); + Optional brokerDataOptional = messageQueueWrapper.getTopicRouteData().getBrokerDatas().stream().findAny(); + if (!brokerDataOptional.isPresent()) { + throw new MQClientException(ResponseCode.TOPIC_NOT_EXIST, + "No topic route info in name server for the topic: " + topic); + } + + String brokerAddress = brokerDataOptional.get().selectBrokerAddr(); + return defaultClient.getTopicConfig(brokerAddress, topic); + } + + @Override + protected void onErr(String key, Exception e) { + log.error("load topic config failed. topic:{}", key, e); + } + } +}