mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] support topic message type
This commit is contained in:
@@ -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<T extends Comparable<T>> 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;
|
||||
}
|
||||
}
|
||||
@@ -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<Integer> {
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -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<Long> {
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<String /* topicName */, TopicConfigAndQueueMapping> 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<String, TopicConfigAndQueueMapping> {
|
||||
|
||||
public TopicConfigCacheLoader() {
|
||||
super(cacheRefreshExecutor);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected TopicConfigAndQueueMapping getDirectly(String topic) throws Exception {
|
||||
MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topic);
|
||||
Optional<BrokerData> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user