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 ffa2628  [INLONG-3035][TubeMQ] Optimize the MetaStoreService 
implementation class (#3142)
ffa2628 is described below

commit ffa262844373874858352ed773def71de758ab16
Author: gosonzhang <[email protected]>
AuthorDate: Tue Mar 15 16:55:02 2022 +0800

    [INLONG-3035][TubeMQ] Optimize the MetaStoreService implementation class 
(#3142)
---
 .../tubemq/server/common/fileconfig/ZKConfig.java  |  10 +
 .../tubemq/server/common/zookeeper/ZKUtil.java     | 204 +++++++++++++
 .../server/master/metamanage/MetaDataManager.java  |   3 +-
 .../metamanage/metastore/KeepAliveService.java     |   4 +-
 .../metastore/impl/AbsMetaConfigMapperImpl.java    |  49 ++-
 .../impl/bdbimpl/BdbMetaConfigMapperImpl.java      |  85 ++----
 .../impl/bdbimpl/BdbMetaStoreServiceImpl.java      |   4 +-
 .../impl/zkimpl/ZKMetaConfigMapperImpl.java        | 332 +++++++++++++++++++++
 .../server/master/web/MasterStatusCheckFilter.java |  14 +-
 .../server/master/web/action/screen/Tubeweb.java   |   7 +-
 10 files changed, 632 insertions(+), 80 deletions(-)

diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/fileconfig/ZKConfig.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/fileconfig/ZKConfig.java
index 42385cd..06f58c9 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/fileconfig/ZKConfig.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/fileconfig/ZKConfig.java
@@ -28,6 +28,7 @@ public class ZKConfig {
     private int zkSyncTimeMs = 1000;
     private long zkCommitPeriodMs = 5000L;
     private int zkCommitFailRetries = 
TServerConstants.CFG_ZK_COMMIT_DEFAULT_RETRIES;
+    private long zkMasterCheckPeriodMs = 5000L;
 
     public ZKConfig() {
 
@@ -89,6 +90,14 @@ public class ZKConfig {
         this.zkCommitPeriodMs = zkCommitPeriodMs;
     }
 
+    public long getZkMasterCheckPeriodMs() {
+        return zkMasterCheckPeriodMs;
+    }
+
+    public void setZkMasterCheckPeriodMs(long zkMasterCheckPeriodMs) {
+        this.zkMasterCheckPeriodMs = zkMasterCheckPeriodMs;
+    }
+
     @Override
     public String toString() {
         return new StringBuilder(512)
@@ -99,6 +108,7 @@ public class ZKConfig {
                 .append(",\"zkSyncTimeMs\":").append(zkSyncTimeMs)
                 .append(",\"zkCommitPeriodMs\":").append(zkCommitPeriodMs)
                 
.append(",\"zkCommitFailRetries\":").append(zkCommitFailRetries)
+                
.append(",\"zkMasterCheckPeriodMs\":").append(zkMasterCheckPeriodMs)
                 .append("}").toString();
     }
 }
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/zookeeper/ZKUtil.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/zookeeper/ZKUtil.java
index 2d48c28..e12c1db 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/zookeeper/ZKUtil.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/zookeeper/ZKUtil.java
@@ -527,4 +527,208 @@ public class ZKUtil {
         }
     }
 
+    /**
+     * Simple class to hold a node path and node data.
+     */
+    public static class NodeAndData {
+        private final String node;
+        private final byte[] data;
+
+        public NodeAndData(String node, byte[] data) {
+            this.node = node;
+            this.data = data;
+        }
+
+        public String getNode() {
+            return node;
+        }
+
+        public byte[] getData() {
+            return data;
+        }
+
+        public boolean isEmpty() {
+            return (data.length == 0);
+        }
+    }
+
+    /**
+     * Returns the date of child zNodes of the specified zNode. Also sets a 
watch
+     * on the specified zNode which will capture a NodeDeleted event on the
+     * specified zNode as well as NodeChildrenChanged if any children of the
+     * specified zNode are created or deleted.
+     *
+     * Returns null if the specified node does not exist. Otherwise returns a 
list
+     * of children of the specified node. If the node exists but it has no
+     * children, an empty list will be returned.
+     *
+     * @param zkw
+     *          zk reference
+     * @param baseNode
+     *          path of node to list and watch children of
+     * @return list of data of children of the specified node, an empty list if
+     *         the node exists but has no children, and null if the node does 
not
+     *         exist
+     * @throws KeeperException
+     *           if unexpected zookeeper exception
+     */
+    public static List<NodeAndData> getChildDataAndWatchForNewChildren(
+            ZooKeeperWatcher zkw, String baseNode) throws KeeperException {
+        List<String> nodes = listChildrenAndWatchForNewChildren(zkw, baseNode);
+        List<NodeAndData> newNodes = new ArrayList<NodeAndData>();
+        if (nodes != null) {
+            for (String node : nodes) {
+                String nodePath = joinZNode(baseNode, node);
+                byte[] data = getDataAndWatch(zkw, nodePath);
+                if (data != null) {
+                    newNodes.add(new NodeAndData(nodePath, data));
+                } else {
+                    logger.error("Get data is null for nodePath " + nodePath);
+                }
+            }
+        }
+        return newNodes;
+    }
+
+    /**
+     * Lists the children zNodes of the specified zNode. Also sets a watch on 
the
+     * specified zNode which will capture a NodeDeleted event on the specified
+     * zNode as well as NodeChildrenChanged if any children of the specified 
zNode
+     * are created or deleted.
+     *
+     * Returns null if the specified node does not exist. Otherwise returns a 
list
+     * of children of the specified node. If the node exists but it has no
+     * children, an empty list will be returned.
+     *
+     * @param zkw zk reference
+     * @param zNode
+     *          path of node to list and watch children of
+     * @return list of children of the specified node, an empty list if the 
node
+     *         exists but has no children, and null if the node does not exist
+     * @throws KeeperException
+     *           if unexpected zookeeper exception
+     */
+    public static List<String> listChildrenAndWatchForNewChildren(
+            ZooKeeperWatcher zkw, String zNode) throws KeeperException {
+        try {
+            List<String> children = 
zkw.getRecoverableZooKeeper().getChildren(zNode,
+                    zkw);
+            return children;
+        } catch (KeeperException.NoNodeException ke) {
+            if (logger.isDebugEnabled()) {
+                logger.debug(zkw.prefix("Unable to list children of zNode " + 
zNode + " "
+                        + "because node does not exist (not an error)"));
+            }
+            return null;
+        } catch (KeeperException e) {
+            logger.warn(zkw.prefix("Unable to list children of zNode " + zNode 
+ " "), e);
+            zkw.keeperException(e);
+            return null;
+        } catch (InterruptedException e) {
+            logger.warn(zkw.prefix("Unable to list children of zNode " + zNode 
+ " "), e);
+            zkw.interruptedException(e);
+            return null;
+        }
+    }
+
+    /**
+     *
+     * Set the specified zNode to be an ephemeral node carrying the specified
+     * data.
+     *
+     * If the node is created successfully, a watcher is also set on the node.
+     *
+     * If the node is not created successfully because it already exists, this
+     * method will also set a watcher on the node.
+     *
+     * If there is another problem, a KeeperException will be thrown.
+     *
+     * @param zkw
+     *          zk reference
+     * @param zNode
+     *          path of node
+     * @param data
+     *          data of node
+     * @return true if node created, false if not, watch set in both cases
+     * @throws KeeperException
+     *           if unexpected zookeeper exception
+     */
+    public static boolean createEphemeralNodeAndWatch(ZooKeeperWatcher zkw,
+                                                      String zNode, byte[] 
data) throws KeeperException {
+        return createEphemeralNodeAndWatch(zkw, zNode, data, 
CreateMode.EPHEMERAL);
+    }
+
+    /**
+     * Set the specified zNode to be an ephemeral node carrying the specified
+     * data.
+     *
+     * If the node is created successfully, a watcher is also set on the node.
+     *
+     * If the node is not created successfully because it already exists, this
+     * method will also set a watcher on the node.
+     *
+     * If there is another problem, a KeeperException will be thrown.
+     *
+     * @param zkw     zk reference
+     * @param zNode   path of node to watch
+     * @param data    write data
+     * @param mode    create mode
+     * @return  whether success
+     * @throws KeeperException process exception
+     */
+    public static boolean createEphemeralNodeAndWatch(ZooKeeperWatcher zkw, 
String zNode,
+                                                      byte[] data, CreateMode 
mode) throws KeeperException {
+        try {
+            waitForZKConnectionIfAuthenticating(zkw);
+            zkw.getRecoverableZooKeeper().create(zNode, data, createACL(zkw, 
zNode),
+                    mode);
+        } catch (KeeperException.NodeExistsException nee) {
+            if (!watchAndCheckExists(zkw, zNode)) {
+                // It did exist but now it doesn't, try again
+                return createEphemeralNodeAndWatch(zkw, zNode, data, mode);
+            }
+            return false;
+        } catch (InterruptedException e) {
+            logger.info("Interrupted", e);
+            Thread.currentThread().interrupt();
+        }
+        return true;
+    }
+
+    /**
+     * Watch the specified zNode for delete/create/change events. The watcher 
is
+     * set whether or not the node exists. If the node already exists, the 
method
+     * returns true. If the node does not exist, the method returns false.
+     *
+     * @param zkw
+     *          zk reference
+     * @param zNode
+     *          path of node to watch
+     * @return true if zNode exists, false if does not exist or error
+     * @throws KeeperException
+     *           if unexpected zookeeper exception
+     */
+    public static boolean watchAndCheckExists(ZooKeeperWatcher zkw, String 
zNode)
+            throws KeeperException {
+        try {
+            Stat s = zkw.getRecoverableZooKeeper().exists(zNode, zkw);
+            boolean exists = (s != null);
+            if (logger.isDebugEnabled()) {
+                if (exists) {
+                    logger.debug(zkw.prefix("Set watcher on existing zNode " + 
zNode));
+                } else {
+                    logger.debug(zkw.prefix(zNode + " does not exist. Watcher 
is set."));
+                }
+            }
+            return exists;
+        } catch (KeeperException e) {
+            logger.warn(zkw.prefix("Unable to set watcher on zNode " + zNode), 
e);
+            zkw.keeperException(e);
+            return false;
+        } catch (InterruptedException e) {
+            logger.warn(zkw.prefix("Unable to set watcher on zNode " + zNode), 
e);
+            zkw.interruptedException(e);
+            return false;
+        }
+    }
 }
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 75f17ab..9351bc4 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
@@ -17,7 +17,6 @@
 
 package org.apache.inlong.tubemq.server.master.metamanage;
 
-import java.net.InetSocketAddress;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
@@ -173,7 +172,7 @@ public class MetaDataManager implements Server {
         }
     }
 
-    public InetSocketAddress getMasterAddress() {
+    public String getMasterAddress() {
         return metaStoreService.getMasterAddress();
     }
 
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
index ecce26e..9d9fd8c 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/KeepAliveService.java
@@ -17,8 +17,6 @@
 
 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;
@@ -44,7 +42,7 @@ public interface KeepAliveService extends Server {
      *
      * @return  the current Master address
      */
-    InetSocketAddress getMasterAddress();
+    String getMasterAddress();
 
     /**
      * Whether the Master node in active
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
index d5ce9e7..07fb9bb 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsMetaConfigMapperImpl.java
@@ -23,6 +23,9 @@ import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+
 import org.apache.commons.codec.binary.StringUtils;
 import org.apache.inlong.tubemq.corebase.TBaseConstants;
 import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
@@ -30,6 +33,7 @@ import 
org.apache.inlong.tubemq.server.common.fielddef.WebFieldDef;
 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.MasterConfig;
 import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
 import 
org.apache.inlong.tubemq.server.master.metamanage.metastore.MetaConfigObserver;
 import 
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.MetaConfigMapper;
@@ -53,6 +57,14 @@ import org.slf4j.LoggerFactory;
 public abstract class AbsMetaConfigMapperImpl implements MetaConfigMapper {
     protected static final Logger logger =
             LoggerFactory.getLogger(AbsMetaConfigMapperImpl.class);
+    // master configure
+    protected final MasterConfig masterConfig;
+    // 0 stopped, 1 starting, 2 started, 3 stopping
+    protected final AtomicInteger srvStatus = new AtomicInteger(0);
+    // master role flag
+    protected volatile boolean isMaster = false;
+    // time since node become active
+    protected final AtomicLong masterSinceTime = new 
AtomicLong(Long.MAX_VALUE);
     // row lock.
     private final RowLock metaRowLock;
     // default cluster setting
@@ -73,9 +85,10 @@ public abstract class AbsMetaConfigMapperImpl implements 
MetaConfigMapper {
     // the observers focusing on active-standby switching
     private final List<MetaConfigObserver> eventObservers = new ArrayList<>();
 
-    public AbsMetaConfigMapperImpl(int rowLockWaiDurMs) {
+    public AbsMetaConfigMapperImpl(MasterConfig masterConfig) {
+        this.masterConfig = masterConfig;
         this.metaRowLock =
-                new RowLock("MetaData-RowLock", rowLockWaiDurMs);
+                new RowLock("MetaData-RowLock", 
masterConfig.getRowLockWaitDurMs());
     }
 
     @Override
@@ -86,6 +99,22 @@ public abstract class AbsMetaConfigMapperImpl implements 
MetaConfigMapper {
     }
 
     @Override
+    public boolean checkStoreStatus(boolean checkIsMaster, ProcessResult 
result) {
+        if (!isServiceStarted()) {
+            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 addOrUpdClusterDefSetting(BaseEntity opEntity,
                                              int brokerPort, int brokerTlsPort,
                                              int brokerWebPort, int 
maxMsgSizeMB,
@@ -1159,6 +1188,22 @@ public abstract class AbsMetaConfigMapperImpl implements 
MetaConfigMapper {
     }
 
     /**
+     * Whether service started
+     *
+     * @return true for started, false for other cases
+     */
+    protected boolean isServiceStarted() {
+        return (this.srvStatus.get() == 2);
+    }
+
+    /**
+     * Initial meta-data stores.
+     *
+     * @param strBuff  the string buffer
+     */
+    protected abstract void initMetaStore(StringBuilder strBuff);
+
+    /**
      * Reload meta-data stores.
      *
      * @param strBuff  the string buffer
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
index f35a85a..a0cd5f7 100644
--- 
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
@@ -27,8 +27,7 @@ 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;
@@ -51,14 +50,12 @@ 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;
@@ -72,11 +69,7 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
             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
@@ -87,10 +80,6 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
     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
@@ -98,13 +87,13 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
     // replication nodes
     private Set<String> replicas4Transfer = new HashSet<>();
     private final Listener listener = new Listener();
+    // ha check thread
     private ExecutorService executorService = null;
     // bdb data store configure
     private final StoreConfig storeConfig = new StoreConfig();
 
     public BdbMetaConfigMapperImpl(MasterConfig masterConfig) {
-        super(masterConfig.getRowLockWaitDurMs());
-        this.masterConfig = masterConfig;
+        super(masterConfig);
         MasterReplicationConfig replicationConfig =
                 masterConfig.getReplicationConfig();
         // build replicationGroupAdmin info
@@ -155,12 +144,12 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
 
     @Override
     public void start() throws Exception {
-        logger.info("[BDB Impl] Start StoreManagerService, begin");
+        logger.info("[BDB Impl] Start MetaConfigService, begin");
         if (!srvStatus.compareAndSet(0, 1)) {
-            logger.info("[BDB Impl] Start StoreManagerService, started");
+            logger.info("[BDB Impl] Start MetaConfigService, started");
             return;
         }
-        logger.info("[BDB Impl] Starting StoreManagerService...");
+        logger.info("[BDB Impl] Starting MetaConfigService...");
         try {
             if (executorService != null) {
                 executorService.shutdownNow();
@@ -170,25 +159,25 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
             // build envHome file
             envHome = new File(masterConfig.getMetaDataPath());
             repEnv = getEnvironment();
-            initMetaStore();
+            initMetaStore(null);
             repEnv.setStateChangeListener(listener);
             srvStatus.compareAndSet(1, 2);
         } catch (Throwable ee) {
             srvStatus.compareAndSet(1, 0);
-            logger.error("[BDB Impl] Start StoreManagerService failure, 
error", ee);
+            logger.error("[BDB Impl] Start MetaConfigService failure, error", 
ee);
             return;
         }
-        logger.info("[BDB Impl] Start StoreManagerService, success");
+        logger.info("[BDB Impl] Start MetaConfigService, success");
     }
 
     @Override
     public void stop() throws Exception {
-        logger.info("[BDB Impl] Stop StoreManagerService, begin");
+        logger.info("[BDB Impl] Stop MetaConfigService, begin");
         if (!srvStatus.compareAndSet(2, 3)) {
-            logger.info("[BDB Impl] Stop StoreManagerService, stopped");
+            logger.info("[BDB Impl] Stop MetaConfigService, stopped");
             return;
         }
-        logger.info("[BDB Impl] Stopping StoreManagerService...");
+        logger.info("[BDB Impl] Stopping MetaConfigService...");
         // close bdb configure
         closeMetaStore();
         /* evn close */
@@ -205,23 +194,7 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
             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;
+        logger.info("[BDB Impl] Stop MetaConfigService, success");
     }
 
     @Override
@@ -235,7 +208,7 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
     }
 
     @Override
-    public InetSocketAddress getMasterAddress() {
+    public String getMasterAddress() {
         ReplicationGroup replicationGroup = getCurrReplicationGroup();
         if (replicationGroup == null) {
             logger.info("[BDB Impl] ReplicationGroup is null...please check 
the group status!");
@@ -247,7 +220,7 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
                         replicationGroupAdmin.getNodeState(node, 2000);
                 if (nodeState != null) {
                     if (nodeState.getNodeState().isMaster()) {
-                        return node.getSocketAddress();
+                        return 
node.getSocketAddress().getAddress().getHostAddress();
                     }
                 }
             } catch (Throwable e) {
@@ -268,7 +241,7 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
 
     @Override
     public void transferMaster() throws Exception {
-        if (!isStarted()) {
+        if (!isServiceStarted()) {
             throw new Exception("The BDB store StoreService is reboot now!");
         }
         if (isMasterNow()) {
@@ -411,6 +384,15 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
         return masterGroupStatus;
     }
 
+    protected void initMetaStore(StringBuilder strBuff) {
+        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);
+    }
+
     /**
      * State Change Listener,
      * through this object, it complete the metadata cache cleaning
@@ -481,23 +463,6 @@ public class BdbMetaConfigMapperImpl extends 
AbsMetaConfigMapperImpl {
         }
     }
 
-    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
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 975b7c4..63313ef 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
@@ -909,7 +909,7 @@ public class BdbMetaStoreServiceImpl implements 
MetaStoreService {
      * @return  the current master address
      */
     @Override
-    public InetSocketAddress getMasterAddress() {
+    public String getMasterAddress() {
         ReplicationGroup replicationGroup = getCurrReplicationGroup();
         if (replicationGroup == null) {
             logger.info("[BDB Impl] ReplicationGroup is null...please check 
the group status!");
@@ -921,7 +921,7 @@ public class BdbMetaStoreServiceImpl implements 
MetaStoreService {
                         replicationGroupAdmin.getNodeState(node, 2000);
                 if (nodeState != null) {
                     if (nodeState.getNodeState().isMaster()) {
-                        return node.getSocketAddress();
+                        return 
node.getSocketAddress().getAddress().getHostAddress();
                     }
                 }
             } catch (Throwable e) {
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKMetaConfigMapperImpl.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKMetaConfigMapperImpl.java
new file mode 100644
index 0000000..82738af
--- /dev/null
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKMetaConfigMapperImpl.java
@@ -0,0 +1,332 @@
+/**
+ * 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.zkimpl;
+
+import java.net.BindException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.inlong.tubemq.corebase.TBaseConstants;
+import org.apache.inlong.tubemq.corebase.TokenConstants;
+import org.apache.inlong.tubemq.corebase.cluster.NodeAddrInfo;
+import org.apache.inlong.tubemq.corebase.utils.Tuple2;
+import org.apache.inlong.tubemq.server.common.zookeeper.ZKUtil;
+import org.apache.inlong.tubemq.server.common.zookeeper.ZooKeeperWatcher;
+import org.apache.inlong.tubemq.server.master.MasterConfig;
+import org.apache.inlong.tubemq.server.master.bdbstore.MasterGroupStatus;
+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.apache.zookeeper.CreateMode;
+import org.apache.zookeeper.KeeperException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class ZKMetaConfigMapperImpl extends AbsMetaConfigMapperImpl {
+    private static final Logger logger =
+            LoggerFactory.getLogger(ZKMetaConfigMapperImpl.class);
+    private final MetaConfigSamplePrint metaSamplePrint =
+            new MetaConfigSamplePrint(logger);
+    // the meta data path in ZooKeeper
+    private final String metaZkRoot;
+    // the ha path in ZooKeeper
+    private final String haNodesPath;
+    // whether is the first start
+    private volatile boolean isFirstChk = true;
+    // the ZooKeeper watcher
+    private ZooKeeperWatcher zkWatcher;
+    // the local node address
+    private final NodeAddrInfo localNodeAdd;
+    private final ScheduledExecutorService executorService;
+
+    static {
+        if (Thread.getDefaultUncaughtExceptionHandler() == null) {
+            Thread.setDefaultUncaughtExceptionHandler((t, e) -> {
+                if (e instanceof BindException) {
+                    logger.error("[ZK Impl] Bind failed.", e);
+                    // System.exit(1);
+                }
+                if (e instanceof IllegalStateException
+                        && e.getMessage().contains("Shutdown in progress")) {
+                    return;
+                }
+                logger.warn("[ZK Impl] Thread terminated with exception: " + 
t.getName(), e);
+            });
+        }
+    }
+
+    public ZKMetaConfigMapperImpl(MasterConfig masterConfig) {
+        super(masterConfig);
+        this.localNodeAdd = new NodeAddrInfo(
+                masterConfig.getHostName(), masterConfig.getPort());
+        String tubeZkRoot = 
ZKUtil.normalizePath(masterConfig.getZkConfig().getZkNodeRoot());
+        StringBuilder strBuff =
+                new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE);
+        this.metaZkRoot = 
strBuff.append(tubeZkRoot).append(TokenConstants.SLASH)
+                .append(TZKNodeKeys.ZK_BRANCH_META_DATA).toString();
+        strBuff.delete(0, strBuff.length());
+        this.haNodesPath = 
strBuff.append(tubeZkRoot).append(TokenConstants.SLASH)
+                
.append(TZKNodeKeys.ZK_BRANCH_HA).append("/nodeIds").toString();
+        strBuff.delete(0, strBuff.length());
+        initMetaStore(strBuff);
+        this.executorService =
+                Executors.newSingleThreadScheduledExecutor(new ThreadFactory() 
{
+                    @Override
+                    public Thread newThread(Runnable r) {
+                        return new Thread(r, "Master selector thread");
+                    }
+                });
+    }
+
+    @Override
+    public void start() throws Exception {
+        logger.info("[ZK Impl] Start MetaConfigService, begin");
+        if (!srvStatus.compareAndSet(0, 1)) {
+            logger.info("[ZK Impl] Start MetaConfigService, started");
+            return;
+        }
+        try {
+            logger.info("[ZK Impl] Starting MetaConfigService...");
+            // start Master select thread
+            executorService.scheduleWithFixedDelay(new MasterSelectorTask(), 
5L,
+                    masterConfig.getZkConfig().getZkMasterCheckPeriodMs(), 
TimeUnit.MILLISECONDS);
+            // sleep 1 second for select
+            Thread.sleep(1000);
+            srvStatus.compareAndSet(1, 2);
+        } catch (Throwable ee) {
+            srvStatus.compareAndSet(1, 0);
+            logger.error("[ZK Impl] Start MetaConfigService failure, error", 
ee);
+            return;
+        }
+        logger.info("[ZK Impl] Start MetaConfigService, success");
+    }
+
+    @Override
+    public void stop() throws Exception {
+        logger.info("[ZK Impl] Stop MetaConfigService, begin");
+        if (!srvStatus.compareAndSet(2, 3)) {
+            logger.info("[ZK Impl] Stop MetaConfigService, stopped");
+            return;
+        }
+        logger.info("[ZK Impl] Stopping MetaConfigService...");
+        // close Master select thread
+        executorService.shutdownNow();
+        // close ZooKeeper watcher
+        zkWatcher.close();
+        // clear configure
+        closeMetaStore();
+        // set status
+        srvStatus.set(0);
+        logger.info("[BDB Impl] Stop MetaConfigService, success");
+    }
+
+    @Override
+    public boolean isMasterNow() {
+        return isMaster;
+    }
+
+    @Override
+    public long getMasterSinceTime() {
+        return masterSinceTime.get();
+    }
+
+    @Override
+    public String getMasterAddress() {
+        Tuple2<Boolean, Long> queryResult;
+        Map<Long, String> clusterNodeMap = new HashMap<>();
+        try {
+            queryResult = getCusterNodes(clusterNodeMap);
+        } catch (Throwable e) {
+            logger.error("[ZK Impl] Get Master Address Throwable error", e);
+            return null;
+        }
+        if (!clusterNodeMap.isEmpty()) {
+            return clusterNodeMap.get(queryResult.getF1()).split(":")[0];
+        }
+        return null;
+    }
+
+    @Override
+    public boolean isPrimaryNodeActive() {
+        return false;
+    }
+
+    @Override
+    public void transferMaster() {
+        // ignore
+    }
+
+    @Override
+    public ClusterGroupVO getGroupAddressStrInfo() {
+        ClusterGroupVO clusterGroupVO = new ClusterGroupVO();
+        clusterGroupVO.setGroupStatus("Abnormal");
+        clusterGroupVO.setGroupName("ZooKeeper HA Cluster");
+        Tuple2<Boolean, Long> queryResult;
+        Map<Long, String> clusterNodeMap = new HashMap<>();
+        try {
+            queryResult = getCusterNodes(clusterNodeMap);
+        } catch (Throwable e) {
+            logger.error("[ZK Impl] Get Master Address Throwable error", e);
+            return clusterGroupVO;
+        }
+        if (clusterNodeMap.isEmpty()) {
+            return clusterGroupVO;
+        }
+        String nodeAdd;
+        List<ClusterNodeVO> clusterNodeVOs = new ArrayList<>();
+        for (Map.Entry<Long, String> entry : clusterNodeMap.entrySet()) {
+            nodeAdd = entry.getValue();
+            clusterNodeVOs.add(new ClusterNodeVO(nodeAdd, 
nodeAdd.split(":")[0],
+                    Integer.parseInt(nodeAdd.split(":")[1]),
+                    entry.getKey().equals(queryResult.getF1()) ? "Master" : 
"Slave", 0));
+        }
+        if (clusterNodeMap.isEmpty()) {
+            return clusterGroupVO;
+        }
+        clusterGroupVO.setNodeData(clusterNodeVOs);
+        clusterGroupVO.setGroupStatus("Running-ReadWrite");
+        return clusterGroupVO;
+    }
+
+    @Override
+    public MasterGroupStatus getMasterGroupStatus(boolean isFromHeartbeat) {
+        return null;
+    }
+
+    protected void initMetaStore(StringBuilder strBuff) {
+        clusterConfigMapper = new ZKClusterConfigMapperImpl(metaZkRoot, 
zkWatcher, strBuff);
+        brokerConfigMapper = new ZKBrokerConfigMapperImpl(metaZkRoot, 
zkWatcher, strBuff);
+        topicDeployMapper =  new ZKTopicDeployMapperImpl(metaZkRoot, 
zkWatcher, strBuff);
+        groupResCtrlMapper = new ZKGroupResCtrlMapperImpl(metaZkRoot, 
zkWatcher, strBuff);
+        topicCtrlMapper = new ZKTopicCtrlMapperImpl(metaZkRoot, zkWatcher, 
strBuff);
+        consumeCtrlMapper = new ZKConsumeCtrlMapperImpl(metaZkRoot, zkWatcher, 
strBuff);
+    }
+
+    /**
+     * Master selector logic
+     */
+    private class MasterSelectorTask implements Runnable {
+        @Override
+        public void run() {
+            long startTime = System.currentTimeMillis();
+            StringBuilder strBuff = new 
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE);
+            try {
+                if (isFirstChk) {
+                    // check whether the HA directory already exists on ZK
+                    if (ZKUtil.checkExists(zkWatcher, haNodesPath) == -1) {
+                        // create path if not exists
+                        ZKUtil.createWithParents(zkWatcher, haNodesPath);
+                    } else {
+                        isFirstChk = false;
+                    }
+                }
+                int totalCnt = 0;
+                Tuple2<Boolean, Long> queryResult;
+                Map<Long, String> clusterNodeMap = new HashMap<>();
+                do {
+                    queryResult = getCusterNodes(clusterNodeMap);
+                    if (!queryResult.getF0()) {
+                        ZKUtil.createEphemeralNodeAndWatch(zkWatcher, 
haNodesPath,
+                                localNodeAdd.getHostPortStr().getBytes(),
+                                CreateMode.EPHEMERAL_SEQUENTIAL);
+                    }
+                    if (totalCnt++ > 1) {
+                        break;
+                    }
+                } while (clusterNodeMap.isEmpty());
+                // judge whether the current master node is the current node
+                if (!clusterNodeMap.isEmpty()) {
+                    String masterAdd = clusterNodeMap.get(queryResult.getF1());
+                    if (localNodeAdd.getHostPortStr().equals(masterAdd)) {
+                        if (!isMaster) {
+                            isMaster = true;
+                            reloadMetaStore(strBuff);
+                            masterSinceTime.set(System.currentTimeMillis());
+                            logger.warn(strBuff.append("[ZK Impl] HA switched, 
")
+                                    .append(localNodeAdd.getHostPortStr())
+                                    .append(" has changed to Master 
role.").toString());
+                            strBuff.delete(0, strBuff.length());
+                        }
+                    } else {
+                        if (isMaster) {
+                            isMaster = false;
+                            closeMetaStore();
+                            logger.warn(strBuff.append("[ZK Impl] HA switched, 
")
+                                    .append(localNodeAdd.getHostPortStr())
+                                    .append(" has changed to Slave 
role.").toString());
+                            strBuff.delete(0, strBuff.length());
+                        }
+                    }
+                }
+            } catch (KeeperException e) {
+                metaSamplePrint.printExceptionCaught(e,
+                        masterConfig.getHostName(), 
masterConfig.getHostName());
+                if (isMaster) {
+                    isMaster = false;
+                    closeMetaStore();
+                    logger.warn(strBuff.append("[ZK Impl] HA select exception, 
")
+                            .append(localNodeAdd.getHostPortStr())
+                            .append(" has changed to Slave role.").toString());
+                    strBuff.delete(0, strBuff.length());
+                }
+            }
+            // log the case which delta time over 30 seconds
+            if (System.currentTimeMillis() - startTime > 30000) {
+                logger.warn(strBuff.append("[ZK Impl] HA Select cost:")
+                        .append(System.currentTimeMillis() - startTime)
+                        .append("ms, please make sure the time cost is below 
30seconds(zk session timeout).")
+                        .toString());
+                strBuff.delete(0, strBuff.length());
+            }
+        }
+    }
+
+    private Tuple2<Boolean, Long> getCusterNodes(
+            Map<Long, String> clusterNodeMap) throws KeeperException {
+        String nodeAdd;
+        long materNodeId;
+        boolean foundSelf = false;
+        long minNodeId = Long.MAX_VALUE;
+        List<ZKUtil.NodeAndData> childNodes =
+                ZKUtil.getChildDataAndWatchForNewChildren(zkWatcher, 
haNodesPath);
+        for (ZKUtil.NodeAndData child : childNodes) {
+            // select the first registered node as Master
+            if (child == null) {
+                continue;
+            }
+            nodeAdd = new String(child.getData());
+            materNodeId = Long.parseLong(child.getNode().split(":")[1]);
+            if (minNodeId > materNodeId) {
+                minNodeId = materNodeId;
+            }
+            if (localNodeAdd.getHostPortStr().equals(nodeAdd)) {
+                foundSelf = true;
+            }
+            clusterNodeMap.put(materNodeId, nodeAdd);
+        }
+        return new Tuple2<>(foundSelf, minNodeId);
+    }
+
+}
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/MasterStatusCheckFilter.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/MasterStatusCheckFilter.java
index e413274..19593df 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/MasterStatusCheckFilter.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/MasterStatusCheckFilter.java
@@ -18,7 +18,6 @@
 package org.apache.inlong.tubemq.server.master.web;
 
 import java.io.IOException;
-import java.net.InetSocketAddress;
 import javax.servlet.Filter;
 import javax.servlet.FilterChain;
 import javax.servlet.FilterConfig;
@@ -53,15 +52,16 @@ public class MasterStatusCheckFilter implements Filter {
         HttpServletRequest req = (HttpServletRequest) request;
         HttpServletResponse resp = (HttpServletResponse) response;
         if (!metaDataManager.isSelfMaster()) {
-            InetSocketAddress masterAddr =
+            String masterAdd =
                     metaDataManager.getMasterAddress();
-            if (masterAddr == null) {
+            if (masterAdd == null) {
                 throw new IOException("Not found the master node address!");
             }
-            StringBuilder sBuilder = new StringBuilder(512).append("http://";)
-                    .append(masterAddr.getAddress().getHostAddress())
-                    .append(":").append(master.getMasterConfig().getWebPort())
-                    .append(req.getRequestURI());
+            StringBuilder sBuilder =
+                    new StringBuilder(512).append("http://";)
+                            .append(masterAdd).append(":")
+                            .append(master.getMasterConfig().getWebPort())
+                            .append(req.getRequestURI());
             if (TStringUtils.isNotBlank(req.getQueryString())) {
                 sBuilder.append("?").append(req.getQueryString());
             }
diff --git 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/action/screen/Tubeweb.java
 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/action/screen/Tubeweb.java
index db46107..48e6956 100644
--- 
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/action/screen/Tubeweb.java
+++ 
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/web/action/screen/Tubeweb.java
@@ -17,7 +17,6 @@
 
 package org.apache.inlong.tubemq.server.master.web.action.screen;
 
-import java.net.InetSocketAddress;
 import org.apache.inlong.tubemq.server.master.TMaster;
 import org.apache.inlong.tubemq.server.master.metamanage.MetaDataManager;
 import org.apache.inlong.tubemq.server.master.web.simplemvc.Action;
@@ -34,13 +33,13 @@ public class Tubeweb implements Action {
     @Override
     public void execute(RequestContext context) {
         MetaDataManager metaDataManager = this.master.getDefMetaDataManager();
-        InetSocketAddress masterAddr = metaDataManager.getMasterAddress();
-        if (master.getMasterConfig().isUseWebProxy() || masterAddr == null) {
+        String masterAdd = metaDataManager.getMasterAddress();
+        if (master.getMasterConfig().isUseWebProxy() || masterAdd == null) {
             // use absolute path
             context.put("tubemqRemoteAddr", "");
         } else {
             // use the whole path of the active master
-            context.put("tubemqRemoteAddr", "http://"; + 
masterAddr.getAddress().getHostAddress() + ":"
+            context.put("tubemqRemoteAddr", "http://"; + masterAdd + ":"
                     + master.getMasterConfig().getWebPort());
         }
     }

Reply via email to