This is an automated email from the ASF dual-hosted git repository.
gosonzhang 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 6277e0d [INLONG-3093][TubeMQ] Optimize the AbsXXXMapperImpl
implementation classes (#3094)
6277e0d is described below
commit 6277e0d1c31315b1b9f76ca750c87de83426b57a
Author: gosonzhang <[email protected]>
AuthorDate: Sun Mar 13 07:19:53 2022 +0800
[INLONG-3093][TubeMQ] Optimize the AbsXXXMapperImpl implementation classes
(#3094)
---
.../server/common/utils/WebParameterUtils.java | 21 +-
.../server/master/metamanage/DataOpErrCode.java | 22 +-
.../server/master/metamanage/MetaDataManager.java | 15 +-
.../master/metamanage/keepalive/KeepAlive.java | 15 ++
.../metastore/dao/mapper/BrokerConfigMapper.java | 16 ++
.../metastore/impl/AbsBrokerConfigMapperImpl.java | 190 +++++++++++----
.../metastore/impl/AbsClusterConfigMapperImpl.java | 49 +++-
.../metastore/impl/AbsConsumeCtrlMapperImpl.java | 56 +++--
.../metastore/impl/AbsGroupResCtrlMapperImpl.java | 41 ++--
.../metastore/impl/AbsTopicCtrlMapperImpl.java | 40 ++--
.../metastore/impl/AbsTopicDeployMapperImpl.java | 263 +++++++++++++++++++--
.../impl/bdbimpl/BdbBrokerConfigMapperImpl.java | 5 +-
.../impl/bdbimpl/BdbClusterConfigMapperImpl.java | 2 +-
.../impl/bdbimpl/BdbConsumeCtrlMapperImpl.java | 4 +-
.../impl/bdbimpl/BdbGroupResCtrlMapperImpl.java | 2 +-
.../bdbimpl}/BdbMetaStoreServiceImpl.java | 78 ++----
.../impl/bdbimpl/BdbTopicCtrlMapperImpl.java | 4 +-
.../impl/bdbimpl/BdbTopicDeployMapperImpl.java | 2 +-
.../metastore/impl/zkimpl/TZKNodeKeys.java | 1 +
.../impl/zkimpl/ZKBrokerConfigMapperImpl.java | 6 +-
.../impl/zkimpl/ZKClusterConfigMapperImpl.java | 2 +-
.../impl/zkimpl/ZKConsumeCtrlMapperImpl.java | 6 +-
.../impl/zkimpl/ZKGroupResCtrlMapperImpl.java | 6 +-
.../impl/zkimpl/ZKTopicCtrlMapperImpl.java | 2 +-
.../impl/zkimpl/ZKTopicDeployMapperImpl.java | 2 +-
25 files changed, 609 insertions(+), 241 deletions(-)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/utils/WebParameterUtils.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/utils/WebParameterUtils.java
index f7b3978..0bb3c71 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/utils/WebParameterUtils.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/common/utils/WebParameterUtils.java
@@ -371,29 +371,30 @@ public class WebParameterUtils {
}
/**
- * Compare whether the configured port values conflict
+ * Compare the configured ports for conflicts
*
* @param brokerPort broker port
* @param brokerTlsPort broker tls port
* @param brokerWebPort broker web port
- * @param sBuffer string buffer
+ * @param strBuff string buffer
* @param result check result of parameter value
- * @return process result
+ * @return true for illegal, false for legal
*/
- public static boolean isLegallyPortValueSet(int brokerPort, int
brokerTlsPort,
- int brokerWebPort,
StringBuilder sBuffer,
- ProcessResult result) {
- result.setSuccResult(null);
+ public static boolean isConflictedPortsSet(int brokerPort, int
brokerTlsPort,
+ int brokerWebPort,
StringBuilder strBuff,
+ ProcessResult result) {
if (brokerPort == brokerWebPort || brokerTlsPort == brokerWebPort) {
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
- sBuffer.append("Illegal port value configuration, the
value of ")
+
strBuff.append(DataOpErrCode.DERR_CONFLICT_VALUE.getDescription())
+ .append(", the value of ")
.append(WebFieldDef.BROKERPORT.name).append(" or ")
.append(WebFieldDef.BROKERTLSPORT.name)
.append(" cannot be the same as the value of")
.append(WebFieldDef.BROKERWEBPORT.name).toString());
- sBuffer.delete(0, sBuffer.length());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
}
- return result.isSuccess();
+ return false;
}
/**
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/DataOpErrCode.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/DataOpErrCode.java
index fdcbb8a..2e9223e 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/DataOpErrCode.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/DataOpErrCode.java
@@ -18,20 +18,22 @@
package org.apache.inlong.tubemq.server.master.metamanage;
public enum DataOpErrCode {
- DERR_SUCCESS(200, "Success."),
+ DERR_SUCCESS(200, "Success"),
DERR_SUCCESS_UNCHANGED(201, "Success, but unchanged"),
- DERR_NOT_EXIST(401, "Record not exist."),
- DERR_EXISTED(402, "Record has existed."),
- DERR_UNCHANGED(403, "Record not changed."),
- DERR_UNCLEANED(404, "Related configuration is not cleaned up."),
+ DERR_NOT_EXIST(401, "Record not exist"),
+ DERR_EXISTED(402, "Record has existed"),
+ DERR_UNCHANGED(403, "Record not changed"),
+ DERR_UNCLEANED(404, "Related configuration is not cleaned up"),
DERR_CONDITION_LACK(405, "The preconditions are not met"),
DERR_ILLEGAL_STATUS(406, "Illegal operate status"),
DERR_ILLEGAL_VALUE(407, "Illegal data format or value"),
- DERR_STORE_ABNORMAL(501, "Store layer throw exception."),
- DERR_UPD_NOT_EXIST(502, "Record updated but not exist."),
- DERR_STORE_STOPPED(510, "Store stopped."),
- DERR_STORE_NOT_MASTER(511, "Store not active master."),
- DERR_MASTER_UNKNOWN(599, "Unknown error.");
+ DERR_CONFLICT_VALUE(408, "Conflicted configure value"),
+ DERR_STORE_ABNORMAL(501, "Store layer throw exception"),
+ DERR_UPD_NOT_EXIST(502, "Record updated but not exist"),
+ DERR_STORE_STOPPED(510, "Store stopped"),
+ DERR_STORE_NOT_MASTER(511, "Store not active master"),
+ DERR_STORE_LOCK_FAILURE(512, "Failed to lock metadata lock"),
+ DERR_MASTER_UNKNOWN(599, "Unknown error");
private int code;
private String description;
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 f7f16f4..6da5457 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
@@ -48,7 +48,7 @@ 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.BdbMetaStoreServiceImpl;
+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;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
@@ -88,8 +88,7 @@ public class MetaDataManager implements Server {
MasterConfig masterConfig = this.tMaster.getMasterConfig();
this.replicationConfig = masterConfig.getReplicationConfig();
this.metaStoreService =
- new BdbMetaStoreServiceImpl(masterConfig.getHostName(),
- masterConfig.getMetaDataPath(),
this.replicationConfig);
+ new BdbMetaStoreServiceImpl(tMaster.getMasterConfig());
this.scheduledExecutorService =
Executors.newSingleThreadScheduledExecutor(new ThreadFactory()
{
@@ -426,7 +425,7 @@ public class MetaDataManager implements Server {
if (isAddOp) {
if (metaStoreService.getBrokerConfByBrokerId(entity.getBrokerId())
== null
&&
metaStoreService.getBrokerConfByBrokerIp(entity.getBrokerIp()) == null) {
- if
(WebParameterUtils.isLegallyPortValueSet(entity.getBrokerPort(),
+ if
(!WebParameterUtils.isConflictedPortsSet(entity.getBrokerPort(),
entity.getBrokerTLSPort(), entity.getBrokerWebPort(),
sBuffer, result)) {
if (metaStoreService.addBrokerConf(entity, sBuffer,
result)) {
this.tMaster.getBrokerRunManager().updBrokerStaticInfo(entity);
@@ -454,7 +453,7 @@ public class MetaDataManager implements Server {
entity.getBrokerTLSPort(), entity.getBrokerWebPort(),
entity.getRegionId(), entity.getGroupId(),
entity.getManageStatus(), entity.getTopicProps())) {
- if
(WebParameterUtils.isLegallyPortValueSet(newEntity.getBrokerPort(),
+ if
(!WebParameterUtils.isConflictedPortsSet(newEntity.getBrokerPort(),
newEntity.getBrokerTLSPort(),
newEntity.getBrokerWebPort(),
sBuffer, result)) {
if (metaStoreService.updBrokerConf(newEntity, sBuffer,
result)) {
@@ -503,7 +502,7 @@ public class MetaDataManager implements Server {
BrokerConfEntity curEntry;
BrokerConfEntity newEntry;
List<BrokerProcessResult> retInfo = new ArrayList<>();
- // check target broker configure's status
+ // check target broker configures status
for (Integer brokerId : brokerIdSet) {
curEntry = metaStoreService.getBrokerConfByBrokerId(brokerId);
if (curEntry == null) {
@@ -1581,7 +1580,7 @@ public class MetaDataManager implements Server {
newConf.updModifyInfo(opEntity.getDataVerId(), brokerPort,
brokerTlsPort, brokerWebPort, maxMsgSizeMB, qryPriorityId,
flowCtrlEnable, flowRuleCnt, flowCtrlInfo, topicProps);
- if
(WebParameterUtils.isLegallyPortValueSet(newConf.getBrokerPort(),
+ if
(!WebParameterUtils.isConflictedPortsSet(newConf.getBrokerPort(),
newConf.getBrokerTLSPort(), newConf.getBrokerWebPort(),
sBuffer, result)) {
metaStoreService.addClusterConfig(newConf, sBuffer, result);
}
@@ -1591,7 +1590,7 @@ public class MetaDataManager implements Server {
if (newConf.updModifyInfo(opEntity.getDataVerId(), brokerPort,
brokerTlsPort, brokerWebPort, maxMsgSizeMB, qryPriorityId,
flowCtrlEnable, flowRuleCnt, flowCtrlInfo, topicProps)) {
- if
(WebParameterUtils.isLegallyPortValueSet(newConf.getBrokerPort(),
+ if
(!WebParameterUtils.isConflictedPortsSet(newConf.getBrokerPort(),
newConf.getBrokerTLSPort(), newConf.getBrokerWebPort(),
sBuffer, result)) {
metaStoreService.updClusterConfig(newConf, sBuffer,
result);
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/keepalive/KeepAlive.java
index 9d4bd84..57f1ae0 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/keepalive/KeepAlive.java
@@ -23,16 +23,31 @@ import
org.apache.inlong.tubemq.server.master.web.model.ClusterGroupVO;
public interface KeepAlive {
+ /**
+ * Whether this node is the master role
+ *
+ * @return true if is master role or else
+ */
boolean isMasterNow();
long getMasterSinceTime();
InetSocketAddress getMasterAddress();
+ /**
+ * Whether the primary node in active
+ *
+ * @return true for active, false for inactive
+ */
boolean isPrimaryNodeActive();
void transferMaster() throws Exception;
+ /**
+ * Register node role switching event observer
+ *
+ * @param eventObserver the event observer
+ */
void registerObserver(AliveObserver eventObserver);
ClusterGroupVO getGroupAddressStrInfo();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/BrokerConfigMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/BrokerConfigMapper.java
index 86f87b8..2a9ff53 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/BrokerConfigMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/BrokerConfigMapper.java
@@ -20,6 +20,8 @@ package
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper;
import java.util.Map;
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.master.metamanage.metastore.dao.entity.BaseEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
public interface BrokerConfigMapper extends AbstractMapper {
@@ -47,6 +49,20 @@ public interface BrokerConfigMapper extends AbstractMapper {
StringBuilder strBuff, ProcessResult result);
/**
+ * Update a broker manage status
+ *
+ * @param opEntity the operator information
+ * @param brokerId the broker id need to updated
+ * @param newMngStatus the new manage status
+ * @param strBuff the string buffer
+ * @param result process result with old value
+ * @return the process result
+ */
+ boolean updBrokerMngStatus(BaseEntity opEntity,
+ Integer brokerId, ManageStatus newMngStatus,
+ StringBuilder strBuff, ProcessResult result);
+
+ /**
* delete broker configure info from store
*
* @param brokerId the broker id to be deleted
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsBrokerConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsBrokerConfigMapperImpl.java
index 817b247..9ee35c3 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsBrokerConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsBrokerConfigMapperImpl.java
@@ -24,7 +24,11 @@ import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
import org.apache.inlong.tubemq.corebase.utils.ConcurrentHashSet;
+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.utils.WebParameterUtils;
import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
+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.mapper.BrokerConfigMapper;
import org.slf4j.Logger;
@@ -48,27 +52,27 @@ public abstract class AbsBrokerConfigMapperImpl implements
BrokerConfigMapper {
@Override
public boolean addBrokerConf(BrokerConfEntity entity,
StringBuilder strBuff, ProcessResult result) {
- BrokerConfEntity curEntity =
- brokerConfCache.get(entity.getBrokerId());
- if (curEntity != null) {
+ // Check whether the brokerId or broker Ip conflict with existing
records
+ if (brokerConfCache.get(entity.getBrokerId()) != null
+ || brokerIpIndexCache.get(entity.getBrokerIp()) != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The broker's brokerId
").append(entity.getBrokerId())
- .append(" has already exists, the value must be
unique!")
+ strBuff.append("Existed record found for ")
+ .append(WebFieldDef.BROKERID.name).append("(")
+ .append(entity.getBrokerId()).append(") or ")
+ .append(WebFieldDef.BROKERIP.name).append("(")
+ .append(entity.getBrokerIp()).append(") value!")
.toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- Integer curBrokerId = brokerIpIndexCache.get(entity.getBrokerIp());
- if (curBrokerId != null) {
- result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The broker's brokerIp
").append(entity.getBrokerIp())
- .append(" has already exists, the value must be
unique!")
- .toString());
- strBuff.delete(0, strBuff.length());
+ // Check whether the configured ports conflict in the record
+ if (WebParameterUtils.isConflictedPortsSet(entity.getBrokerPort(),
+ entity.getBrokerTLSPort(), entity.getBrokerWebPort(), strBuff,
result)) {
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ putRecord2Caches(entity);
}
return result.isSuccess();
}
@@ -76,26 +80,80 @@ public abstract class AbsBrokerConfigMapperImpl implements
BrokerConfigMapper {
@Override
public boolean updBrokerConf(BrokerConfEntity entity,
StringBuilder strBuff, ProcessResult result) {
+ // Check the existence of records by brokerId
BrokerConfEntity curEntity =
brokerConfCache.get(entity.getBrokerId());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- strBuff.append("The broker configure
").append(entity.getBrokerIp())
- .append(" is not exists, please add it first!")
- .toString());
+ strBuff.append("Not found broker configure for ")
+ .append(WebFieldDef.BROKERID.name).append("(")
+
.append(entity.getBrokerId()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ BrokerConfEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.getBrokerPort(), entity.getBrokerTLSPort(),
+ entity.getBrokerWebPort(), entity.getRegionId(),
+ entity.getGroupId(), entity.getManageStatus(),
+ entity.getTopicProps())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- strBuff.append("The broker configure
").append(entity.getBrokerIp())
- .append(" have not changed, please confirm it
first!")
- .toString());
+ "Broker configure not changed!");
+ return result.isSuccess();
+ }
+ // Check whether the configured ports conflict in the record
+ if (WebParameterUtils.isConflictedPortsSet(newEntity.getBrokerPort(),
+ newEntity.getBrokerTLSPort(), newEntity.getBrokerWebPort(),
+ strBuff, result)) {
+ return result.isSuccess();
+ }
+ // Check manage status
+ if (isIllegalManageStatusChange(newEntity, curEntity, strBuff,
result)) {
+ return result.isSuccess();
+ }
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ putRecord2Caches(newEntity);
+ result.setSuccResult(curEntity);
+ }
+ return result.isSuccess();
+ }
+
+ @Override
+ public boolean updBrokerMngStatus(BaseEntity opEntity,
+ Integer brokerId, ManageStatus
newMngStatus,
+ StringBuilder strBuff, ProcessResult
result) {
+ // Check the existence of records by brokerId
+ BrokerConfEntity curEntity = brokerConfCache.get(brokerId);
+ if (curEntity == null) {
+ result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
+ strBuff.append("Not found broker configure for ")
+ .append(WebFieldDef.BROKERID.name).append("(")
+ .append(brokerId).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ // Build the entity that need to be updated
+ BrokerConfEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(opEntity);
+ if (!newEntity.updModifyInfo(opEntity.getDataVerId(),
+ curEntity.getBrokerPort(), curEntity.getBrokerTLSPort(),
+ curEntity.getBrokerWebPort(), curEntity.getRegionId(),
+ curEntity.getGroupId(), newMngStatus,
+ curEntity.getTopicProps())) {
+ result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
+ "Broker configure not changed!");
+ return result.isSuccess();
+ }
+ // Check manage status
+ if (isIllegalManageStatusChange(newEntity, curEntity, strBuff,
result)) {
+ return result.isSuccess();
+ }
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ putRecord2Caches(newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -105,12 +163,24 @@ public abstract class AbsBrokerConfigMapperImpl
implements BrokerConfigMapper {
public boolean delBrokerConf(int brokerId, StringBuilder strBuff,
ProcessResult result) {
BrokerConfEntity curEntity =
brokerConfCache.get(brokerId);
+ // Check the existence of records by brokerId
if (curEntity == null) {
result.setSuccResult(null);
return result.isSuccess();
}
+ // Check broker's manage status
+ if (curEntity.getManageStatus().isOnlineStatus()) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("Illegal manage status, please offline the
broker(")
+ .append(WebFieldDef.BROKERID.name).append("=")
+ .append(curEntity.getBrokerId()).append(")
first!").toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ // Delete record from persistent
delConfigFromPersistent(brokerId, strBuff);
- delCacheRecord(brokerId);
+ // Clear cache data
+ delRecordFromCaches(brokerId);
result.setSuccResult(curEntity);
return result.isSuccess();
}
@@ -250,11 +320,11 @@ public abstract class AbsBrokerConfigMapperImpl
implements BrokerConfigMapper {
}
/**
- * Add or update a record
+ * Add or update a record to caches
*
* @param entity need added or updated entity
*/
- protected void addOrUpdCacheRecord(BrokerConfEntity entity) {
+ protected void putRecord2Caches(BrokerConfEntity entity) {
brokerConfCache.put(entity.getBrokerId(), entity);
// add brokerId info
Integer brokerId = brokerIpIndexCache.get(entity.getBrokerIp());
@@ -272,21 +342,6 @@ public abstract class AbsBrokerConfigMapperImpl implements
BrokerConfigMapper {
brokerIdSet.add(entity.getBrokerId());
}
- private void delCacheRecord(int brokerId) {
- BrokerConfEntity curEntity =
- brokerConfCache.remove(brokerId);
- if (curEntity == null) {
- return;
- }
- brokerIpIndexCache.remove(curEntity.getBrokerIp());
- ConcurrentHashSet<Integer> brokerIdSet =
- regionIndexCache.get(curEntity.getRegionId());
- if (brokerIdSet == null) {
- return;
- }
- brokerIdSet.remove(brokerId);
- }
-
/**
* Put broker configure information into persistent storage
*
@@ -296,7 +351,8 @@ public abstract class AbsBrokerConfigMapperImpl implements
BrokerConfigMapper {
* @return the process result
*/
protected abstract boolean putConfig2Persistent(BrokerConfEntity entity,
- StringBuilder strBuff,
ProcessResult result);
+ StringBuilder strBuff,
+ ProcessResult result);
/**
* Delete broker configure information from persistent storage
@@ -307,4 +363,56 @@ public abstract class AbsBrokerConfigMapperImpl implements
BrokerConfigMapper {
*/
protected abstract boolean delConfigFromPersistent(int brokerId,
StringBuilder strBuff);
+ /**
+ * Delete the record from caches
+ *
+ * @param brokerId need deleted broker id
+ */
+ private void delRecordFromCaches(int brokerId) {
+ BrokerConfEntity curEntity =
+ brokerConfCache.remove(brokerId);
+ if (curEntity == null) {
+ return;
+ }
+ brokerIpIndexCache.remove(curEntity.getBrokerIp());
+ ConcurrentHashSet<Integer> brokerIdSet =
+ regionIndexCache.get(curEntity.getRegionId());
+ if (brokerIdSet == null) {
+ return;
+ }
+ brokerIdSet.remove(brokerId);
+ }
+
+ /**
+ * Check whether the management status change is illegal
+ *
+ * @param newEntity the entity to be updated
+ * @param curEntity the current entity
+ * @param strBuff string buffer
+ * @param result check result of parameter value
+ * @return true for illegal, false for legal
+ */
+ private boolean isIllegalManageStatusChange(BrokerConfEntity newEntity,
+ BrokerConfEntity curEntity,
+ StringBuilder strBuff,
+ ProcessResult result) {
+ if (newEntity.getManageStatus() == curEntity.getManageStatus()) {
+ return false;
+ }
+ if (((newEntity.getManageStatus().getCode() <
ManageStatus.STATUS_MANAGE_ONLINE.getCode())
+ && (curEntity.getManageStatus().getCode() >=
ManageStatus.STATUS_MANAGE_ONLINE.getCode()))
+ || ((newEntity.getManageStatus().getCode() >
ManageStatus.STATUS_MANAGE_ONLINE.getCode())
+ && (curEntity.getManageStatus().getCode() <
ManageStatus.STATUS_MANAGE_ONLINE.getCode()))) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
+ strBuff.append("Illegal manage status, cannot reverse ")
+ .append(WebFieldDef.MANAGESTATUS.name).append("
from ")
+
.append(curEntity.getManageStatus().getDescription())
+ .append(" to
").append(newEntity.getManageStatus().getDescription())
+ .append(" for the
broker(").append(WebFieldDef.BROKERID.name).append("=")
+
.append(curEntity.getBrokerId()).append(")!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ return false;
+ }
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsClusterConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsClusterConfigMapperImpl.java
index 6f47cb2..827c44a 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsClusterConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsClusterConfigMapperImpl.java
@@ -20,6 +20,7 @@ package
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
+import org.apache.inlong.tubemq.server.common.utils.WebParameterUtils;
import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.TStoreConstants;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.ClusterSettingEntity;
@@ -40,11 +41,22 @@ public abstract class AbsClusterConfigMapperImpl implements
ClusterConfigMapper
@Override
public boolean addClusterConfig(ClusterSettingEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- if (!metaDataCache.isEmpty()) {
+ // Check whether the configure record already exist
+ ClusterSettingEntity curEntity =
metaDataCache.get(entity.getRecordKey());
+ if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- "The cluster configure already exists, please delete or
update it first!");
+ strBuff.append("Existed record found for ")
+ .append(entity.getRecordKey()).toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ // Check whether the configured ports conflict in the record
+ if (WebParameterUtils.isConflictedPortsSet(entity.getBrokerPort(),
+ entity.getBrokerTLSPort(), entity.getBrokerWebPort(),
+ strBuff, result)) {
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
metaDataCache.put(entity.getRecordKey(), entity);
}
@@ -54,19 +66,36 @@ public abstract class AbsClusterConfigMapperImpl implements
ClusterConfigMapper
@Override
public boolean updClusterConfig(ClusterSettingEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- if (metaDataCache.isEmpty()) {
+ // Check for the record to be updated
+ ClusterSettingEntity curEntity =
metaDataCache.get(entity.getRecordKey());
+ if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- "The cluster configure is null, please add it first!");
+ strBuff.append("Not found cluster configure for (")
+
.append(entity.getRecordKey()).append(")!").toString());
return result.isSuccess();
}
- ClusterSettingEntity curEntity =
metaDataCache.get(entity.getRecordKey());
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ ClusterSettingEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.getBrokerPort(), entity.getBrokerTLSPort(),
+ entity.getBrokerWebPort(), entity.getMaxMsgSizeInMB(),
+ entity.getQryPriorityId(), entity.enableFlowCtrl(),
+ entity.getGloFlowCtrlRuleCnt(),
entity.getGloFlowCtrlRuleInfo(),
+ entity.getClsDefTopicProps())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- "The cluster configure have not changed!");
+ "Cluster configure not changed!");
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- metaDataCache.put(entity.getRecordKey(), entity);
+ // Check whether the configured ports conflict in the record
+ if (WebParameterUtils.isConflictedPortsSet(newEntity.getBrokerPort(),
+ newEntity.getBrokerTLSPort(), newEntity.getBrokerWebPort(),
+ strBuff, result)) {
+ return result.isSuccess();
+ }
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ metaDataCache.put(newEntity.getRecordKey(), newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -103,7 +132,7 @@ public abstract class AbsClusterConfigMapperImpl implements
ClusterConfigMapper
*
* @param entity need added or updated entity
*/
- protected void addOrUpdCacheRecord(ClusterSettingEntity entity) {
+ protected void putRecord2Caches(ClusterSettingEntity entity) {
metaDataCache.put(entity.getRecordKey(), entity);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsConsumeCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsConsumeCtrlMapperImpl.java
index 4235080..f2a03d3 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsConsumeCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsConsumeCtrlMapperImpl.java
@@ -52,18 +52,18 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
@Override
public boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
StringBuilder strBuff,
ProcessResult result) {
- GroupConsumeCtrlEntity curEntity =
- grpConsumeCtrlCache.get(entity.getRecordKey());
+ // Checks whether the record already exists
+ GroupConsumeCtrlEntity curEntity =
grpConsumeCtrlCache.get(entity.getRecordKey());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The group consume configure
").append(entity.getRecordKey())
- .append(" already exists, please delete it first!")
- .toString());
+ strBuff.append("Existed record found for
groupName-topicName(")
+
.append(entity.getRecordKey()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ putRecord2Caches(entity);
}
return result.isSuccess();
}
@@ -71,26 +71,28 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
@Override
public boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
StringBuilder strBuff,
ProcessResult result) {
- GroupConsumeCtrlEntity curEntity =
- grpConsumeCtrlCache.get(entity.getRecordKey());
+ // Checks whether the record already exists
+ GroupConsumeCtrlEntity curEntity =
grpConsumeCtrlCache.get(entity.getRecordKey());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- strBuff.append("The group consume configure
").append(entity.getRecordKey())
- .append(" is not exists, please add it first!")
- .toString());
+ strBuff.append("Not found consume control for through
groupName-topicName(")
+
.append(entity.getRecordKey()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ GroupConsumeCtrlEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.isEnableConsume(), entity.getDisableReason(),
+ entity.isEnableFilterConsume(), entity.getFilterCondStr())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- strBuff.append("The group consume configure
").append(entity.getRecordKey())
- .append(" have not changed, please confirm it
first!")
- .toString());
- strBuff.delete(0, strBuff.length());
+ "Consume control configure not changed!");
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ putRecord2Caches(newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -121,7 +123,7 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
return true;
}
delConfigFromPersistent(recordKey, strBuff);
- delCacheRecord(recordKey);
+ delRecordFromCaches(recordKey);
result.setSuccResult(curEntity);
return true;
}
@@ -207,8 +209,8 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
@Override
public GroupConsumeCtrlEntity getConsumeCtrlByGroupAndTopic(
String groupName, String topicName) {
- String recKey = KeyBuilderUtils.buildGroupTopicRecKey(groupName,
topicName);
- return grpConsumeCtrlCache.get(recKey);
+ return grpConsumeCtrlCache.get(
+ KeyBuilderUtils.buildGroupTopicRecKey(groupName, topicName));
}
@Override
@@ -244,8 +246,7 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
}
@Override
- public List<GroupConsumeCtrlEntity> getGroupConsumeCtrlConf(
- GroupConsumeCtrlEntity qryEntity) {
+ public List<GroupConsumeCtrlEntity>
getGroupConsumeCtrlConf(GroupConsumeCtrlEntity qryEntity) {
List<GroupConsumeCtrlEntity> retEntities = new ArrayList<>();
if (qryEntity == null) {
retEntities.addAll(grpConsumeCtrlCache.values());
@@ -325,7 +326,7 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
*
* @param entity need added or updated entity
*/
- protected void addOrUpdCacheRecord(GroupConsumeCtrlEntity entity) {
+ protected void putRecord2Caches(GroupConsumeCtrlEntity entity) {
grpConsumeCtrlCache.put(entity.getRecordKey(), entity);
// add topic index map
ConcurrentHashSet<String> keySet =
@@ -370,7 +371,12 @@ public abstract class AbsConsumeCtrlMapperImpl implements
ConsumeCtrlMapper {
*/
protected abstract boolean delConfigFromPersistent(String recordKey,
StringBuilder strBuff);
- private void delCacheRecord(String recordKey) {
+ /**
+ * Delete the cached record
+ *
+ * @param recordKey the record key to be deleted
+ */
+ private void delRecordFromCaches(String recordKey) {
GroupConsumeCtrlEntity curEntity =
grpConsumeCtrlCache.remove(recordKey);
if (curEntity == null) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsGroupResCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsGroupResCtrlMapperImpl.java
index ebe1892..25fbe43 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsGroupResCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsGroupResCtrlMapperImpl.java
@@ -41,16 +41,16 @@ public abstract class AbsGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
@Override
public boolean addGroupResCtrlConf(GroupResCtrlEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- GroupResCtrlEntity curEntity =
- groupBaseCtrlCache.get(entity.getGroupName());
+ // Checks whether the record already exists
+ GroupResCtrlEntity curEntity =
groupBaseCtrlCache.get(entity.getGroupName());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The group control configure
").append(entity.getGroupName())
- .append(" already exists, please delete it first!")
- .toString());
+ strBuff.append("Existed record found for groupName(")
+
.append(entity.getGroupName()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
groupBaseCtrlCache.put(entity.getGroupName(), entity);
}
@@ -60,26 +60,29 @@ public abstract class AbsGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
@Override
public boolean updGroupResCtrlConf(GroupResCtrlEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- GroupResCtrlEntity curEntity =
- groupBaseCtrlCache.get(entity.getGroupName());
+ // Checks whether the record already exists
+ GroupResCtrlEntity curEntity =
groupBaseCtrlCache.get(entity.getGroupName());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- strBuff.append("The group control configure
").append(entity.getGroupName())
- .append(" is not exists, please add it first!")
- .toString());
+ strBuff.append("Not found group control configure for
groupName(")
+
.append(entity.getGroupName()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ GroupResCtrlEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.isEnableResCheck(), entity.getAllowedBrokerClientRate(),
+ entity.getQryPriorityId(), entity.isFlowCtrlEnable(),
+ entity.getRuleCnt(), entity.getFlowCtrlInfo())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- strBuff.append("The group control configure
").append(entity.getGroupName())
- .append(" have not changed, please confirm it
first!")
- .toString());
- strBuff.delete(0, strBuff.length());
+ "Group control configure not changed!");
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- groupBaseCtrlCache.put(entity.getGroupName(), entity);
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ groupBaseCtrlCache.put(newEntity.getGroupName(), newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -138,9 +141,9 @@ public abstract class AbsGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
/**
* Add or update a record
*
- * @param entity need added or updated entity
+ * @param entity the entity to be added or updated
*/
- protected void addOrUpdCacheRecord(GroupResCtrlEntity entity) {
+ protected void putRecord2Caches(GroupResCtrlEntity entity) {
groupBaseCtrlCache.put(entity.getGroupName(), entity);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicCtrlMapperImpl.java
index e95fdf5..9504ba1 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicCtrlMapperImpl.java
@@ -45,16 +45,16 @@ public abstract class AbsTopicCtrlMapperImpl implements
TopicCtrlMapper {
@Override
public boolean addTopicCtrlConf(TopicCtrlEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- TopicCtrlEntity curEntity =
- topicCtrlCache.get(entity.getTopicName());
+ // Checks whether the record already exists
+ TopicCtrlEntity curEntity = topicCtrlCache.get(entity.getTopicName());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The topic control configure
").append(entity.getTopicName())
- .append(" already exists, please delete it first!")
- .toString());
+ strBuff.append("Existed record found for topicName(")
+
.append(entity.getTopicName()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
topicCtrlCache.put(entity.getTopicName(), entity);
}
@@ -64,26 +64,28 @@ public abstract class AbsTopicCtrlMapperImpl implements
TopicCtrlMapper {
@Override
public boolean updTopicCtrlConf(TopicCtrlEntity entity,
StringBuilder strBuff, ProcessResult
result) {
- TopicCtrlEntity curEntity =
- topicCtrlCache.get(entity.getTopicName());
+ // Checks whether the record already exists
+ TopicCtrlEntity curEntity = topicCtrlCache.get(entity.getTopicName());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- strBuff.append("The topic control configure
").append(entity.getTopicName())
- .append(" is not exists, please add it first!")
- .toString());
+ strBuff.append("Not found topic control configure for
topicName(")
+
.append(entity.getTopicName()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ TopicCtrlEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.getTopicId(), entity.getMaxMsgSizeInMB(),
+ entity.isAuthCtrlEnable())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- strBuff.append("The topic control configure
").append(entity.getTopicName())
- .append(" have not changed, please confirm it
first!")
- .toString());
- strBuff.delete(0, strBuff.length());
+ "Topic control configure not changed!");
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- topicCtrlCache.put(entity.getTopicName(), entity);
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ topicCtrlCache.put(newEntity.getTopicName(), newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -155,9 +157,9 @@ public abstract class AbsTopicCtrlMapperImpl implements
TopicCtrlMapper {
/**
* Add or update a record
*
- * @param entity need added or updated entity
+ * @param entity the entity need to added or updated
*/
- protected void addOrUpdCacheRecord(TopicCtrlEntity entity) {
+ protected void putRecord2Caches(TopicCtrlEntity entity) {
topicCtrlCache.put(entity.getTopicName(), entity);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicDeployMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicDeployMapperImpl.java
index 0a00e93..0951c16 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicDeployMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/AbsTopicDeployMapperImpl.java
@@ -25,9 +25,12 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
+import org.apache.inlong.tubemq.corebase.TBaseConstants;
import org.apache.inlong.tubemq.corebase.rv.ProcessResult;
import org.apache.inlong.tubemq.corebase.utils.ConcurrentHashSet;
import org.apache.inlong.tubemq.corebase.utils.KeyBuilderUtils;
+import org.apache.inlong.tubemq.server.common.TServerConstants;
+import org.apache.inlong.tubemq.server.common.statusdef.TopicStatus;
import org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.TopicDeployEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.mapper.TopicDeployMapper;
@@ -54,18 +57,40 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
@Override
public boolean addTopicConf(TopicDeployEntity entity,
StringBuilder strBuff, ProcessResult result) {
+ // Checks whether the record already exists
TopicDeployEntity curEntity =
topicConfCache.get(entity.getRecordKey());
if (curEntity != null) {
- result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- strBuff.append("The topic deploy configure
").append(entity.getRecordKey())
- .append(" already exists, please delete it first!")
- .toString());
+ if (curEntity.isValidTopicStatus()) {
+ result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
+ strBuff.append("Existed record found for
brokerId-topicName(")
+
.append(curEntity.getRecordKey()).append(")!").toString());
+ } else {
+ result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
+ strBuff.append("Softly deleted record found for
brokerId-topicName(")
+ .append(curEntity.getRecordKey())
+ .append("), please resume or remove it
first!").toString());
+ }
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ // valid whether system topic
+ if (!validSysTopicConfigure(entity, strBuff, result)) {
+ return result.isSuccess();
+ }
+ // check deploy status if still accept publish and subscribe
+ if (!entity.isValidTopicStatus()
+ && (entity.isAcceptPublish() || entity.isAcceptSubscribe())) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" when add brokerId-topicName(")
+ .append(entity.getRecordKey()).append(")
record!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
+ // Store data to persistent
if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ putRecord2Caches(entity);
}
return result.isSuccess();
}
@@ -73,26 +98,38 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
@Override
public boolean updTopicConf(TopicDeployEntity entity,
StringBuilder strBuff, ProcessResult result) {
+ // Checks whether the record already exists
TopicDeployEntity curEntity =
topicConfCache.get(entity.getRecordKey());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- strBuff.append("The topic deploy configure
").append(entity.getRecordKey())
- .append(" is not exists, please add it first!")
- .toString());
+ strBuff.append("Not found topic deploy configure for
brokerId-topicName(")
+
.append(entity.getRecordKey()).append(")!").toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (curEntity.equals(entity)) {
+ // Build the entity that need to be updated
+ TopicDeployEntity newEntity = curEntity.clone();
+ newEntity.updBaseModifyInfo(entity);
+ if (!newEntity.updModifyInfo(entity.getDataVerId(),
+ entity.getTopicId(), entity.getBrokerPort(),
+ entity.getBrokerIp(), entity.getDeployStatus(),
+ entity.getTopicProps())) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- strBuff.append("The topic deploy configure
").append(entity.getRecordKey())
- .append(" have not changed, please confirm it
first!")
- .toString());
- strBuff.delete(0, strBuff.length());
+ "Topic deploy configure not changed!");
return result.isSuccess();
}
- if (putConfig2Persistent(entity, strBuff, result)) {
- addOrUpdCacheRecord(entity);
+ // valid whether system topic
+ if (!validSysTopicConfigure(newEntity, strBuff, result)) {
+ return result.isSuccess();
+ }
+ // check deploy status
+ if (isIllegalValuesChange(newEntity, curEntity, strBuff, result)) {
+ return result.isSuccess();
+ }
+ // Store data to persistent
+ if (putConfig2Persistent(newEntity, strBuff, result)) {
+ putRecord2Caches(newEntity);
result.setSuccResult(curEntity);
}
return result.isSuccess();
@@ -106,8 +143,17 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
result.setSuccResult(null);
return result.isSuccess();
}
+ // check deploy status if still accept publish and subscribe
+ if (curEntity.isAcceptPublish() || curEntity.isAcceptSubscribe()) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" before delete brokerId-topicName(")
+ .append(curEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
delConfigFromPersistent(recordKey, strBuff);
- delCacheRecord(recordKey);
+ delRecordFromCaches(recordKey);
result.setSuccResult(curEntity);
return result.isSuccess();
}
@@ -116,13 +162,30 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
public boolean delTopicConfByBrokerId(Integer brokerId, StringBuilder
strBuff, ProcessResult result) {
ConcurrentHashSet<String> recordKeySet =
brokerIdCacheIndex.get(brokerId);
- if (recordKeySet == null) {
+ if (recordKeySet == null || recordKeySet.isEmpty()) {
result.setSuccResult(null);
return result.isSuccess();
}
+ // check deploy status if still accept publish and subscribe
+ TopicDeployEntity curEntity;
+ for (String recordKey : recordKeySet) {
+ curEntity = topicConfCache.get(recordKey);
+ if (curEntity == null) {
+ continue;
+ }
+ if (curEntity.isAcceptPublish() || curEntity.isAcceptSubscribe()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" before delete brokerId-topicName(")
+ .append(curEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ }
+ // delete records
for (String recordKey : recordKeySet) {
delConfigFromPersistent(recordKey, strBuff);
- delCacheRecord(recordKey);
+ delRecordFromCaches(recordKey);
}
result.setSuccResult(null);
return result.isSuccess();
@@ -366,7 +429,7 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
*
* @param entity need added or updated entity
*/
- protected void addOrUpdCacheRecord(TopicDeployEntity entity) {
+ protected void putRecord2Caches(TopicDeployEntity entity) {
topicConfCache.put(entity.getRecordKey(), entity);
// add topic index map
ConcurrentHashSet<String> keySet =
@@ -421,7 +484,7 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
*/
protected abstract boolean delConfigFromPersistent(String recordKey,
StringBuilder strBuff);
- private void delCacheRecord(String recordKey) {
+ private void delRecordFromCaches(String recordKey) {
TopicDeployEntity curEntity =
topicConfCache.remove(recordKey);
if (curEntity == null) {
@@ -505,4 +568,164 @@ public abstract class AbsTopicDeployMapperImpl implements
TopicDeployMapper {
}
return matchedKeySet;
}
+
+ /**
+ * Check whether the change of deploy values is illegal
+ * Attention, the newEntity and newEntity must not equal
+ *
+ * @param newEntity the entity to be updated
+ * @param curEntity the current entity
+ * @param strBuff string buffer
+ * @param result check result of parameter value
+ * @return true for illegal, false for legal
+ */
+ private boolean isIllegalValuesChange(TopicDeployEntity newEntity,
+ TopicDeployEntity curEntity,
+ StringBuilder strBuff,
+ ProcessResult result) {
+ // check if shrink data store block
+ if (newEntity.getNumPartitions() != TBaseConstants.META_VALUE_UNDEFINED
+ && newEntity.getNumPartitions() <
curEntity.getNumPartitions()) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
+ strBuff.append("Partition number less than before,")
+ .append(" new value is
").append(newEntity.getNumPartitions())
+ .append(", current value is
").append(curEntity.getNumPartitions())
+ .append("in
brokerId-topicName(").append(curEntity.getRecordKey())
+ .append(") record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ if (newEntity.getNumTopicStores() !=
TBaseConstants.META_VALUE_UNDEFINED
+ && newEntity.getNumTopicStores() <
curEntity.getNumTopicStores()) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
+ strBuff.append("TopicStores number less than before,")
+ .append(" new value is
").append(newEntity.getNumTopicStores())
+ .append(", current value is
").append(curEntity.getNumTopicStores())
+ .append("in
brokerId-topicName(").append(curEntity.getRecordKey())
+ .append(") record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ // check whether the deploy status is equal
+ if (newEntity.getTopicStatus() == curEntity.getTopicStatus()) {
+ if (!newEntity.isValidTopicStatus()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("Softly deleted record cannot be
changed,")
+ .append(" please resume or hard remove for
brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ return false;
+ }
+ // check deploy status case from valid to invalid
+ if (curEntity.isValidTopicStatus() && !newEntity.isValidTopicStatus())
{
+ if (curEntity.isAcceptPublish()
+ || curEntity.isAcceptSubscribe()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" before change status of
brokerId-topicName(")
+ .append(curEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ if (newEntity.getTopicStatus().getCode()
+ > TopicStatus.STATUS_TOPIC_SOFT_DELETE.getCode()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("Please softly deleted the
brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record first!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ return false;
+ }
+ // check deploy status case from invalid to invalid
+ if (!curEntity.isValidTopicStatus() &&
!newEntity.isValidTopicStatus()) {
+ if (!((curEntity.getTopicStatus() ==
TopicStatus.STATUS_TOPIC_SOFT_DELETE
+ && newEntity.getTopicStatus() ==
TopicStatus.STATUS_TOPIC_SOFT_REMOVE)
+ || (curEntity.getTopicStatus() ==
TopicStatus.STATUS_TOPIC_SOFT_REMOVE
+ && newEntity.getTopicStatus() ==
TopicStatus.STATUS_TOPIC_HARD_REMOVE))) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("Illegal transfer status from ")
+
.append(curEntity.getTopicStatus().getDescription())
+ .append(" to
").append(newEntity.getTopicStatus().getDescription())
+ .append(" for the brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ if (newEntity.isAcceptPublish()
+ || newEntity.isAcceptSubscribe()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" before change status of
brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ return false;
+ }
+ // check deploy status case from invalid to valid
+ if (!curEntity.isValidTopicStatus() && newEntity.isValidTopicStatus())
{
+ if (curEntity.getTopicStatus() !=
TopicStatus.STATUS_TOPIC_SOFT_DELETE) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("Illegal transfer status from ")
+
.append(curEntity.getTopicStatus().getDescription())
+ .append(" to
").append(newEntity.getTopicStatus().getDescription())
+ .append(" for the brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ if (newEntity.isAcceptPublish()
+ || newEntity.isAcceptSubscribe()) {
+
result.setFailResult(DataOpErrCode.DERR_ILLEGAL_STATUS.getCode(),
+ strBuff.append("The values of acceptPublish and
acceptSubscribe must be false")
+ .append(" before change status of
brokerId-topicName(")
+ .append(newEntity.getRecordKey()).append(")
record!").toString());
+ strBuff.delete(0, strBuff.length());
+ return !result.isSuccess();
+ }
+ return false;
+ }
+ return false;
+ }
+
+ /**
+ * Verify the validity of the configuration value for the system topic
+ *
+ * @param deployEntity the topic configuration that needs to be added or
updated
+ * @param strBuff the print info string buffer
+ * @param result the process result return
+ * @return true if success otherwise false
+ */
+ private boolean validSysTopicConfigure(TopicDeployEntity deployEntity,
+ StringBuilder strBuff,
ProcessResult result) {
+ if
(!TServerConstants.OFFSET_HISTORY_NAME.equals(deployEntity.getTopicName())) {
+ return true;
+ }
+ if (deployEntity.getNumTopicStores()
+ != TServerConstants.OFFSET_HISTORY_NUMSTORES) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
+ strBuff.append("For system topic")
+ .append(TServerConstants.OFFSET_HISTORY_NAME)
+ .append(", the TopicStores value(")
+ .append(TServerConstants.OFFSET_HISTORY_NUMSTORES)
+ .append(") cannot be changed!").toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ if (deployEntity.getNumPartitions()
+ != TServerConstants.OFFSET_HISTORY_NUMPARTS) {
+ result.setFailResult(DataOpErrCode.DERR_ILLEGAL_VALUE.getCode(),
+ strBuff.append("For system topic")
+ .append(TServerConstants.OFFSET_HISTORY_NAME)
+ .append(", the Partition value(")
+ .append(TServerConstants.OFFSET_HISTORY_NUMPARTS)
+ .append(") cannot be changed!").toString());
+ strBuff.delete(0, strBuff.length());
+ return result.isSuccess();
+ }
+ return true;
+ }
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbBrokerConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbBrokerConfigMapperImpl.java
index fba7ea6..0ecf7bf 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbBrokerConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbBrokerConfigMapperImpl.java
@@ -71,7 +71,7 @@ public class BdbBrokerConfigMapperImpl extends
AbsBrokerConfigMapperImpl {
logger.warn("[BDB Impl] found Null data while loading
broker configure!");
continue;
}
- addOrUpdCacheRecord(new BrokerConfEntity(bdbEntity));
+ putRecord2Caches(new BrokerConfEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
@@ -88,7 +88,8 @@ public class BdbBrokerConfigMapperImpl extends
AbsBrokerConfigMapperImpl {
}
protected boolean putConfig2Persistent(BrokerConfEntity entity,
- StringBuilder strBuff,
ProcessResult result) {
+ StringBuilder strBuff,
+ ProcessResult result) {
BdbBrokerConfEntity bdbEntity =
entity.buildBdbBrokerConfEntity();
try {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbClusterConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbClusterConfigMapperImpl.java
index 26489a8..569bc15 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbClusterConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbClusterConfigMapperImpl.java
@@ -73,7 +73,7 @@ public class BdbClusterConfigMapperImpl extends
AbsClusterConfigMapperImpl {
logger.warn("[BDB Impl] found Null data while loading
cluster configure!");
continue;
}
- addOrUpdCacheRecord(new ClusterSettingEntity(bdbEntity));
+ putRecord2Caches(new ClusterSettingEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbConsumeCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbConsumeCtrlMapperImpl.java
index 722f19f..8343077 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbConsumeCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbConsumeCtrlMapperImpl.java
@@ -69,7 +69,7 @@ public class BdbConsumeCtrlMapperImpl extends
AbsConsumeCtrlMapperImpl {
logger.warn("[BDB Impl] found Null data while loading
consume control configure!");
continue;
}
- addOrUpdCacheRecord(new GroupConsumeCtrlEntity(bdbEntity));
+ putRecord2Caches(new GroupConsumeCtrlEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
@@ -94,7 +94,7 @@ public class BdbConsumeCtrlMapperImpl extends
AbsConsumeCtrlMapperImpl {
} catch (Throwable e) {
logger.error("[BDB Impl] put consume control configure failure ",
e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- strBuff.append("Put consume control configure failure: ")
+ strBuff.append("Put consume control configure failure: ")
.append(e.getMessage()).toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupResCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupResCtrlMapperImpl.java
index 3a29011..9ecdfab 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupResCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupResCtrlMapperImpl.java
@@ -69,7 +69,7 @@ public class BdbGroupResCtrlMapperImpl extends
AbsGroupResCtrlMapperImpl {
logger.warn("[BDB Impl] null data while loading group
control configure!");
continue;
}
- addOrUpdCacheRecord(new GroupResCtrlEntity(bdbEntity));
+ putRecord2Caches(new GroupResCtrlEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/BdbMetaStoreServiceImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
similarity index 93%
rename from
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/BdbMetaStoreServiceImpl.java
rename to
inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
index cce7eba..b2917b7 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/BdbMetaStoreServiceImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbMetaStoreServiceImpl.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.inlong.tubemq.server.master.metamanage.metastore;
+package
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl;
import com.sleepycat.je.DatabaseException;
import com.sleepycat.je.Durability;
@@ -58,10 +58,12 @@ 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.keepalive.AliveObserver;
+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;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.GroupConsumeCtrlEntity;
@@ -74,12 +76,6 @@ 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.metamanage.metastore.impl.bdbimpl.BdbBrokerConfigMapperImpl;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbClusterConfigMapperImpl;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbConsumeCtrlMapperImpl;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbGroupResCtrlMapperImpl;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbTopicCtrlMapperImpl;
-import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.bdbimpl.BdbTopicDeployMapperImpl;
import org.apache.inlong.tubemq.server.master.utils.BdbStoreSamplePrint;
import org.apache.inlong.tubemq.server.master.web.model.ClusterGroupVO;
import org.apache.inlong.tubemq.server.master.web.model.ClusterNodeVO;
@@ -144,11 +140,10 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
// group consume control configure
private ConsumeCtrlMapper consumeCtrlMapper;
- public BdbMetaStoreServiceImpl(String nodeHost, String metaDataPath,
- MasterReplicationConfig replicationConfig) {
- this.nodeHost = nodeHost;
- this.metaDataPath = metaDataPath;
- this.replicationConfig = replicationConfig;
+ public BdbMetaStoreServiceImpl(MasterConfig masterConfig) {
+ this.nodeHost = masterConfig.getHostName();
+ this.metaDataPath = masterConfig.getMetaDataPath();
+ this.replicationConfig = masterConfig.getReplicationConfig();
// build replicationGroupAdmin info
Set<InetSocketAddress> helpers = new HashSet<>();
for (int i = 1; i <= 3; i++) {
@@ -223,10 +218,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (clusterConfigMapper.addClusterConfig(entity, strBuff, result)) {
- strBuff.append("[addClusterConfig], ")
- .append(entity.getCreateUser())
- .append(" added cluster setting record :")
- .append(entity.toString());
+ strBuff.append("[addClusterConfig],
").append(entity.getCreateUser())
+ .append(" added cluster setting record :").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addClusterConfig], ")
@@ -281,7 +274,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
(ClusterSettingEntity) result.getRetData();
if (entity != null) {
strBuff.append("[delClusterConfig], ").append(operator)
- .append(" deleted cluster setting record
:").append(entity.toString());
+ .append(" deleted cluster setting record
:").append(entity);
logger.info(strBuff.toString());
}
} else {
@@ -303,10 +296,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (brokerConfigMapper.addBrokerConf(entity, strBuff, result)) {
- strBuff.append("[addBrokerConf], ")
- .append(entity.getCreateUser())
- .append(" added broker configure record :")
- .append(entity.toString());
+ strBuff.append("[addBrokerConf], ").append(entity.getCreateUser())
+ .append(" added broker configure record :").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addBrokerConf], ")
@@ -357,8 +348,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
BrokerConfEntity entity = (BrokerConfEntity) result.getRetData();
if (entity != null) {
strBuffer.append("[delBrokerConf], ").append(operator)
- .append(" deleted broker configure record :")
- .append(entity.toString());
+ .append(" deleted broker configure record
:").append(entity);
logger.info(strBuffer.toString());
}
} else {
@@ -402,10 +392,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (topicDeployMapper.addTopicConf(entity, strBuff, result)) {
- strBuff.append("[addTopicConf], ")
- .append(entity.getCreateUser())
- .append(" added topic configure record :")
- .append(entity.toString());
+ strBuff.append("[addTopicConf], ").append(entity.getCreateUser())
+ .append(" added topic configure record :").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addTopicConf], ")
@@ -457,8 +445,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
(GroupResCtrlEntity) result.getRetData();
if (entity != null) {
strBuff.append("[delTopicConf], ").append(operator)
- .append(" deleted topic configure record :")
- .append(entity.toString());
+ .append(" deleted topic configure record
:").append(entity);
logger.info(strBuff.toString());
}
} else {
@@ -566,10 +553,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (topicCtrlMapper.addTopicCtrlConf(entity, strBuff, result)) {
- strBuff.append("[addTopicCtrlConf], ")
- .append(entity.getCreateUser())
- .append(" added topic control record :")
- .append(entity.toString());
+ strBuff.append("[addTopicCtrlConf],
").append(entity.getCreateUser())
+ .append(" added topic control record :").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addTopicCtrlConf], ")
@@ -620,7 +605,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
(TopicCtrlEntity) result.getRetData();
if (entity != null) {
strBuff.append("[delTopicCtrlConf], ").append(operator)
- .append(" deleted topic control record
:").append(entity.toString());
+ .append(" deleted topic control record
:").append(entity);
logger.info(strBuff.toString());
}
} else {
@@ -658,10 +643,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (groupResCtrlMapper.addGroupResCtrlConf(entity, strBuff, result)) {
- strBuff.append("[addGroupResCtrlConf], ")
- .append(entity.getCreateUser())
- .append(" added group resource control record :")
- .append(entity.toString());
+ strBuff.append("[addGroupResCtrlConf],
").append(entity.getCreateUser())
+ .append(" added group resource control record
:").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addGroupResCtrlConf], ")
@@ -713,8 +696,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
(GroupResCtrlEntity) result.getRetData();
if (entity != null) {
strBuff.append("[delGroupResCtrlConf], ").append(operator)
- .append(" deleted group resource control record :")
- .append(entity.toString());
+ .append(" deleted group resource control record
:").append(entity);
logger.info(strBuff.toString());
}
} else {
@@ -746,10 +728,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (consumeCtrlMapper.addGroupConsumeCtrlConf(entity, strBuff,
result)) {
- strBuff.append("[addGroupConsumeCtrlConf], ")
- .append(entity.getCreateUser())
- .append(" added group consume control record :")
- .append(entity.toString());
+ strBuff.append("[addGroupConsumeCtrlConf],
").append(entity.getCreateUser())
+ .append(" added group consume control record
:").append(entity);
logger.info(strBuff.toString());
} else {
strBuff.append("[addGroupConsumeCtrlConf], ")
@@ -979,7 +959,6 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
}
} catch (Throwable e) {
logger.error("[BDB Impl] Get nodeState Throwable error", e);
- continue;
}
}
return null;
@@ -1216,7 +1195,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
* Initialize configuration for BDB-JE replication environment.
*
* */
- private void initEnvConfig() throws InterruptedException {
+ private void initEnvConfig() {
// Set envHome and generate a ReplicationConfig. Note that
ReplicationConfig and
// EnvironmentConfig values could all be specified in the
je.properties file,
@@ -1271,7 +1250,6 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
*/
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++) {
@@ -1305,11 +1283,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
}
}
// Failed despite retries.
- if (exception != null) {
- throw exception;
- }
- // Don't expect to get here.
- throw new IllegalStateException("Failed despite retries");
+ throw exception;
}
/* initial metadata */
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicCtrlMapperImpl.java
index 7cb983a..714b7a5 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicCtrlMapperImpl.java
@@ -71,7 +71,7 @@ public class BdbTopicCtrlMapperImpl extends
AbsTopicCtrlMapperImpl {
logger.warn("[BDB Impl] found Null data while loading
topic control!");
continue;
}
- addOrUpdCacheRecord(new TopicCtrlEntity(bdbEntity));
+ putRecord2Caches(new TopicCtrlEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
@@ -96,7 +96,7 @@ public class BdbTopicCtrlMapperImpl extends
AbsTopicCtrlMapperImpl {
} catch (Throwable e) {
logger.error("[BDB Impl] put topic control failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- strBuff.append("Put topic control failure: ")
+ strBuff.append("Put topic control configure failure: ")
.append(e.getMessage()).toString());
strBuff.delete(0, strBuff.length());
return result.isSuccess();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicDeployMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicDeployMapperImpl.java
index b4b05b9..8cff9a2 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicDeployMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbTopicDeployMapperImpl.java
@@ -69,7 +69,7 @@ public class BdbTopicDeployMapperImpl extends
AbsTopicDeployMapperImpl {
logger.warn("[BDB Impl] found Null data while loading
topic deploy configure!");
continue;
}
- addOrUpdCacheRecord(new TopicDeployEntity(bdbEntity));
+ putRecord2Caches(new TopicDeployEntity(bdbEntity));
totalCnt++;
}
} catch (Exception e) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/TZKNodeKeys.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/TZKNodeKeys.java
index f49fb3e..a5b21d7 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/TZKNodeKeys.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/TZKNodeKeys.java
@@ -18,6 +18,7 @@
package
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.zkimpl;
public class TZKNodeKeys {
+ public static final String ZK_BRANCH_HA = "masterHA";
public static final String ZK_BRANCH_META_DATA = "metaData";
public static final String ZK_LEAF_CLUSTER_CONFIG = "clusterConfig";
public static final String ZK_LEAF_BROKER_CONFIG = "brokerConfig";
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKBrokerConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKBrokerConfigMapperImpl.java
index a9d4233..703c13b 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKBrokerConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKBrokerConfigMapperImpl.java
@@ -32,12 +32,8 @@ import
org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.BrokerConfEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.AbsBrokerConfigMapperImpl;
import org.apache.zookeeper.KeeperException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
public class ZKBrokerConfigMapperImpl extends AbsBrokerConfigMapperImpl {
- private static final Logger logger =
- LoggerFactory.getLogger(ZKBrokerConfigMapperImpl.class);
private final ZooKeeperWatcher zkWatcher;
private final String brokerCfgRootDir;
@@ -89,7 +85,7 @@ public class ZKBrokerConfigMapperImpl extends
AbsBrokerConfigMapperImpl {
if (confStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(confStr, type));
+ putRecord2Caches(gson.fromJson(confStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKClusterConfigMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKClusterConfigMapperImpl.java
index 8feeb64..377d70c 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKClusterConfigMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKClusterConfigMapperImpl.java
@@ -89,7 +89,7 @@ public class ZKClusterConfigMapperImpl extends
AbsClusterConfigMapperImpl {
if (clrConfigureStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(clrConfigureStr, type));
+ putRecord2Caches(gson.fromJson(clrConfigureStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKConsumeCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKConsumeCtrlMapperImpl.java
index 8850caa..a2fbd2a 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKConsumeCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKConsumeCtrlMapperImpl.java
@@ -32,12 +32,8 @@ import
org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.GroupConsumeCtrlEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.AbsConsumeCtrlMapperImpl;
import org.apache.zookeeper.KeeperException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
public class ZKConsumeCtrlMapperImpl extends AbsConsumeCtrlMapperImpl {
- private static final Logger logger =
- LoggerFactory.getLogger(ZKConsumeCtrlMapperImpl.class);
private final ZooKeeperWatcher zkWatcher;
private final String csmCtrlRootDir;
@@ -85,7 +81,7 @@ public class ZKConsumeCtrlMapperImpl extends
AbsConsumeCtrlMapperImpl {
if (recordStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(recordStr, type));
+ putRecord2Caches(gson.fromJson(recordStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKGroupResCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKGroupResCtrlMapperImpl.java
index 3bb8475..27d29f7 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKGroupResCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKGroupResCtrlMapperImpl.java
@@ -32,12 +32,8 @@ import
org.apache.inlong.tubemq.server.master.metamanage.DataOpErrCode;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.GroupResCtrlEntity;
import
org.apache.inlong.tubemq.server.master.metamanage.metastore.impl.AbsGroupResCtrlMapperImpl;
import org.apache.zookeeper.KeeperException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
public class ZKGroupResCtrlMapperImpl extends AbsGroupResCtrlMapperImpl {
- private static final Logger logger =
- LoggerFactory.getLogger(ZKGroupResCtrlMapperImpl.class);
private final ZooKeeperWatcher zkWatcher;
private final String groupCtrlRootDir;
@@ -88,7 +84,7 @@ public class ZKGroupResCtrlMapperImpl extends
AbsGroupResCtrlMapperImpl {
if (recordStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(recordStr, type));
+ putRecord2Caches(gson.fromJson(recordStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicCtrlMapperImpl.java
index fdd8e99..7a4febb 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicCtrlMapperImpl.java
@@ -84,7 +84,7 @@ public class ZKTopicCtrlMapperImpl extends
AbsTopicCtrlMapperImpl {
if (recordStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(recordStr, type));
+ putRecord2Caches(gson.fromJson(recordStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicDeployMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicDeployMapperImpl.java
index 5e3ae30..d8aedd5 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicDeployMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/zkimpl/ZKTopicDeployMapperImpl.java
@@ -85,7 +85,7 @@ public class ZKTopicDeployMapperImpl extends
AbsTopicDeployMapperImpl {
if (recordStr == null) {
continue;
}
- addOrUpdCacheRecord(gson.fromJson(recordStr, type));
+ putRecord2Caches(gson.fromJson(recordStr, type));
totalCnt++;
}
logger.info(strBuff.append("[ZK Impl] loaded ").append(totalCnt)