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());
}
}