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) {
+        //
+    }
+}

Reply via email to