This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 20bba97 [INLONG-3108][TubeMQ] Optimize the implementation of
KeepAlive Interface (#3123)
20bba97 is described below
commit 20bba9710831182a2ca70d83e47317407e1ddc13
Author: gosonzhang <[email protected]>
AuthorDate: Mon Mar 14 17:37:23 2022 +0800
[INLONG-3108][TubeMQ] Optimize the implementation of KeepAlive Interface
(#3123)
---
.../server/master/metamanage/MetaDataManager.java | 4 +-
.../KeepAliveService.java} | 41 +-
.../MetaConfigObserver.java} | 4 +-
.../metamanage/metastore/MetaStoreService.java | 6 +-
...{MetaStoreMapper.java => MetaConfigMapper.java} | 13 +-
...apperImpl.java => AbsMetaConfigMapperImpl.java} | 54 +-
.../impl/bdbimpl/BdbMetaConfigMapperImpl.java | 597 +++++++++++++++++++++
.../impl/bdbimpl/BdbMetaStoreServiceImpl.java | 22 +-
.../nodemanage/nodebroker/DefBrokerRunManager.java | 4 +-
...SamplePrint.java => MetaConfigSamplePrint.java} | 214 ++++----
10 files changed, 790 insertions(+), 169 deletions(-)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/MetaDataManager.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/MetaDataManager.java
index f35104e..75f17ab 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/MetaDataManager.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/MetaDataManager.java
@@ -47,7 +47,7 @@ import
org.apache.inlong.tubemq.server.common.utils.WebParameterUtils;
import org.apache.inlong.tubemq.server.master.MasterConfig;
import org.apache.inlong.tubemq.server.master.TMaster;
import org.apache.inlong.tubemq.server.master.bdbstore.MasterGroupStatus;
-import
org.apache.inlong.tubemq.server.master.metamanage.keepalive.AliveObserver;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbMetaStoreServiceImpl;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaStoreService;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BaseEntity;
@@ -144,7 +144,7 @@ public class MetaDataManager implements Server {
logger.info("BrokerConfManager StoreService stopped");
}
- public void registerObserver(AliveObserver eventObserver) {
+ public void registerObserver(MetaConfigObserver eventObserver) {
metaStoreService.registerObserver(eventObserver);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/KeepAlive.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
similarity index 57%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/KeepAlive.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
index 57f1ae0..ecce26e 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/KeepAlive.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
@@ -15,42 +15,63 @@
* limitations under the License.
*/
-package org.apache.inlong.tubemq.server.master.metamanage.keepalive;
+package org.apache.inlong.tubemq.server.master.metamanage.metastore;
import java.net.InetSocketAddress;
+
+import org.apache.inlong.tubemq.server.Server;
import org.apache.inlong.tubemq.server.master.bdbstore.MasterGroupStatus;
import org.apache.inlong.tubemq.server.master.web.model.ClusterGroupVO;
-public interface KeepAlive {
+public interface KeepAliveService extends Server {
/**
- * Whether this node is the master role
+ * Whether this node is the Master role now
*
- * @return true if is master role or else
+ * @return true if is Master role or else
*/
boolean isMasterNow();
+ /**
+ * Get the timestamp when the current node becomes the Master role
+ *
+ * @return the since time
+ */
long getMasterSinceTime();
+ /**
+ * Get current Master address
+ *
+ * @return the current Master address
+ */
InetSocketAddress getMasterAddress();
/**
- * Whether the primary node in active
+ * Whether the Master node in active
*
- * @return true for active, false for inactive
+ * @return true for Master, false for Slave
*/
boolean isPrimaryNodeActive();
+ /**
+ * Transfer Master role to other replica node
+ *
+ * @throws Exception the exception information
+ */
void transferMaster() throws Exception;
/**
- * Register node role switching event observer
+ * Get group address info
*
- * @param eventObserver the event observer
+ * @return the group address information
*/
- void registerObserver(AliveObserver eventObserver);
-
ClusterGroupVO getGroupAddressStrInfo();
+ /**
+ * Get Master group status
+ *
+ * @param isFromHeartbeat whether called by hb thread
+ * @return the master group status
+ */
MasterGroupStatus getMasterGroupStatus(boolean isFromHeartbeat);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/AliveObserver.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaConfigObserver.java
similarity index 86%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/AliveObserver.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaConfigObserver.java
index ca8a4b0..f428b51 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/keepalive/AliveObserver.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaConfigObserver.java
@@ -15,9 +15,9 @@
* limitations under the License.
*/
-package org.apache.inlong.tubemq.server.master.metamanage.keepalive;
+package org.apache.inlong.tubemq.server.master.metamanage.metastore;
-public interface AliveObserver {
+public interface MetaConfigObserver {
void clearCacheData();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
index fb7821b..1a3073c 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
@@ -22,8 +22,6 @@ import java.util.Map;
import java.util.Set;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
import org.apache.inlong.tubemq.server.Server;
-import
org.apache.inlong.tubemq.server.master.metamanage.keepalive.AliveObserver;
-import org.apache.inlong.tubemq.server.master.metamanage.keepalive.KeepAlive;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.ClusterSettingEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.GroupConsumeCtrlEntity;
@@ -31,7 +29,7 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.Gr
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.TopicCtrlEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.TopicDeployEntity;
-public interface MetaStoreService extends KeepAlive, Server {
+public interface MetaStoreService extends KeepAliveService, Server {
boolean checkStoreStatus(boolean checkIsMaster, ProcessResult result);
@@ -279,7 +277,7 @@ public interface MetaStoreService extends KeepAlive, Server
{
boolean delGroupConsumeCtrlConf(String operator, String recordKey,
StringBuilder strBuff, ProcessResult
result);
- void registerObserver(AliveObserver eventObserver);
+ void registerObserver(MetaConfigObserver eventObserver);
boolean isTopicNameInUsed(String topicName);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaStoreMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaConfigMapper.java
similarity index 95%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaStoreMapper.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaConfigMapper.java
index 1470b7f..f987509 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaStoreMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/MetaConfigMapper.java
@@ -23,6 +23,8 @@ import java.util.Set;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
import org.apache.inlong.tubemq.server.common.statusdef.ManageStatus;
import org.apache.inlong.tubemq.server.common.statusdef.TopicStatus;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.KeepAliveService;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BaseEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.ClusterSettingEntity;
@@ -32,7 +34,16 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.To
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.TopicDeployEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.TopicPropGroup;
-public interface MetaStoreMapper {
+public interface MetaConfigMapper extends KeepAliveService {
+
+ /**
+ * Register meta configure change observer
+ *
+ * @param eventObserver the event observer
+ */
+ void regMetaConfigObserver(MetaConfigObserver eventObserver);
+
+ boolean checkStoreStatus(boolean checkIsMaster, ProcessResult result);
/**
* Add or update cluster default setting
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaStoreMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
similarity index 95%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaStoreMapperImpl.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
index 933cd86..d5ce9e7 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaStoreMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
@@ -23,7 +23,6 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
-import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.codec.binary.StringUtils;
import org.apache.inlong.tubemq.corebase.TBaseConstants;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
@@ -32,8 +31,8 @@ import
org.apache.inlong.tubemq.server.common.statusdef.ManageStatus;
import org.apache.inlong.tubemq.server.common.statusdef.TopicStatus;
import org.apache.inlong.tubemq.server.common.utils.RowLock;
import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
-import
org.apache.inlong.tubemq.server.master.metamanage.keepalive.AliveObserver;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.MetaStoreMapper;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.MetaConfigMapper;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BaseEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.ClusterSettingEntity;
@@ -51,39 +50,42 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.To
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-public class AbsMetaStoreMapperImpl implements MetaStoreMapper {
+public abstract class AbsMetaConfigMapperImpl implements MetaConfigMapper {
protected static final Logger logger =
- LoggerFactory.getLogger(AbsMetaStoreMapperImpl.class);
- // service status
- // 0 stopped, 1 starting, 2 started, 3 stopping
- private final AtomicInteger srvStatus = new AtomicInteger(0);
- // the observers focusing on active-standby switching
- private final List<AliveObserver> eventObservers = new ArrayList<>();
-
+ LoggerFactory.getLogger(AbsMetaConfigMapperImpl.class);
// row lock.
private final RowLock metaRowLock;
// default cluster setting
private static final ClusterSettingEntity defClusterSetting =
new ClusterSettingEntity().fillDefaultValue();
// cluster default setting
- private ClusterConfigMapper clusterConfigMapper;
+ protected ClusterConfigMapper clusterConfigMapper;
// broker configure
- private BrokerConfigMapper brokerConfigMapper;
+ protected BrokerConfigMapper brokerConfigMapper;
// topic deployment configure
- private TopicDeployMapper topicDeployMapper;
+ protected TopicDeployMapper topicDeployMapper;
// topic control configure
- private TopicCtrlMapper topicCtrlMapper;
+ protected TopicCtrlMapper topicCtrlMapper;
// group resource control configure
- private GroupResCtrlMapper groupResCtrlMapper;
+ protected GroupResCtrlMapper groupResCtrlMapper;
// group consume control configure
- private ConsumeCtrlMapper consumeCtrlMapper;
+ protected ConsumeCtrlMapper consumeCtrlMapper;
+ // the observers focusing on active-standby switching
+ private final List<MetaConfigObserver> eventObservers = new ArrayList<>();
- public AbsMetaStoreMapperImpl(int rowLockWaiDurMs) {
+ public AbsMetaConfigMapperImpl(int rowLockWaiDurMs) {
this.metaRowLock =
new RowLock("MetaData-RowLock", rowLockWaiDurMs);
}
@Override
+ public void regMetaConfigObserver(MetaConfigObserver eventObserver) {
+ if (eventObserver != null) {
+ eventObservers.add(eventObserver);
+ }
+ }
+
+ @Override
public boolean addOrUpdClusterDefSetting(BaseEntity opEntity,
int brokerPort, int brokerTlsPort,
int brokerWebPort, int
maxMsgSizeMB,
@@ -1157,21 +1159,13 @@ public class AbsMetaStoreMapperImpl implements
MetaStoreMapper {
}
/**
- * Initial meta-data stores.
- *
- */
- protected void initMetaStore() {
-
- }
-
- /**
* Reload meta-data stores.
*
* @param strBuff the string buffer
*/
- private void reloadMetaStore(StringBuilder strBuff) {
+ protected void reloadMetaStore(StringBuilder strBuff) {
// Clear observers' cache data.
- for (AliveObserver observer : eventObservers) {
+ for (MetaConfigObserver observer : eventObservers) {
observer.clearCacheData();
}
// Load the latest meta-data from persistent
@@ -1182,7 +1176,7 @@ public class AbsMetaStoreMapperImpl implements
MetaStoreMapper {
groupResCtrlMapper.loadConfig(strBuff);
consumeCtrlMapper.loadConfig(strBuff);
// load the latest meta-data to observers
- for (AliveObserver observer : eventObservers) {
+ for (MetaConfigObserver observer : eventObservers) {
observer.reloadCacheData();
}
}
@@ -1191,7 +1185,7 @@ public class AbsMetaStoreMapperImpl implements
MetaStoreMapper {
* Close meta-data stores.
*
*/
- private void closeMetaStore() {
+ protected void closeMetaStore() {
brokerConfigMapper.close();
topicDeployMapper.close();
groupResCtrlMapper.close();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaConfigMapperImpl.java
new file mode 100644
index 0000000..f35a85a
--- /dev/null
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaConfigMapperImpl.java
@@ -0,0 +1,597 @@
+/**
+ * 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.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl;
+
+import java.io.File;
+import java.io.IOException;
+import java.net.InetSocketAddress;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import com.sleepycat.je.DatabaseException;
+import com.sleepycat.je.Durability;
+import com.sleepycat.je.EnvironmentConfig;
+import com.sleepycat.je.EnvironmentFailureException;
+import com.sleepycat.je.rep.InsufficientLogException;
+import com.sleepycat.je.rep.NetworkRestore;
+import com.sleepycat.je.rep.NetworkRestoreConfig;
+import com.sleepycat.je.rep.NodeState;
+import com.sleepycat.je.rep.ReplicatedEnvironment;
+import com.sleepycat.je.rep.ReplicationConfig;
+import com.sleepycat.je.rep.ReplicationGroup;
+import com.sleepycat.je.rep.ReplicationMutableConfig;
+import com.sleepycat.je.rep.ReplicationNode;
+import com.sleepycat.je.rep.StateChangeEvent;
+import com.sleepycat.je.rep.StateChangeListener;
+import com.sleepycat.je.rep.TimeConsistencyPolicy;
+import com.sleepycat.je.rep.UnknownMasterException;
+import com.sleepycat.je.rep.util.ReplicationGroupAdmin;
+import com.sleepycat.je.rep.utilint.ServiceDispatcher;
+import com.sleepycat.persist.StoreConfig;
+import org.apache.inlong.tubemq.corebase.TBaseConstants;
+import org.apache.inlong.tubemq.corebase.TokenConstants;
+import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
+import org.apache.inlong.tubemq.corebase.utils.TStringUtils;
+import org.apache.inlong.tubemq.corebase.utils.Tuple2;
+import
org.apache.inlong.tubemq.server.common.fileconfig.MasterReplicationConfig;
+import org.apache.inlong.tubemq.server.master.MasterConfig;
+import org.apache.inlong.tubemq.server.master.bdbstore.MasterGroupStatus;
+import org.apache.inlong.tubemq.server.master.bdbstore.MasterNodeInfo;
+import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.AbsMetaConfigMapperImpl;
+import org.apache.inlong.tubemq.server.master.utils.MetaConfigSamplePrint;
+import org.apache.inlong.tubemq.server.master.web.model.ClusterGroupVO;
+import org.apache.inlong.tubemq.server.master.web.model.ClusterNodeVO;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class BdbMetaConfigMapperImpl extends AbsMetaConfigMapperImpl {
+ private static final int REP_HANDLE_RETRY_MAX = 1;
+ protected static final Logger logger =
+ LoggerFactory.getLogger(BdbMetaConfigMapperImpl.class);
+ private final MetaConfigSamplePrint metaSamplePrint =
+ new MetaConfigSamplePrint(logger);
+ // 0 stopped, 1 starting, 2 started, 3 stopping
+ private final AtomicInteger srvStatus = new AtomicInteger(0);
+
+ // master configure
+ private final MasterConfig masterConfig;
+ // bdb environment configure
+ private final EnvironmentConfig envConfig;
+ // meta data store file
+ private File envHome;
+ // bdb replication configure
+ private final ReplicationConfig repConfig;
+ // bdb replicated environment
+ private ReplicatedEnvironment repEnv;
+ // bdb replication group admin info
+ private final ReplicationGroupAdmin replicationGroupAdmin;
+ // master role flag
+ private volatile boolean isMaster = false;
+ // time since node become active
+ private final AtomicLong masterSinceTime = new AtomicLong(Long.MAX_VALUE);
+ // master node name
+ private String masterNodeName;
+ // node connect failure count
+ private int connectNodeFailCount = 0;
+ // replication nodes
+ private Set<String> replicas4Transfer = new HashSet<>();
+ private final Listener listener = new Listener();
+ private ExecutorService executorService = null;
+ // bdb data store configure
+ private final StoreConfig storeConfig = new StoreConfig();
+
+ public BdbMetaConfigMapperImpl(MasterConfig masterConfig) {
+ super(masterConfig.getRowLockWaitDurMs());
+ this.masterConfig = masterConfig;
+ MasterReplicationConfig replicationConfig =
+ masterConfig.getReplicationConfig();
+ // build replicationGroupAdmin info
+ Set<InetSocketAddress> helpers = new HashSet<>();
+ for (int i = 1; i <= 3; i++) {
+ helpers.add(new InetSocketAddress(this.masterConfig.getHostName(),
+ replicationConfig.getRepNodePort() + i));
+ }
+ this.replicationGroupAdmin =
+ new ReplicationGroupAdmin(replicationConfig.getRepGroupName(),
helpers);
+ // Initialize configuration for BDB-JE replication environment.
+ // Set envHome and generate a ReplicationConfig. Note that
ReplicationConfig and
+ // EnvironmentConfig values could all be specified in the
je.properties file,
+ // as is shown in the properties file included in the example.
+ this.repConfig = new ReplicationConfig();
+ // Set consistency policy for replica.
+ this.repConfig.setConsistencyPolicy(new TimeConsistencyPolicy(3,
+ TimeUnit.SECONDS, 3, TimeUnit.SECONDS));
+ // Wait up to 3 seconds for commitConsumed acknowledgments.
+ this.repConfig.setReplicaAckTimeout(3, TimeUnit.SECONDS);
+ this.repConfig.setConfigParam(ReplicationConfig.TXN_ROLLBACK_LIMIT,
"1000");
+ this.repConfig.setGroupName(replicationConfig.getRepGroupName());
+ this.repConfig.setNodeName(replicationConfig.getRepNodeName());
+ this.repConfig.setNodeHostPort(this.masterConfig.getHostName() +
TokenConstants.ATTR_SEP
+ + replicationConfig.getRepNodePort());
+ if (TStringUtils.isNotEmpty(replicationConfig.getRepHelperHost())) {
+ logger.info("[BDB Impl] ADD HELP HOST");
+
this.repConfig.setHelperHosts(replicationConfig.getRepHelperHost());
+ }
+ // A replicated environment must be opened with transactions enabled.
+ // Environments on a master must be read/write, while environments
+ // on a client can be read/write or read/only. Since the master's
+ // identity may change, it's most convenient to open the environment
in the default
+ // read/write mode. All write operations will be refused on the client
though.
+ this.envConfig = new EnvironmentConfig();
+ this.envConfig.setTransactional(true);
+ this.envConfig.setDurability(new Durability(
+ replicationConfig.getMetaLocalSyncPolicy(),
+ replicationConfig.getMetaReplicaSyncPolicy(),
+ replicationConfig.getRepReplicaAckPolicy()));
+ this.envConfig.setAllowCreate(true);
+ // Set transactional for the replicated environment.
+ this.storeConfig.setTransactional(true);
+ // Set both Master and Replica open the store for write.
+ this.storeConfig.setReadOnly(false);
+ this.storeConfig.setAllowCreate(true);
+ }
+
+ @Override
+ public void start() throws Exception {
+ logger.info("[BDB Impl] Start StoreManagerService, begin");
+ if (!srvStatus.compareAndSet(0, 1)) {
+ logger.info("[BDB Impl] Start StoreManagerService, started");
+ return;
+ }
+ logger.info("[BDB Impl] Starting StoreManagerService...");
+ try {
+ if (executorService != null) {
+ executorService.shutdownNow();
+ executorService = null;
+ }
+ executorService = Executors.newSingleThreadExecutor();
+ // build envHome file
+ envHome = new File(masterConfig.getMetaDataPath());
+ repEnv = getEnvironment();
+ initMetaStore();
+ repEnv.setStateChangeListener(listener);
+ srvStatus.compareAndSet(1, 2);
+ } catch (Throwable ee) {
+ srvStatus.compareAndSet(1, 0);
+ logger.error("[BDB Impl] Start StoreManagerService failure,
error", ee);
+ return;
+ }
+ logger.info("[BDB Impl] Start StoreManagerService, success");
+ }
+
+ @Override
+ public void stop() throws Exception {
+ logger.info("[BDB Impl] Stop StoreManagerService, begin");
+ if (!srvStatus.compareAndSet(2, 3)) {
+ logger.info("[BDB Impl] Stop StoreManagerService, stopped");
+ return;
+ }
+ logger.info("[BDB Impl] Stopping StoreManagerService...");
+ // close bdb configure
+ closeMetaStore();
+ /* evn close */
+ if (repEnv != null) {
+ try {
+ repEnv.close();
+ repEnv = null;
+ } catch (Throwable ee) {
+ logger.error("[BDB Impl] Close repEnv throw error ", ee);
+ }
+ }
+ if (executorService != null) {
+ executorService.shutdownNow();
+ executorService = null;
+ }
+ srvStatus.set(0);
+ logger.info("[BDB Impl] Stop StoreManagerService, success");
+ }
+
+ @Override
+ public boolean checkStoreStatus(boolean checkIsMaster, ProcessResult
result) {
+ if (!isStarted()) {
+ result.setFailResult(DataOpErrCode.DERR_STORE_STOPPED.getCode(),
+ "Meta store service stopped!");
+ return result.isSuccess();
+ }
+ if (checkIsMaster && !isMasterNow()) {
+ result.setFailResult(DataOpErrCode.DERR_STORE_NOT_MASTER.getCode(),
+ "Current node not active, please send your request to the
active Node!");
+ return result.isSuccess();
+ }
+ result.setSuccResult(null);
+ return true;
+ }
+
+ @Override
+ public boolean isMasterNow() {
+ return isMaster;
+ }
+
+ @Override
+ public long getMasterSinceTime() {
+ return this.masterSinceTime.get();
+ }
+
+ @Override
+ public InetSocketAddress getMasterAddress() {
+ ReplicationGroup replicationGroup = getCurrReplicationGroup();
+ if (replicationGroup == null) {
+ logger.info("[BDB Impl] ReplicationGroup is null...please check
the group status!");
+ return null;
+ }
+ for (ReplicationNode node : replicationGroup.getNodes()) {
+ try {
+ NodeState nodeState =
+ replicationGroupAdmin.getNodeState(node, 2000);
+ if (nodeState != null) {
+ if (nodeState.getNodeState().isMaster()) {
+ return node.getSocketAddress();
+ }
+ }
+ } catch (Throwable e) {
+ logger.error("[BDB Impl] Get nodeState Throwable error", e);
+ }
+ }
+ return null;
+ }
+
+ @Override
+ public boolean isPrimaryNodeActive() {
+ if (repEnv == null) {
+ return false;
+ }
+ ReplicationMutableConfig tmpConfig = repEnv.getRepMutableConfig();
+ return tmpConfig != null && tmpConfig.getDesignatedPrimary();
+ }
+
+ @Override
+ public void transferMaster() throws Exception {
+ if (!isStarted()) {
+ throw new Exception("The BDB store StoreService is reboot now!");
+ }
+ if (isMasterNow()) {
+ if (!isPrimaryNodeActive()) {
+ if ((replicas4Transfer != null) &&
(!replicas4Transfer.isEmpty())) {
+ logger.info(new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
+ .append("[BDB Impl] start transferMaster to
replicas: ")
+ .append(replicas4Transfer).toString());
+ repEnv.transferMaster(replicas4Transfer, 5,
TimeUnit.MINUTES);
+ logger.info("[BDB Impl] transferMaster end...");
+ } else {
+ throw new Exception("The replicate nodes is empty!");
+ }
+ } else {
+ throw new Exception("DesignatedPrimary happened...please check
if the other member is down!");
+ }
+ } else {
+ throw new Exception("Please send your request to the master
Node!");
+ }
+ }
+
+ @Override
+ public ClusterGroupVO getGroupAddressStrInfo() {
+ ClusterGroupVO clusterGroupVO = new ClusterGroupVO();
+ clusterGroupVO.setGroupStatus("Abnormal");
+ clusterGroupVO.setGroupName(replicationGroupAdmin.getGroupName());
+ // query current replication group info
+ ReplicationGroup replicationGroup = getCurrReplicationGroup();
+ if (replicationGroup == null) {
+ return clusterGroupVO;
+ }
+ // translate replication group info to ClusterGroupVO structure
+ Tuple2<Boolean, List<ClusterNodeVO>> transResult =
+ transReplicateNodes(replicationGroup);
+ clusterGroupVO.setNodeData(transResult.getF1());
+ clusterGroupVO.setPrimaryNodeActive(isPrimaryNodeActive());
+ if (transResult.getF0()) {
+ if (isPrimaryNodeActive()) {
+ clusterGroupVO.setGroupStatus("Running-ReadOnly");
+ } else {
+ clusterGroupVO.setGroupStatus("Running-ReadWrite");
+ }
+ }
+ return clusterGroupVO;
+ }
+
+ @Override
+ public MasterGroupStatus getMasterGroupStatus(boolean isFromHeartbeat) {
+ // #lizard forgives
+ if (repEnv == null) {
+ return null;
+ }
+ ReplicationGroup replicationGroup = null;
+ try {
+ replicationGroup = repEnv.getGroup();
+ } catch (DatabaseException e) {
+ if (e instanceof EnvironmentFailureException) {
+ if (isFromHeartbeat) {
+ logger.error("[BDB Error] Check found
EnvironmentFailureException", e);
+ try {
+ stop();
+ start();
+ replicationGroup = repEnv.getGroup();
+ } catch (Throwable e1) {
+ logger.error("[BDB Error] close and reopen
storeManager error", e1);
+ }
+ } else {
+ logger.error(
+ "[BDB Error] Get EnvironmentFailureException error
while non heartBeat request", e);
+ }
+ } else {
+ logger.error("[BDB Error] Get replication group info error",
e);
+ }
+ } catch (Throwable ee) {
+ logger.error("[BDB Error] Get replication group throw error", ee);
+ }
+ if (replicationGroup == null) {
+ logger.error(
+ "[BDB Error] ReplicationGroup is null...please check the
status of the group!");
+ return null;
+ }
+ int activeNodes = 0;
+ boolean isMasterActive = false;
+ Set<String> tmp = new HashSet<>();
+ for (ReplicationNode node : replicationGroup.getNodes()) {
+ MasterNodeInfo masterNodeInfo =
+ new MasterNodeInfo(replicationGroup.getName(),
+ node.getName(), node.getHostName(),
node.getPort());
+ try {
+ NodeState nodeState = replicationGroupAdmin.getNodeState(node,
2000);
+ if (nodeState != null) {
+ if (nodeState.getNodeState().isActive()) {
+ activeNodes++;
+ if (nodeState.getNodeName().equals(masterNodeName)) {
+ isMasterActive = true;
+ masterNodeInfo.setNodeStatus(1);
+ }
+ }
+ if (nodeState.getNodeState().isReplica()) {
+ tmp.add(nodeState.getNodeName());
+ replicas4Transfer = tmp;
+ masterNodeInfo.setNodeStatus(0);
+ }
+ }
+ } catch (IOException e) {
+ connectNodeFailCount++;
+ masterNodeInfo.setNodeStatus(-1);
+ metaSamplePrint.printExceptionCaught(e, node.getHostName(),
node.getName());
+ continue;
+ } catch (ServiceDispatcher.ServiceConnectFailedException e) {
+ masterNodeInfo.setNodeStatus(-2);
+ metaSamplePrint.printExceptionCaught(e, node.getHostName(),
node.getName());
+ continue;
+ } catch (Throwable ee) {
+ masterNodeInfo.setNodeStatus(-3);
+ metaSamplePrint.printExceptionCaught(ee, node.getHostName(),
node.getName());
+ continue;
+ }
+ }
+ MasterGroupStatus masterGroupStatus = new
MasterGroupStatus(isMasterActive);
+ int groupSize = replicationGroup.getElectableNodes().size();
+ int majoritySize = groupSize / 2 + 1;
+ if ((activeNodes >= majoritySize) && isMasterActive) {
+ masterGroupStatus.setMasterGroupStatus(true, true, true);
+ connectNodeFailCount = 0;
+ if (isPrimaryNodeActive()) {
+
repEnv.setRepMutableConfig(repEnv.getRepMutableConfig().setDesignatedPrimary(false));
+ }
+ }
+ if (groupSize == 2 && connectNodeFailCount >= 3) {
+ masterGroupStatus.setMasterGroupStatus(true, false, true);
+ if (connectNodeFailCount > 1000) {
+ connectNodeFailCount = 3;
+ }
+ if (!isPrimaryNodeActive()) {
+ logger.error("[BDB Error] DesignatedPrimary happened...please
check if the other member is down");
+
repEnv.setRepMutableConfig(repEnv.getRepMutableConfig().setDesignatedPrimary(true));
+ }
+ }
+ return masterGroupStatus;
+ }
+
+ /**
+ * State Change Listener,
+ * through this object, it complete the metadata cache cleaning
+ * and loading of the latest data.
+ *
+ * */
+ public class Listener implements StateChangeListener {
+ @Override
+ public void stateChange(StateChangeEvent stateChangeEvent) throws
RuntimeException {
+ if (repConfig != null) {
+ logger.warn(new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
+ .append("[BDB Impl][").append(repConfig.getGroupName())
+ .append("Receive a group status changed event]...
stateChangeEventTime: ")
+ .append(stateChangeEvent.getEventTime()).toString());
+ }
+ doWork(stateChangeEvent);
+ }
+
+ /**
+ * process replicate nodes status event
+ *
+ * @param stateChangeEvent status change event
+ */
+ public void doWork(final StateChangeEvent stateChangeEvent) {
+
+ final String currentNode = new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
+ .append("GroupName:").append(repConfig.getGroupName())
+ .append(",nodeName:").append(repConfig.getNodeName())
+
.append(",hostName:").append(repConfig.getNodeHostPort()).toString();
+ if (executorService == null) {
+ logger.error("[BDB Impl] found executorService is null while
doWork!");
+ return;
+ }
+ executorService.submit(() -> {
+ StringBuilder sBuilder =
+ new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE);
+ switch (stateChangeEvent.getState()) {
+ case MASTER:
+ if (!isMaster) {
+ try {
+ reloadMetaStore(sBuilder);
+ isMaster = true;
+
masterSinceTime.set(System.currentTimeMillis());
+ masterNodeName =
stateChangeEvent.getMasterNodeName();
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode).append(" is a
master.").toString());
+ } catch (Throwable e) {
+ isMaster = false;
+ logger.error("[BDB Impl] fatal error when
Reloading Info ", e);
+ }
+ }
+ break;
+ case REPLICA:
+ isMaster = false;
+ masterNodeName = stateChangeEvent.getMasterNodeName();
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode).append(" is a
slave.").toString());
+ break;
+ default:
+ isMaster = false;
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode).append(" is Unknown state
")
+
.append(stateChangeEvent.getState().name()).toString());
+ break;
+ }
+ sBuilder.delete(0, sBuilder.length());
+ });
+ }
+ }
+
+ private boolean isStarted() {
+ return (this.srvStatus.get() == 2);
+ }
+
+ /**
+ * Initial meta-data stores.
+ *
+ */
+ protected void initMetaStore() {
+ clusterConfigMapper = new BdbClusterConfigMapperImpl(repEnv,
storeConfig);
+ brokerConfigMapper = new BdbBrokerConfigMapperImpl(repEnv,
storeConfig);
+ topicDeployMapper = new BdbTopicDeployMapperImpl(repEnv, storeConfig);
+ groupResCtrlMapper = new BdbGroupResCtrlMapperImpl(repEnv,
storeConfig);
+ topicCtrlMapper = new BdbTopicCtrlMapperImpl(repEnv, storeConfig);
+ consumeCtrlMapper = new BdbConsumeCtrlMapperImpl(repEnv, storeConfig);
+ }
+
+ /**
+ * Creates the replicated environment handle and returns it. It will retry
indefinitely if a
+ * master could not be established because a sufficient number of nodes
were not available, or
+ * there were networking issues, etc.
+ *
+ * @return the newly created replicated environment handle
+ * @throws InterruptedException if the operation was interrupted
+ */
+ private ReplicatedEnvironment getEnvironment() throws InterruptedException
{
+ DatabaseException exception = null;
+ //In this example we retry REP_HANDLE_RETRY_MAX times, but a
production HA application may
+ //retry indefinitely.
+ for (int i = 0; i < REP_HANDLE_RETRY_MAX; i++) {
+ try {
+ return new ReplicatedEnvironment(envHome, repConfig,
envConfig);
+ } catch (UnknownMasterException unknownMaster) {
+ exception = unknownMaster;
+ // Indicates there is a group level problem: insufficient
nodes for an election,
+ // network connectivity issues, etc. Wait and retry to allow
the problem
+ // to be resolved.
+ logger.error(new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
+ .append("[BDB Impl] master could not be established. ")
+ .append("Exception
message:").append(unknownMaster.getMessage())
+ .append(" Will retry after 5 seconds.").toString());
+ Thread.sleep(5 * 1000);
+ continue;
+ } catch (InsufficientLogException insufficientLogEx) {
+ logger.error(new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
+ .append("[BDB Impl] [Restoring data please wait....] ")
+ .append("Obtains logger files for a Replica from other
members of ")
+ .append("the replication group. A Replica may need to
do so if it ")
+ .append("has been offline for some time, and has
fallen behind in ")
+ .append("its execution of the replication
stream.").toString());
+ NetworkRestore restore = new NetworkRestore();
+ NetworkRestoreConfig config = new NetworkRestoreConfig();
+ // delete obsolete logger files.
+ config.setRetainLogFiles(false);
+ restore.execute(insufficientLogEx, config);
+ // retry
+ return new ReplicatedEnvironment(envHome, repConfig,
envConfig);
+ }
+ }
+ // Failed despite retries.
+ throw exception;
+ }
+
+ private ReplicationGroup getCurrReplicationGroup() {
+ ReplicationGroup replicationGroup;
+ try {
+ replicationGroup = repEnv.getGroup();
+ } catch (Throwable e) {
+ logger.error("[BDB Impl] get current master group info error", e);
+ return null;
+ }
+ return replicationGroup;
+ }
+
+ /**
+ * Query replication group nodes status and translate to ClusterNodeVO type
+ *
+ * @param replicationGroup the replication group
+ * @return if has master, replication nodes info
+ */
+ private Tuple2<Boolean, List<ClusterNodeVO>> transReplicateNodes(
+ ReplicationGroup replicationGroup) {
+ boolean hasMaster = false;
+ List<ClusterNodeVO> clusterNodeVOList = new ArrayList<>();
+ for (ReplicationNode node : replicationGroup.getNodes()) {
+ ClusterNodeVO clusterNodeVO = new ClusterNodeVO();
+ clusterNodeVO.setHostName(node.getHostName());
+ clusterNodeVO.setNodeName(node.getName());
+ clusterNodeVO.setPort(node.getPort());
+ try {
+ NodeState nodeState =
+ replicationGroupAdmin.getNodeState(node, 2000);
+ if (nodeState != null) {
+ if (nodeState.getNodeState() ==
ReplicatedEnvironment.State.MASTER) {
+ hasMaster = true;
+ }
+
clusterNodeVO.setNodeStatus(nodeState.getNodeState().toString());
+ clusterNodeVO.setJoinTime(nodeState.getJoinTime());
+ } else {
+ clusterNodeVO.setNodeStatus("Not-found");
+ clusterNodeVO.setJoinTime(0);
+ }
+ } catch (IOException e) {
+ clusterNodeVO.setNodeStatus("Error");
+ clusterNodeVO.setJoinTime(0);
+ } catch (ServiceDispatcher.ServiceConnectFailedException e) {
+ clusterNodeVO.setNodeStatus("Unconnected");
+ clusterNodeVO.setJoinTime(0);
+ }
+ clusterNodeVOList.add(clusterNodeVO);
+ }
+ return new Tuple2<>(hasMaster, clusterNodeVOList);
+ }
+}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
index aa4e13d..975b7c4 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
@@ -62,7 +62,7 @@ import org.apache.inlong.tubemq.server.master.MasterConfig;
import org.apache.inlong.tubemq.server.master.bdbstore.MasterGroupStatus;
import org.apache.inlong.tubemq.server.master.bdbstore.MasterNodeInfo;
import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
-import
org.apache.inlong.tubemq.server.master.metamanage.keepalive.AliveObserver;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaStoreService;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.ClusterSettingEntity;
@@ -76,7 +76,7 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.Co
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.GroupResCtrlMapper;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.TopicCtrlMapper;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.TopicDeployMapper;
-import org.apache.inlong.tubemq.server.master.utils.BdbStoreSamplePrint;
+import org.apache.inlong.tubemq.server.master.utils.MetaConfigSamplePrint;
import org.apache.inlong.tubemq.server.master.web.model.ClusterGroupVO;
import org.apache.inlong.tubemq.server.master.web.model.ClusterNodeVO;
import org.slf4j.Logger;
@@ -88,8 +88,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
private static final Logger logger =
LoggerFactory.getLogger(BdbMetaStoreServiceImpl.class);
- private final BdbStoreSamplePrint bdbStoreSamplePrint =
- new BdbStoreSamplePrint(logger);
+ private final MetaConfigSamplePrint metaConfigSamplePrint =
+ new MetaConfigSamplePrint(logger);
// parameters need input
// local host name
private final String nodeHost;
@@ -99,7 +99,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
private final MasterReplicationConfig replicationConfig;
private final Listener listener = new Listener();
private ExecutorService executorService = null;
- private final List<AliveObserver> eventObservers = new ArrayList<>();
+ private final List<MetaConfigObserver> eventObservers = new ArrayList<>();
// service status
// 0 stopped, 1 starting, 2 started, 3 stopping
private final AtomicInteger srvStatus = new AtomicInteger(0);
@@ -839,7 +839,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
}
@Override
- public void registerObserver(AliveObserver eventObserver) {
+ public void registerObserver(MetaConfigObserver eventObserver) {
if (eventObserver != null) {
eventObservers.add(eventObserver);
}
@@ -1028,15 +1028,15 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
} catch (IOException e) {
connectNodeFailCount++;
masterNodeInfo.setNodeStatus(-1);
- bdbStoreSamplePrint.printExceptionCaught(e,
node.getHostName(), node.getName());
+ metaConfigSamplePrint.printExceptionCaught(e,
node.getHostName(), node.getName());
continue;
} catch (ServiceDispatcher.ServiceConnectFailedException e) {
masterNodeInfo.setNodeStatus(-2);
- bdbStoreSamplePrint.printExceptionCaught(e,
node.getHostName(), node.getName());
+ metaConfigSamplePrint.printExceptionCaught(e,
node.getHostName(), node.getName());
continue;
} catch (Throwable ee) {
masterNodeInfo.setNodeStatus(-3);
- bdbStoreSamplePrint.printExceptionCaught(ee,
node.getHostName(), node.getName());
+ metaConfigSamplePrint.printExceptionCaught(ee,
node.getHostName(), node.getName());
continue;
}
}
@@ -1143,7 +1143,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
*
* */
private void clearCachedRunData() {
- for (AliveObserver observer : eventObservers) {
+ for (MetaConfigObserver observer : eventObservers) {
observer.clearCacheData();
}
}
@@ -1153,7 +1153,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
*
* */
private void reloadRunData() {
- for (AliveObserver observer : eventObservers) {
+ for (MetaConfigObserver observer : eventObservers) {
observer.reloadCacheData();
}
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
index 186920e..3a4918f 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
@@ -43,7 +43,7 @@ import
org.apache.inlong.tubemq.server.common.utils.SerialIdUtils;
import org.apache.inlong.tubemq.server.master.MasterConfig;
import org.apache.inlong.tubemq.server.master.TMaster;
import org.apache.inlong.tubemq.server.master.metamanage.MetaDataManager;
-import
org.apache.inlong.tubemq.server.master.metamanage.keepalive.AliveObserver;
+import
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import org.apache.inlong.tubemq.server.master.stats.MasterSrvStatsHolder;
import org.slf4j.Logger;
@@ -52,7 +52,7 @@ import org.slf4j.LoggerFactory;
/*
* Broker run manager
*/
-public class DefBrokerRunManager implements BrokerRunManager, AliveObserver {
+public class DefBrokerRunManager implements BrokerRunManager,
MetaConfigObserver {
private static final Logger logger =
LoggerFactory.getLogger(DefBrokerRunManager.class);
// meta data manager
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/BdbStoreSamplePrint.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/MetaConfigSamplePrint.java
similarity index 81%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/BdbStoreSamplePrint.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/MetaConfigSamplePrint.java
index 94e1adf..c277f2c 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/BdbStoreSamplePrint.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/utils/MetaConfigSamplePrint.java
@@ -1,107 +1,107 @@
-/**
- * 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.inlong.tubemq.server.master.utils;
-
-import com.sleepycat.je.rep.utilint.ServiceDispatcher;
-import java.io.IOException;
-import org.apache.inlong.tubemq.corebase.utils.AbstractSamplePrint;
-import org.slf4j.Logger;
-
-public class BdbStoreSamplePrint extends AbstractSamplePrint {
- /**
- * Log limit class
- */
- private final Logger logger;
-
- public BdbStoreSamplePrint(final Logger logger) {
- super();
- this.logger = logger;
- }
-
- public BdbStoreSamplePrint(final Logger logger,
- long sampleDetailDur, long sampleResetDur,
- long maxDetailCount, long maxTotalCount) {
- super(sampleDetailDur, sampleResetDur, maxDetailCount, maxTotalCount);
- this.logger = logger;
- }
-
- @Override
- public void printExceptionCaught(Throwable e) {
- //
- }
-
- @Override
- public void printExceptionCaught(Throwable e, String hostName, String
nodeName) {
- if (e != null) {
- if ((e instanceof IOException)) {
- // IOException log limit
- final long now = System.currentTimeMillis();
- final long diffTime = now - lastLogTime.get();
- final long curPrintCnt = totalPrintCount.incrementAndGet();
- if (curPrintCnt < maxTotalCount) {
- if (diffTime < sampleDetailDur && curPrintCnt <
maxDetailCount) {
- logger.error(sBuilder.append("[BDB Error] Connect to
node:[")
- .append(hostName).append(",").append(nodeName)
- .append("] IOException error").toString(), e);
- } else {
- logger.error(sBuilder.append("[BDB Error] Connect to
node:[")
- .append(hostName).append(",").append(nodeName)
- .append("] IOException error is
").append(e.toString()).toString());
- }
- sBuilder.delete(0, sBuilder.length());
- }
- if (diffTime > sampleResetDur) {
- if (this.lastLogTime.compareAndSet(now - diffTime, now)) {
- totalPrintCount.set(0);
- }
- }
- } else {
- if (e instanceof
ServiceDispatcher.ServiceConnectFailedException) {
- // Print some part log
- final long curPrintCnt =
totalUncheckCount.incrementAndGet();
- if (curPrintCnt < maxUncheckDetailCount) {
- logger.error(sBuilder.append("[BDB Error] Connect to
node:[")
- .append(hostName).append(",").append(nodeName)
- .append("] ServiceConnectFailedException
error").toString(), e);
- } else {
- logger.error(sBuilder.append("[BDB Error] Connect to
node:[")
- .append(hostName).append(",").append(nodeName)
- .append("] ServiceConnectFailedException error
is ")
- .append(e.toString()).toString());
- }
- sBuilder.delete(0, sBuilder.length());
- } else {
- logger.error(sBuilder.append("[BDB Error] Connect to
node:[")
- .append(hostName).append(",").append(nodeName)
- .append("] throw error").toString(), e);
- sBuilder.delete(0, sBuilder.length());
- }
- }
- }
- }
-
- @Override
- public void printWarn(String err) {
- //
- }
-
- @Override
- public void printError(String err) {
- //
- }
-}
+/**
+ * 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.inlong.tubemq.server.master.utils;
+
+import com.sleepycat.je.rep.utilint.ServiceDispatcher;
+import java.io.IOException;
+import org.apache.inlong.tubemq.corebase.utils.AbstractSamplePrint;
+import org.slf4j.Logger;
+
+public class MetaConfigSamplePrint extends AbstractSamplePrint {
+ /**
+ * Log limit class
+ */
+ private final Logger logger;
+
+ public MetaConfigSamplePrint(final Logger logger) {
+ super();
+ this.logger = logger;
+ }
+
+ public MetaConfigSamplePrint(final Logger logger,
+ long sampleDetailDur, long sampleResetDur,
+ long maxDetailCount, long maxTotalCount) {
+ super(sampleDetailDur, sampleResetDur, maxDetailCount, maxTotalCount);
+ this.logger = logger;
+ }
+
+ @Override
+ public void printExceptionCaught(Throwable e) {
+ //
+ }
+
+ @Override
+ public void printExceptionCaught(Throwable e, String hostName, String
nodeName) {
+ if (e != null) {
+ if ((e instanceof IOException)) {
+ // IOException log limit
+ final long now = System.currentTimeMillis();
+ final long diffTime = now - lastLogTime.get();
+ final long curPrintCnt = totalPrintCount.incrementAndGet();
+ if (curPrintCnt < maxTotalCount) {
+ if (diffTime < sampleDetailDur && curPrintCnt <
maxDetailCount) {
+ logger.error(sBuilder.append("[MetaConfig Error]
Connect to node:[")
+ .append(hostName).append(",").append(nodeName)
+ .append("] IOException error").toString(), e);
+ } else {
+ logger.error(sBuilder.append("[MetaConfig Error]
Connect to node:[")
+ .append(hostName).append(",").append(nodeName)
+ .append("] IOException error is
").append(e.toString()).toString());
+ }
+ sBuilder.delete(0, sBuilder.length());
+ }
+ if (diffTime > sampleResetDur) {
+ if (this.lastLogTime.compareAndSet(now - diffTime, now)) {
+ totalPrintCount.set(0);
+ }
+ }
+ } else {
+ if (e instanceof
ServiceDispatcher.ServiceConnectFailedException) {
+ // Print some part log
+ final long curPrintCnt =
totalUncheckCount.incrementAndGet();
+ if (curPrintCnt < maxUncheckDetailCount) {
+ logger.error(sBuilder.append("[MetaConfig Error]
Connect to node:[")
+ .append(hostName).append(",").append(nodeName)
+ .append("] ServiceConnectFailedException
error").toString(), e);
+ } else {
+ logger.error(sBuilder.append("[MetaConfig Error]
Connect to node:[")
+ .append(hostName).append(",").append(nodeName)
+ .append("] ServiceConnectFailedException error
is ")
+ .append(e.toString()).toString());
+ }
+ sBuilder.delete(0, sBuilder.length());
+ } else {
+ logger.error(sBuilder.append("[MetaConfig Error] Connect
to node:[")
+ .append(hostName).append(",").append(nodeName)
+ .append("] throw error").toString(), e);
+ sBuilder.delete(0, sBuilder.length());
+ }
+ }
+ }
+ }
+
+ @Override
+ public void printWarn(String err) {
+ //
+ }
+
+ @Override
+ public void printError(String err) {
+ //
+ }
+}