This is an automated email from the ASF dual-hosted git repository.
aloyszhang 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 b4b3054 [INLONG-2980][TubeMQ] Modify the code style problems of the
metadata classes (#2984)
b4b3054 is described below
commit b4b30546bf60089791c34a1d4c1ff296ef00d6e7
Author: gosonzhang <[email protected]>
AuthorDate: Mon Mar 7 20:32:47 2022 +0800
[INLONG-2980][TubeMQ] Modify the code style problems of the metadata
classes (#2984)
---
.../metastore/BdbMetaStoreServiceImpl.java | 374 ++++++++++-----------
.../metamanage/metastore/MetaStoreService.java | 115 +++----
.../metastore/dao/entity/BaseEntity.java | 6 +-
.../metastore/dao/entity/BrokerConfEntity.java | 11 +-
.../metastore/dao/entity/ClusterSettingEntity.java | 5 +-
.../dao/entity/GroupConsumeCtrlEntity.java | 2 +-
.../metastore/dao/entity/TopicDeployEntity.java | 2 +-
.../metastore/dao/mapper/BrokerConfigMapper.java | 6 +-
.../metastore/dao/mapper/ClusterConfigMapper.java | 6 +-
.../dao/mapper/GroupConsumeCtrlMapper.java | 6 +-
.../metastore/dao/mapper/GroupResCtrlMapper.java | 6 +-
.../metastore/dao/mapper/TopicCtrlMapper.java | 6 +-
.../metastore/dao/mapper/TopicDeployMapper.java | 6 +-
.../impl/bdbimpl/BdbBrokerConfigMapperImpl.java | 65 ++--
.../impl/bdbimpl/BdbClusterConfigMapperImpl.java | 31 +-
.../bdbimpl/BdbGroupConsumeCtrlMapperImpl.java | 56 +--
.../impl/bdbimpl/BdbGroupResCtrlMapperImpl.java | 42 +--
.../impl/bdbimpl/BdbTopicCtrlMapperImpl.java | 53 +--
.../impl/bdbimpl/BdbTopicDeployMapperImpl.java | 56 +--
.../nodemanage/nodebroker/DefBrokerRunManager.java | 9 +-
20 files changed, 424 insertions(+), 439 deletions(-)
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/BdbMetaStoreServiceImpl.java
index da09924..caaba10 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/BdbMetaStoreServiceImpl.java
@@ -217,53 +217,51 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
// cluster default configure api
@Override
public boolean addClusterConfig(ClusterSettingEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (clusterConfigMapper.addClusterConfig(entity, result)) {
- strBuffer.append("[addClusterConfig], ")
+ if (clusterConfigMapper.addClusterConfig(entity, strBuff, result)) {
+ strBuff.append("[addClusterConfig], ")
.append(entity.getCreateUser())
.append(" added cluster setting record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[addClusterConfig], ")
+ strBuff.append("[addClusterConfig], ")
.append("failure to add cluster setting record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updClusterConfig(ClusterSettingEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (clusterConfigMapper.updClusterConfig(entity, result)) {
+ if (clusterConfigMapper.updClusterConfig(entity, strBuff, result)) {
ClusterSettingEntity oldEntity =
(ClusterSettingEntity) result.getRetData();
ClusterSettingEntity curEntity =
clusterConfigMapper.getClusterConfig();
- strBuffer.append("[updClusterConfig], ")
+ strBuff.append("[updClusterConfig], ")
.append(entity.getModifyUser())
.append(" updated record from
:").append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[updClusterConfig], ")
+ strBuff.append("[updClusterConfig], ")
.append("failure to update record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -274,8 +272,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
@Override
public boolean delClusterConfig(String operator,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
if (!checkStoreStatus(true, result)) {
return false;
}
@@ -283,71 +280,69 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
ClusterSettingEntity entity =
(ClusterSettingEntity) result.getRetData();
if (entity != null) {
- strBuffer.append("[delClusterConfig], ").append(operator)
+ strBuff.append("[delClusterConfig], ").append(operator)
.append(" deleted cluster setting record
:").append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
}
} else {
- strBuffer.append("[delClusterConfig], ")
+ strBuff.append("[delClusterConfig], ")
.append("failure to delete cluster setting record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
// broker configure api
@Override
public boolean addBrokerConf(BrokerConfEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (brokerConfigMapper.addBrokerConf(entity, result)) {
- strBuffer.append("[addBrokerConf], ")
+ if (brokerConfigMapper.addBrokerConf(entity, strBuff, result)) {
+ strBuff.append("[addBrokerConf], ")
.append(entity.getCreateUser())
.append(" added broker configure record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[addBrokerConf], ")
+ strBuff.append("[addBrokerConf], ")
.append("failure to add broker configure record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updBrokerConf(BrokerConfEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (brokerConfigMapper.updBrokerConf(entity, result)) {
+ if (brokerConfigMapper.updBrokerConf(entity, strBuff, result)) {
BrokerConfEntity oldEntity =
(BrokerConfEntity) result.getRetData();
BrokerConfEntity curEntity =
brokerConfigMapper.getBrokerConfByBrokerId(entity.getBrokerId());
- strBuffer.append("[updBrokerConf], ")
+ strBuff.append("[updBrokerConf], ")
.append(entity.getModifyUser())
.append(" updated broker configure record from :")
.append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[updBrokerConf], ")
+ strBuff.append("[updBrokerConf], ")
.append("failure to update broker configure record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -362,7 +357,8 @@ 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.toString());
logger.info(strBuffer.toString());
}
} else {
@@ -400,62 +396,58 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
// topic configure api
@Override
public boolean addTopicConf(TopicDeployEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (topicDeployMapper.addTopicConf(entity, result)) {
- strBuffer.append("[addTopicConf], ")
+ if (topicDeployMapper.addTopicConf(entity, strBuff, result)) {
+ strBuff.append("[addTopicConf], ")
.append(entity.getCreateUser())
.append(" added topic configure record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[addTopicConf], ")
+ strBuff.append("[addTopicConf], ")
.append("failure to add topic configure record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updTopicConf(TopicDeployEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (topicDeployMapper.updTopicConf(entity, result)) {
+ if (topicDeployMapper.updTopicConf(entity, strBuff, result)) {
TopicDeployEntity oldEntity =
(TopicDeployEntity) result.getRetData();
TopicDeployEntity curEntity =
topicDeployMapper.getTopicConfByeRecKey(entity.getRecordKey());
- strBuffer.append("[updTopicConf], ")
+ strBuff.append("[updTopicConf], ")
.append(entity.getModifyUser())
.append(" updated record from :")
.append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[updTopicConf], ")
+ strBuff.append("[updTopicConf], ")
.append("failure to update topic configure record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delTopicConf(String operator,
- String recordKey,
- StringBuilder strBuffer,
- ProcessResult result) {
+ public boolean delTopicConf(String operator, String recordKey,
+ StringBuilder strBuff, ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
@@ -464,43 +456,41 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
GroupResCtrlEntity entity =
(GroupResCtrlEntity) result.getRetData();
if (entity != null) {
- strBuffer.append("[delTopicConf], ").append(operator)
+ strBuff.append("[delTopicConf], ").append(operator)
.append(" deleted topic configure record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
}
} else {
- strBuffer.append("[delTopicConf], ")
+ strBuff.append("[delTopicConf], ")
.append("failure to delete topic configure record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delTopicConfByBrokerId(String operator,
- int brokerId,
- StringBuilder strBuffer,
- ProcessResult result) {
+ public boolean delTopicConfByBrokerId(String operator, int brokerId,
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
if (topicDeployMapper.delTopicConfByBrokerId(brokerId, result)) {
- strBuffer.append("[delTopicConfByBrokerId], ")
+ strBuff.append("[delTopicConfByBrokerId], ")
.append(operator)
.append(" deleted topic deploy record :")
.append(brokerId);
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[delTopicConfByBrokerId], ")
+ strBuff.append("[delTopicConfByBrokerId], ")
.append("failure to delete topic deploy record : ")
.append(brokerId).append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -570,61 +560,57 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
// topic control api
@Override
public boolean addTopicCtrlConf(TopicCtrlEntity entity,
- StringBuilder sBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (topicCtrlMapper.addTopicCtrlConf(entity, result)) {
- sBuffer.append("[addTopicCtrlConf], ")
+ if (topicCtrlMapper.addTopicCtrlConf(entity, strBuff, result)) {
+ strBuff.append("[addTopicCtrlConf], ")
.append(entity.getCreateUser())
.append(" added topic control record :")
.append(entity.toString());
- logger.info(sBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- sBuffer.append("[addTopicCtrlConf], ")
+ strBuff.append("[addTopicCtrlConf], ")
.append("failure to add topic control record : ")
.append(result.getErrMsg());
- logger.warn(sBuffer.toString());
+ logger.warn(strBuff.toString());
}
- sBuffer.delete(0, sBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updTopicCtrlConf(TopicCtrlEntity entity,
- StringBuilder sBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (topicCtrlMapper.updTopicCtrlConf(entity, result)) {
+ if (topicCtrlMapper.updTopicCtrlConf(entity, strBuff, result)) {
TopicCtrlEntity oldEntity =
(TopicCtrlEntity) result.getRetData();
TopicCtrlEntity curEntity =
topicCtrlMapper.getTopicCtrlConf(entity.getTopicName());
- sBuffer.append("[updTopicCtrlConf], ")
+ strBuff.append("[updTopicCtrlConf], ")
.append(entity.getModifyUser())
.append(" updated record from
:").append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(sBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- sBuffer.append("[updTopicCtrlConf], ")
+ strBuff.append("[updTopicCtrlConf], ")
.append("failure to update topic control record : ")
.append(result.getErrMsg());
- logger.warn(sBuffer.toString());
+ logger.warn(strBuff.toString());
}
- sBuffer.delete(0, sBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delTopicCtrlConf(String operator,
- String topicName,
- StringBuilder sBuffer,
- ProcessResult result) {
+ public boolean delTopicCtrlConf(String operator, String topicName,
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
@@ -633,17 +619,17 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
TopicCtrlEntity entity =
(TopicCtrlEntity) result.getRetData();
if (entity != null) {
- sBuffer.append("[delTopicCtrlConf], ").append(operator)
+ strBuff.append("[delTopicCtrlConf], ").append(operator)
.append(" deleted topic control record
:").append(entity.toString());
- logger.info(sBuffer.toString());
+ logger.info(strBuff.toString());
}
} else {
- sBuffer.append("[delTopicCtrlConf], ")
+ strBuff.append("[delTopicCtrlConf], ")
.append("failure to delete topic control record : ")
.append(result.getErrMsg());
- logger.warn(sBuffer.toString());
+ logger.warn(strBuff.toString());
}
- sBuffer.delete(0, sBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -666,62 +652,58 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
// group configure api
@Override
public boolean addGroupResCtrlConf(GroupResCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (groupResCtrlMapper.addGroupResCtrlConf(entity, result)) {
- strBuffer.append("[addGroupResCtrlConf], ")
+ if (groupResCtrlMapper.addGroupResCtrlConf(entity, strBuff, result)) {
+ strBuff.append("[addGroupResCtrlConf], ")
.append(entity.getCreateUser())
.append(" added group resource control record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[addGroupResCtrlConf], ")
+ strBuff.append("[addGroupResCtrlConf], ")
.append("failure to add group resource control record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updGroupResCtrlConf(GroupResCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (groupResCtrlMapper.updGroupResCtrlConf(entity, result)) {
+ if (groupResCtrlMapper.updGroupResCtrlConf(entity, strBuff, result)) {
GroupResCtrlEntity oldEntity =
(GroupResCtrlEntity) result.getRetData();
GroupResCtrlEntity curEntity =
groupResCtrlMapper.getGroupResCtrlConf(entity.getGroupName());
- strBuffer.append("[updGroupResCtrlConf], ")
+ strBuff.append("[updGroupResCtrlConf], ")
.append(entity.getModifyUser())
.append(" updated record from :")
.append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[updGroupResCtrlConf], ")
+ strBuff.append("[updGroupResCtrlConf], ")
.append("failure to update group resource control record :
")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delGroupResCtrlConf(String operator,
- String groupName,
- StringBuilder strBuffer,
- ProcessResult result) {
+ public boolean delGroupResCtrlConf(String operator, String groupName,
+ StringBuilder strBuff, ProcessResult
result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
@@ -730,18 +712,18 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
GroupResCtrlEntity entity =
(GroupResCtrlEntity) result.getRetData();
if (entity != null) {
- strBuffer.append("[delGroupResCtrlConf], ").append(operator)
+ strBuff.append("[delGroupResCtrlConf], ").append(operator)
.append(" deleted group resource control record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
}
} else {
- strBuffer.append("[delGroupResCtrlConf], ")
+ strBuff.append("[delGroupResCtrlConf], ")
.append("failure to delete group resource control record :
")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -758,61 +740,57 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
@Override
public boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff,
ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (groupConsumeCtrlMapper.addGroupConsumeCtrlConf(entity, result)) {
- strBuffer.append("[addGroupConsumeCtrlConf], ")
+ if (groupConsumeCtrlMapper.addGroupConsumeCtrlConf(entity, strBuff,
result)) {
+ strBuff.append("[addGroupConsumeCtrlConf], ")
.append(entity.getCreateUser())
.append(" added group consume control record :")
.append(entity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[addGroupConsumeCtrlConf], ")
+ strBuff.append("[addGroupConsumeCtrlConf], ")
.append("failure to add group consume control record : ")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
public boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result) {
+ StringBuilder strBuff,
ProcessResult result) {
// check current status
if (!checkStoreStatus(true, result)) {
return result.isSuccess();
}
- if (groupConsumeCtrlMapper.updGroupConsumeCtrlConf(entity, result)) {
+ if (groupConsumeCtrlMapper.updGroupConsumeCtrlConf(entity, strBuff,
result)) {
GroupConsumeCtrlEntity oldEntity =
(GroupConsumeCtrlEntity) result.getRetData();
GroupConsumeCtrlEntity curEntity =
groupConsumeCtrlMapper.getGroupConsumeCtrlConfByRecKey(entity.getRecordKey());
- strBuffer.append("[updGroupConsumeCtrlConf], ")
+ strBuff.append("[updGroupConsumeCtrlConf], ")
.append(entity.getModifyUser())
.append(" updated record from
:").append(oldEntity.toString())
.append(" to ").append(curEntity.toString());
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[updGroupConsumeCtrlConf], ")
+ strBuff.append("[updGroupConsumeCtrlConf], ")
.append("failure to update group consume control record :
")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delGroupConsumeCtrlConf(String operator,
- String groupName,
- String topicName,
- StringBuilder sBuffer,
+ public boolean delGroupConsumeCtrlConf(String operator, String groupName,
+ String topicName, StringBuilder
strBuff,
ProcessResult result) {
// check current status
if (groupName == null && topicName == null) {
@@ -823,26 +801,24 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (groupConsumeCtrlMapper.delGroupConsumeCtrlConf(groupName,
topicName, result)) {
- sBuffer.append("[delGroupConsumeCtrlConf], ").append(operator)
+ strBuff.append("[delGroupConsumeCtrlConf], ").append(operator)
.append(" deleted group consume control record by index :
")
.append("groupName=").append(groupName)
.append(", topicName=").append(topicName);
- logger.info(sBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- sBuffer.append("[delGroupConsumeCtrlConf], ")
+ strBuff.append("[delGroupConsumeCtrlConf], ")
.append("failure to delete group consume control record :
")
.append(result.getErrMsg());
- logger.warn(sBuffer.toString());
+ logger.warn(strBuff.toString());
}
- sBuffer.delete(0, sBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@Override
- public boolean delGroupConsumeCtrlConf(String operator,
- String recordKey,
- StringBuilder strBuffer,
- ProcessResult result) {
+ public boolean delGroupConsumeCtrlConf(String operator, String recordKey,
+ StringBuilder strBuff,
ProcessResult result) {
if (recordKey == null) {
result.setSuccResult(null);
return result.isSuccess();
@@ -852,17 +828,17 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return result.isSuccess();
}
if (groupConsumeCtrlMapper.delGroupConsumeCtrlConf(recordKey, result))
{
- strBuffer.append("[delGroupConsumeCtrlConf], ").append(operator)
+ strBuff.append("[delGroupConsumeCtrlConf], ").append(operator)
.append(" deleted group consume control record by index :
")
.append("recordKey=").append(recordKey);
- logger.info(strBuffer.toString());
+ logger.info(strBuff.toString());
} else {
- strBuffer.append("[delGroupConsumeCtrlConf], ")
+ strBuff.append("[delGroupConsumeCtrlConf], ")
.append("failure to delete group consume control record :
")
.append(result.getErrMsg());
- logger.warn(strBuffer.toString());
+ logger.warn(strBuff.toString());
}
- strBuffer.delete(0, strBuffer.length());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
@@ -953,7 +929,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
/**
* Transfer master role to other replica node
*
- * @throws Exception
+ * @throws Exception the exception information
*/
@Override
public void transferMaster() throws Exception {
@@ -1011,7 +987,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
/**
* Get group address info
*
- * @return the group address inforamtion
+ * @return the group address information
*/
@Override
public ClusterGroupVO getGroupAddressStrInfo() {
@@ -1024,11 +1000,11 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
return clusterGroupVO;
}
// translate replication group info to ClusterGroupVO structure
- Tuple2<Boolean, List<ClusterNodeVO>> transReult =
+ Tuple2<Boolean, List<ClusterNodeVO>> transResult =
transReplicateNodes(replicationGroup);
- clusterGroupVO.setNodeData(transReult.getF1());
+ clusterGroupVO.setNodeData(transResult.getF1());
clusterGroupVO.setPrimaryNodeActive(isPrimaryNodeActive());
- if (transReult.getF0()) {
+ if (transResult.getF0()) {
if (isPrimaryNodeActive()) {
clusterGroupVO.setGroupStatus("Running-ReadOnly");
} else {
@@ -1173,44 +1149,40 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
logger.error("[BDB Impl] found executorService is null while
doWork!");
return;
}
- executorService.submit(new Runnable() {
- @Override
- public void run() {
- ProcessResult result = new ProcessResult();
- StringBuilder sBuilder =
- new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE);
- switch (stateChangeEvent.getState()) {
- case MASTER:
- if (!isMaster) {
- try {
- reloadMetaStore();
- isMaster = true;
-
masterSinceTime.set(System.currentTimeMillis());
- masterNodeName =
stateChangeEvent.getMasterNodeName();
- logger.info(sBuilder.append("[BDB Impl] ")
- .append(currentNode)
- .append(" is a
master.").toString());
- } catch (Throwable e) {
- isMaster = false;
- logger.error("[BDB Impl] fatal error when
Reloading Info ", e);
- }
+ executorService.submit(() -> {
+ StringBuilder sBuilder =
+ new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE);
+ switch (stateChangeEvent.getState()) {
+ case MASTER:
+ if (!isMaster) {
+ try {
+ reloadMetaStore();
+ isMaster = true;
+
masterSinceTime.set(System.currentTimeMillis());
+ masterNodeName =
stateChangeEvent.getMasterNodeName();
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode)
+ .append(" is a master.").toString());
+ } catch (Throwable e) {
+ isMaster = false;
+ logger.error("[BDB Impl] fatal error when
Reloading Info ", e);
}
- break;
- case REPLICA:
- isMaster = false;
- masterNodeName =
stateChangeEvent.getMasterNodeName();
- logger.info(sBuilder.append("[BDB Impl] ")
- .append(currentNode).append(" is a
slave.").toString());
- break;
- default:
- isMaster = false;
- logger.info(sBuilder.append("[BDB Impl] ")
- .append(currentNode).append(" is Unknown
state ")
-
.append(stateChangeEvent.getState().name()).toString());
- break;
- }
- sBuilder.delete(0, sBuilder.length());
+ }
+ break;
+ case REPLICA:
+ isMaster = false;
+ masterNodeName = stateChangeEvent.getMasterNodeName();
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode).append(" is a
slave.").toString());
+ break;
+ default:
+ isMaster = false;
+ logger.info(sBuilder.append("[BDB Impl] ")
+ .append(currentNode).append(" is Unknown state
")
+
.append(stateChangeEvent.getState().name()).toString());
+ break;
}
+ sBuilder.delete(0, sBuilder.length());
});
}
}
@@ -1362,7 +1334,7 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
}
private ReplicationGroup getCurrReplicationGroup() {
- ReplicationGroup replicationGroup = null;
+ ReplicationGroup replicationGroup;
try {
replicationGroup = repEnv.getGroup();
} catch (Throwable e) {
@@ -1375,8 +1347,8 @@ public class BdbMetaStoreServiceImpl implements
MetaStoreService {
/**
* Query replication group nodes status and translate to ClusterNodeVO type
*
+ * @param replicationGroup the replication group
* @return if has master, replication nodes info
- * @throws InterruptedException if the operation was interrupted
*/
private Tuple2<Boolean, List<ClusterNodeVO>> transReplicateNodes(
ReplicationGroup replicationGroup) {
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
index 82d5eb2..3daeb15 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/MetaStoreService.java
@@ -41,27 +41,23 @@ public interface MetaStoreService extends KeepAlive, Server
{
* Add or update cluster default setting
*
* @param entity the cluster default setting entity will be add
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
- * @throws Exception
*/
boolean addClusterConfig(ClusterSettingEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Update cluster default setting
*
* @param entity the cluster default setting entity will be add
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
- * @throws Exception
*/
boolean updClusterConfig(ClusterSettingEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
ClusterSettingEntity getClusterConfig();
@@ -69,52 +65,47 @@ public interface MetaStoreService extends KeepAlive, Server
{
* Delete cluster default setting
*
* @param operator operator
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success
*/
boolean delClusterConfig(String operator,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
// broker configure api
/**
* Add broker configure information
*
* @param entity the broker configure entity will be add
- * @param strBuffer the print information string buffer
+ * @param strBuff the print information string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean addBrokerConf(BrokerConfEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Modify broker configure information
*
* @param entity the broker configure entity will be update
- * @param strBuffer the print information string buffer
+ * @param strBuff the print information string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean updBrokerConf(BrokerConfEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Delete broker configure information
*
* @param operator operator
* @param brokerId need deleted broker id
- * @param strBuffer the print information string buffer
+ * @param strBuff the print information string buffer
* @param result the process result return
* @return true if success otherwise false
*/
- boolean delBrokerConf(String operator,
- int brokerId,
- StringBuilder strBuffer,
- ProcessResult result);
+ boolean delBrokerConf(String operator, int brokerId,
+ StringBuilder strBuff, ProcessResult result);
Map<Integer, BrokerConfEntity> getBrokerConfInfo(BrokerConfEntity
qryEntity);
@@ -128,22 +119,16 @@ public interface MetaStoreService extends KeepAlive,
Server {
// topic configure api
boolean addTopicConf(TopicDeployEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
boolean updTopicConf(TopicDeployEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
- boolean delTopicConf(String operator,
- String recordKey,
- StringBuilder strBuffer,
- ProcessResult result);
+ boolean delTopicConf(String operator, String recordKey,
+ StringBuilder strBuff, ProcessResult result);
- boolean delTopicConfByBrokerId(String operator,
- int brokerId,
- StringBuilder strBuffer,
- ProcessResult result);
+ boolean delTopicConfByBrokerId(String operator, int brokerId,
+ StringBuilder strBuff, ProcessResult
result);
boolean hasConfiguredTopics(int brokerId);
@@ -177,39 +162,35 @@ public interface MetaStoreService extends KeepAlive,
Server {
* Add topic control configure info
*
* @param entity the topic control info entity will be add
- * @param sBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean addTopicCtrlConf(TopicCtrlEntity entity,
- StringBuilder sBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Update topic control configure
*
* @param entity the topic control info entity will be update
- * @param sBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean updTopicCtrlConf(TopicCtrlEntity entity,
- StringBuilder sBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Delete topic control configure
*
* @param operator operator
* @param topicName the topicName will be deleted
- * @param sBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
- boolean delTopicCtrlConf(String operator,
- String topicName,
- StringBuilder sBuffer,
- ProcessResult result);
+ boolean delTopicCtrlConf(String operator, String topicName,
+ StringBuilder strBuff, ProcessResult result);
TopicCtrlEntity getTopicCtrlConf(String topicName);
@@ -223,39 +204,35 @@ public interface MetaStoreService extends KeepAlive,
Server {
* Add group resource control configure info
*
* @param entity the group resource control info entity will be add
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean addGroupResCtrlConf(GroupResCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Update group reource control configure
*
* @param entity the group resource control info entity will be update
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean updGroupResCtrlConf(GroupResCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult result);
/**
* Delete group resource control configure
*
* @param operator operator
* @param groupName the group will be deleted
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
- boolean delGroupResCtrlConf(String operator,
- String groupName,
- StringBuilder strBuffer,
- ProcessResult result);
+ boolean delGroupResCtrlConf(String operator, String groupName,
+ StringBuilder strBuff, ProcessResult result);
GroupResCtrlEntity getGroupResCtrlConf(String groupName);
@@ -267,25 +244,23 @@ public interface MetaStoreService extends KeepAlive,
Server {
* Add group consume control configure
*
* @param entity the group consume control info entity will be add
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult
result);
/**
* Modify group consume control configure
*
* @param entity the group consume control info entity will be update
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success otherwise false
*/
boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
- StringBuilder strBuffer,
- ProcessResult result);
+ StringBuilder strBuff, ProcessResult
result);
/**
* Delete group consume control configure
@@ -295,27 +270,25 @@ public interface MetaStoreService extends KeepAlive,
Server {
* @param topicName the blacklist record related to topic
* allow groupName or topicName is null,
* but not all null
+ * @param strBuff the string buffer
+ * @param result the process result
* @return true if success
*/
boolean delGroupConsumeCtrlConf(String operator,
- String groupName,
- String topicName,
- StringBuilder sBuffer,
- ProcessResult result);
+ String groupName, String topicName,
+ StringBuilder strBuff, ProcessResult
result);
/**
* Delete group consume control configure
*
* @param operator operator
* @param recordKey the record key to group consume control record
- * @param strBuffer the print info string buffer
+ * @param strBuff the print info string buffer
* @param result the process result return
* @return true if success
*/
- boolean delGroupConsumeCtrlConf(String operator,
- String recordKey,
- StringBuilder strBuffer,
- ProcessResult result);
+ boolean delGroupConsumeCtrlConf(String operator, String recordKey,
+ StringBuilder strBuff, ProcessResult
result);
void registerObserver(AliveObserver eventObserver);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BaseEntity.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BaseEntity.java
index e287dad..6ee8848 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BaseEntity.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BaseEntity.java
@@ -160,9 +160,9 @@ public class BaseEntity implements Serializable, Cloneable {
return attributes;
}
- public void setCreateInfo(String creater, Date createDate) {
- if (TStringUtils.isNotBlank(creater)) {
- this.createUser = creater;
+ public void setCreateInfo(String creator, Date createDate) {
+ if (TStringUtils.isNotBlank(creator)) {
+ this.createUser = creator;
}
setCreateDate(createDate);
}
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BrokerConfEntity.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BrokerConfEntity.java
index 9a888b4..056ab2b 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BrokerConfEntity.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/BrokerConfEntity.java
@@ -61,7 +61,7 @@ public class BrokerConfEntity extends BaseEntity implements
Cloneable {
}
/**
- * Initial Broker Conigure entity by BdbBrokerConfEntity
+ * Initial Broker Configure entity by BdbBrokerConfEntity
*
* @param bdbEntity need initialed BdbBrokerConfEntity information
*/
@@ -383,11 +383,10 @@ public class BrokerConfEntity extends BaseEntity
implements Cloneable {
/**
* Get broker config string
*
- * @return config string
+ * @param strBuff the string buffer
*/
- public String getBrokerDefaultConfInfo() {
- return new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
-
.append(topicProps.getNumPartitions()).append(TokenConstants.ATTR_SEP)
+ public void getBrokerDefaultConfInfo(StringBuilder strBuff) {
+
strBuff.append(topicProps.getNumPartitions()).append(TokenConstants.ATTR_SEP)
.append(topicProps.isAcceptPublish()).append(TokenConstants.ATTR_SEP)
.append(topicProps.isAcceptSubscribe()).append(TokenConstants.ATTR_SEP)
.append(topicProps.getUnflushThreshold()).append(TokenConstants.ATTR_SEP)
@@ -398,7 +397,7 @@ public class BrokerConfEntity extends BaseEntity implements
Cloneable {
.append(topicProps.getUnflushDataHold()).append(TokenConstants.ATTR_SEP)
.append(topicProps.getMemCacheMsgSizeInMB()).append(TokenConstants.ATTR_SEP)
.append(topicProps.getMemCacheMsgCntInK()).append(TokenConstants.ATTR_SEP)
- .append(topicProps.getMemCacheFlushIntvl()).toString();
+ .append(topicProps.getMemCacheFlushIntvl());
}
/**
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/ClusterSettingEntity.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/ClusterSettingEntity.java
index ef624e1..332a963 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/ClusterSettingEntity.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/ClusterSettingEntity.java
@@ -33,7 +33,7 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.TStoreConstan
*/
public class ClusterSettingEntity extends BaseEntity implements Cloneable {
- private String recordKey =
+ private final String recordKey =
TStoreConstants.TOKEN_DEFAULT_CLUSTER_SETTING;
// broker tcp port
private int brokerPort = TBaseConstants.META_VALUE_UNDEFINED;
@@ -271,7 +271,6 @@ public class ClusterSettingEntity extends BaseEntity
implements Cloneable {
&& maxMsgSizeInB == other.maxMsgSizeInB
&& qryPriorityId == other.qryPriorityId
&& gloFlowCtrlRuleCnt == other.gloFlowCtrlRuleCnt
- && recordKey.equals(other.recordKey)
&& Objects.equals(clsDefTopicProps, other.clsDefTopicProps)
&& gloFlowCtrlStatus == other.gloFlowCtrlStatus
&& Objects.equals(gloFlowCtrlRuleInfo,
other.gloFlowCtrlRuleInfo);
@@ -283,7 +282,7 @@ public class ClusterSettingEntity extends BaseEntity
implements Cloneable {
* @param sBuilder build container
* @param isLongName if return field key is long name
* @param fullFormat if return full format json
- * @return
+ * @return the serialized result
*/
public StringBuilder toWebJsonStr(StringBuilder sBuilder,
boolean isLongName,
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/GroupConsumeCtrlEntity.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/GroupConsumeCtrlEntity.java
index 3b7faa8..983b276 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/GroupConsumeCtrlEntity.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/GroupConsumeCtrlEntity.java
@@ -170,7 +170,7 @@ public class GroupConsumeCtrlEntity extends BaseEntity
implements Cloneable {
* @param consumeEnable new consume enable status
* @param disableRsn new disable reason
* @param filterEnable new filter enable status
- * @param filterCondStr new fliter condition configure
+ * @param filterCondStr new filter condition configure
*
* @return whether data is changed
*/
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/TopicDeployEntity.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/TopicDeployEntity.java
index dcd328d..6df2cfe 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/TopicDeployEntity.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/entity/TopicDeployEntity.java
@@ -286,7 +286,7 @@ public class TopicDeployEntity extends BaseEntity
implements Cloneable {
* @param sBuilder build container
* @param isLongName if return field key is long name
* @param fullFormat if return full format json
- * @return
+ * @return the serialized content
*/
public StringBuilder toWebJsonStr(StringBuilder sBuilder,
boolean isLongName,
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 8e7c513..b39f88b 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
@@ -24,9 +24,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.Br
public interface BrokerConfigMapper extends AbstractMapper {
- boolean addBrokerConf(BrokerConfEntity memEntity, ProcessResult result);
+ boolean addBrokerConf(BrokerConfEntity memEntity,
+ StringBuilder strBuff, ProcessResult result);
- boolean updBrokerConf(BrokerConfEntity memEntity, ProcessResult result);
+ boolean updBrokerConf(BrokerConfEntity memEntity,
+ StringBuilder strBuff, ProcessResult result);
boolean delBrokerConf(int brokerId, ProcessResult result);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/ClusterConfigMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/ClusterConfigMapper.java
index b647944..0089ebb 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/ClusterConfigMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/ClusterConfigMapper.java
@@ -22,9 +22,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.Cl
public interface ClusterConfigMapper extends AbstractMapper {
- boolean addClusterConfig(ClusterSettingEntity memEntity, ProcessResult
result);
+ boolean addClusterConfig(ClusterSettingEntity memEntity,
+ StringBuilder strBuff, ProcessResult result);
- boolean updClusterConfig(ClusterSettingEntity memEntity, ProcessResult
result);
+ boolean updClusterConfig(ClusterSettingEntity memEntity,
+ StringBuilder strBuff, ProcessResult result);
ClusterSettingEntity getClusterConfig();
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupConsumeCtrlMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupConsumeCtrlMapper.java
index f015421..cf3d89e 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupConsumeCtrlMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupConsumeCtrlMapper.java
@@ -25,9 +25,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.Gr
public interface GroupConsumeCtrlMapper extends AbstractMapper {
- boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
ProcessResult result);
+ boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult
result);
- boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
ProcessResult result);
+ boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult
result);
boolean delGroupConsumeCtrlConf(String recordKey, ProcessResult result);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupResCtrlMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupResCtrlMapper.java
index e72b3d2..e6ba466 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupResCtrlMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/GroupResCtrlMapper.java
@@ -24,9 +24,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.Gr
public interface GroupResCtrlMapper extends AbstractMapper {
- boolean addGroupResCtrlConf(GroupResCtrlEntity entity, ProcessResult
result);
+ boolean addGroupResCtrlConf(GroupResCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult result);
- boolean updGroupResCtrlConf(GroupResCtrlEntity entity, ProcessResult
result);
+ boolean updGroupResCtrlConf(GroupResCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult result);
boolean delGroupResCtrlConf(String groupName, ProcessResult result);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicCtrlMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicCtrlMapper.java
index 82a5c6f..b5dd4a6 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicCtrlMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicCtrlMapper.java
@@ -25,9 +25,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.To
public interface TopicCtrlMapper extends AbstractMapper {
- boolean addTopicCtrlConf(TopicCtrlEntity entity, ProcessResult result);
+ boolean addTopicCtrlConf(TopicCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult result);
- boolean updTopicCtrlConf(TopicCtrlEntity entity, ProcessResult result);
+ boolean updTopicCtrlConf(TopicCtrlEntity entity,
+ StringBuilder strBuff, ProcessResult result);
boolean delTopicCtrlConf(String topicName, ProcessResult result);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicDeployMapper.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicDeployMapper.java
index cc1d510..ab2d387 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicDeployMapper.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/dao/mapper/TopicDeployMapper.java
@@ -25,9 +25,11 @@ import
org.apache.inlong.tubemq.server.master.metamanage.metastore.dao.entity.To
public interface TopicDeployMapper extends AbstractMapper {
- boolean addTopicConf(TopicDeployEntity entity, ProcessResult result);
+ boolean addTopicConf(TopicDeployEntity entity,
+ StringBuilder strBuff, ProcessResult result);
- boolean updTopicConf(TopicDeployEntity entity, ProcessResult result);
+ boolean updTopicConf(TopicDeployEntity entity,
+ StringBuilder strBuff, ProcessResult result);
boolean delTopicConf(String recordKey, ProcessResult result);
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 9eb5d49..f25722e 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
@@ -27,7 +27,6 @@ import java.util.HashSet;
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.server.common.exception.LoadMetaException;
@@ -45,13 +44,13 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
// broker config store
private EntityStore brokerConfStore;
- private PrimaryIndex<Integer/* brokerId */, BdbBrokerConfEntity>
brokerConfIndex;
- private ConcurrentHashMap<Integer/* brokerId */, BrokerConfEntity>
brokerConfCache =
- new ConcurrentHashMap<>();
- private ConcurrentHashMap<String/* brokerIP */, Integer/* brokerId */>
brokerIpIndexCache =
- new ConcurrentHashMap<>();
- private ConcurrentHashMap<Integer/* regionId */,
ConcurrentHashSet<Integer>> regionIndexCache =
- new ConcurrentHashMap<>();
+ private final PrimaryIndex<Integer/* brokerId */, BdbBrokerConfEntity>
brokerConfIndex;
+ private final ConcurrentHashMap<Integer/* brokerId */, BrokerConfEntity>
+ brokerConfCache = new ConcurrentHashMap<>();
+ private final ConcurrentHashMap<String/* brokerIP */, Integer/* brokerId
*/>
+ brokerIpIndexCache = new ConcurrentHashMap<>();
+ private final ConcurrentHashMap<Integer/* regionId */,
ConcurrentHashSet<Integer>>
+ regionIndexCache = new ConcurrentHashMap<>();
public BdbBrokerConfigMapperImpl(ReplicatedEnvironment repEnv, StoreConfig
storeConfig) {
brokerConfStore = new EntityStore(repEnv,
@@ -104,53 +103,55 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
}
@Override
- public boolean addBrokerConf(BrokerConfEntity memEntity, ProcessResult
result) {
+ public boolean addBrokerConf(BrokerConfEntity memEntity,
+ StringBuilder strBuff, ProcessResult result) {
BrokerConfEntity curEntity =
brokerConfCache.get(memEntity.getBrokerId());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The broker's brokerId
").append(memEntity.getBrokerId())
+ strBuff.append("The broker's brokerId
").append(memEntity.getBrokerId())
.append(" has already exists, the value must be
unique!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
Integer curBrokerId = brokerIpIndexCache.get(memEntity.getBrokerIp());
if (curBrokerId != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The broker's brokerIp
").append(memEntity.getBrokerIp())
+ strBuff.append("The broker's brokerIp
").append(memEntity.getBrokerIp())
.append(" has already exists, the value must be
unique!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putBrokerConfig2Bdb(memEntity, result)) {
+ if (putBrokerConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
}
return result.isSuccess();
}
@Override
- public boolean updBrokerConf(BrokerConfEntity memEntity, ProcessResult
result) {
+ public boolean updBrokerConf(BrokerConfEntity memEntity,
+ StringBuilder strBuff, ProcessResult result) {
BrokerConfEntity curEntity =
brokerConfCache.get(memEntity.getBrokerId());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The broker
").append(memEntity.getBrokerIp())
+ strBuff.append("The broker
").append(memEntity.getBrokerIp())
.append("'s configure is not exists, please add
record first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
if (curEntity.equals(memEntity)) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The broker
").append(memEntity.getBrokerIp())
+ strBuff.append("The broker
").append(memEntity.getBrokerIp())
.append("'s configure have not changed, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putBrokerConfig2Bdb(memEntity, result)) {
+ if (putBrokerConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
result.setSuccResult(curEntity);
}
@@ -159,7 +160,10 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
/**
* delete broker configure info from bdb store
- * @return
+ *
+ * @param brokerId the broker id to be deleted
+ * @param result the process result
+ * @return whether success
*/
@Override
public boolean delBrokerConf(int brokerId, ProcessResult result) {
@@ -177,6 +181,8 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
/**
* get broker configure info from bdb store
+ *
+ * @param qryEntity the query conditions
* @return result, only read
*/
@Override
@@ -280,6 +286,8 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
/**
* get broker configure info from bdb store
+ *
+ * @param brokerId the broker id to be queried
* @return result, only read
*/
@Override
@@ -289,6 +297,8 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
/**
* get broker configure info from bdb store
+ *
+ * @param brokerIp the broker ip to be queried
* @return result, only read
*/
@Override
@@ -324,21 +334,22 @@ public class BdbBrokerConfigMapperImpl implements
BrokerConfigMapper {
* Put cluster setting info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
- private boolean putBrokerConfig2Bdb(BrokerConfEntity memEntity,
ProcessResult result) {
- BdbBrokerConfEntity retData = null;
+ private boolean putBrokerConfig2Bdb(BrokerConfEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
BdbBrokerConfEntity bdbEntity =
memEntity.buildBdbBrokerConfEntity();
try {
- retData = brokerConfIndex.put(bdbEntity);
+ brokerConfIndex.put(bdbEntity);
} catch (Throwable e) {
logger.error("[BDB Impl] put broker configure failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put broker configure failure: ")
+ strBuff.append("Put broker configure failure: ")
.append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
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 ff4266b..3df718f 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
@@ -24,7 +24,6 @@ import com.sleepycat.persist.PrimaryIndex;
import com.sleepycat.persist.StoreConfig;
import java.util.Map;
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.server.common.exception.LoadMetaException;
import
org.apache.inlong.tubemq.server.master.bdbstore.bdbentitys.BdbClusterSettingEntity;
@@ -41,7 +40,7 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
LoggerFactory.getLogger(BdbClusterConfigMapperImpl.class);
private EntityStore clsDefSettingStore;
- private PrimaryIndex<String, BdbClusterSettingEntity> clsDefSettingIndex;
+ private final PrimaryIndex<String, BdbClusterSettingEntity>
clsDefSettingIndex;
Map<String, ClusterSettingEntity> metaDataCache = new
ConcurrentHashMap<>();
public BdbClusterConfigMapperImpl(ReplicatedEnvironment repEnv,
StoreConfig storeConfig) {
@@ -98,17 +97,19 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
* Put cluster setting info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
@Override
- public boolean addClusterConfig(ClusterSettingEntity memEntity,
ProcessResult result) {
+ public boolean addClusterConfig(ClusterSettingEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
if (!metaDataCache.isEmpty()) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
"The cluster setting already exists, please delete or
update!");
return result.isSuccess();
}
- if (putClusterConfig2Bdb(memEntity, result)) {
+ if (putClusterConfig2Bdb(memEntity, strBuff, result)) {
metaDataCache.put(memEntity.getRecordKey(), memEntity);
}
return result.isSuccess();
@@ -118,11 +119,13 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
* Update cluster setting info in bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
@Override
- public boolean updClusterConfig(ClusterSettingEntity memEntity,
ProcessResult result) {
+ public boolean updClusterConfig(ClusterSettingEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
if (metaDataCache.isEmpty()) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
"The cluster setting is null, please add record first!");
@@ -134,7 +137,7 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
"The cluster settings have not changed!");
return result.isSuccess();
}
- if (putClusterConfig2Bdb(memEntity, result)) {
+ if (putClusterConfig2Bdb(memEntity, strBuff, result)) {
metaDataCache.put(memEntity.getRecordKey(), memEntity);
result.setSuccResult(curEntity);
}
@@ -143,6 +146,7 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
/**
* get current cluster setting from bdb store
+ *
* @return current cluster setting, null or object, only read
*/
@Override
@@ -152,7 +156,9 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
/**
* delete current cluster setting from bdb store
- * @return if success
+ *
+ * @param result the process result
+ * @return the process result
*/
@Override
public boolean delClusterConfig(ProcessResult result) {
@@ -168,7 +174,8 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
return true;
}
- private boolean putClusterConfig2Bdb(ClusterSettingEntity memEntity,
ProcessResult result) {
+ private boolean putClusterConfig2Bdb(ClusterSettingEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
BdbClusterSettingEntity bdbEntity =
memEntity.buildBdbClsDefSettingEntity();
try {
@@ -176,9 +183,9 @@ public class BdbClusterConfigMapperImpl implements
ClusterConfigMapper {
} catch (Throwable e) {
logger.error("[BDB Impl] put cluster configure failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put cluster configure failure: ")
+ strBuff.append("Put cluster configure failure: ")
.append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupConsumeCtrlMapperImpl.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupConsumeCtrlMapperImpl.java
index cef8f3f..b69f321 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupConsumeCtrlMapperImpl.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/metamanage/metastore/impl/bdbimpl/BdbGroupConsumeCtrlMapperImpl.java
@@ -30,7 +30,6 @@ 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;
@@ -48,13 +47,13 @@ public class BdbGroupConsumeCtrlMapperImpl implements
GroupConsumeCtrlMapper {
LoggerFactory.getLogger(BdbGroupConsumeCtrlMapperImpl.class);
// consumer group consume control store
private EntityStore groupConsumeStore;
- private PrimaryIndex<String/* recordKey */, BdbGroupFilterCondEntity>
groupConsumeIndex;
+ private final PrimaryIndex<String/* recordKey */,
BdbGroupFilterCondEntity> groupConsumeIndex;
// configure cache
- private ConcurrentHashMap<String/* recordKey */, GroupConsumeCtrlEntity>
+ private final ConcurrentHashMap<String/* recordKey */,
GroupConsumeCtrlEntity>
grpConsumeCtrlCache = new ConcurrentHashMap<>();
- private ConcurrentHashMap<String/* topicName */, ConcurrentHashSet<String>>
+ private final ConcurrentHashMap<String/* topicName */,
ConcurrentHashSet<String>>
grpConsumeCtrlTopicCache = new ConcurrentHashMap<>();
- private ConcurrentHashMap<String/* groupName */, ConcurrentHashSet<String>>
+ private final ConcurrentHashMap<String/* groupName */,
ConcurrentHashSet<String>>
grpConsumeCtrlGroupCache = new ConcurrentHashMap<>();
public BdbGroupConsumeCtrlMapperImpl(ReplicatedEnvironment repEnv,
StoreConfig storeConfig) {
@@ -108,44 +107,46 @@ public class BdbGroupConsumeCtrlMapperImpl implements
GroupConsumeCtrlMapper {
}
@Override
- public boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity memEntity,
ProcessResult result) {
+ public boolean addGroupConsumeCtrlConf(GroupConsumeCtrlEntity memEntity,
+ StringBuilder strBuff,
ProcessResult result) {
GroupConsumeCtrlEntity curEntity =
grpConsumeCtrlCache.get(memEntity.getRecordKey());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group consume
").append(memEntity.getRecordKey())
+ strBuff.append("The group consume
").append(memEntity.getRecordKey())
.append("'s configure already exists, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putGroupConsumeCtrlConfig2Bdb(memEntity, result)) {
+ if (putGroupConsumeCtrlConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
}
return result.isSuccess();
}
@Override
- public boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity memEntity,
ProcessResult result) {
+ public boolean updGroupConsumeCtrlConf(GroupConsumeCtrlEntity memEntity,
+ StringBuilder strBuff,
ProcessResult result) {
GroupConsumeCtrlEntity curEntity =
grpConsumeCtrlCache.get(memEntity.getRecordKey());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group consume
").append(memEntity.getRecordKey())
+ strBuff.append("The group consume
").append(memEntity.getRecordKey())
.append("'s configure is not exists, please add
record first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
if (curEntity.equals(memEntity)) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group consume
").append(memEntity.getRecordKey())
+ strBuff.append("The group consume
").append(memEntity.getRecordKey())
.append("'s configure have not changed, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putGroupConsumeCtrlConfig2Bdb(memEntity, result)) {
+ if (putGroupConsumeCtrlConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
result.setSuccResult(curEntity);
}
@@ -299,18 +300,19 @@ public class BdbGroupConsumeCtrlMapperImpl implements
GroupConsumeCtrlMapper {
}
@Override
- public List<GroupConsumeCtrlEntity>
getGroupConsumeCtrlConf(GroupConsumeCtrlEntity qryEntity) {
- List<GroupConsumeCtrlEntity> retEntitys = new ArrayList<>();
+ public List<GroupConsumeCtrlEntity> getGroupConsumeCtrlConf(
+ GroupConsumeCtrlEntity qryEntity) {
+ List<GroupConsumeCtrlEntity> retEntities = new ArrayList<>();
if (qryEntity == null) {
- retEntitys.addAll(grpConsumeCtrlCache.values());
+ retEntities.addAll(grpConsumeCtrlCache.values());
} else {
for (GroupConsumeCtrlEntity entity : grpConsumeCtrlCache.values())
{
if (entity != null && entity.isMatched(qryEntity)) {
- retEntitys.add(entity);
+ retEntities.add(entity);
}
}
}
- return retEntitys;
+ return retEntities;
}
@Override
@@ -369,22 +371,22 @@ public class BdbGroupConsumeCtrlMapperImpl implements
GroupConsumeCtrlMapper {
* Put Group consume configure info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return true sucess, false failue
+ * @return true success, false failure
*/
- private boolean putGroupConsumeCtrlConfig2Bdb(
- GroupConsumeCtrlEntity memEntity, ProcessResult result) {
- BdbGroupFilterCondEntity retData = null;
+ private boolean putGroupConsumeCtrlConfig2Bdb(GroupConsumeCtrlEntity
memEntity,
+ StringBuilder strBuff,
ProcessResult result) {
BdbGroupFilterCondEntity bdbEntity =
memEntity.buildBdbGroupFilterCondEntity();
try {
- retData = groupConsumeIndex.put(bdbEntity);
+ groupConsumeIndex.put(bdbEntity);
} catch (Throwable e) {
logger.error("[BDB Impl] put consume configure failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put filter configure failure: ")
+ strBuff.append("Put filter configure failure: ")
.append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
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 913b046..34b17c6 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
@@ -26,7 +26,6 @@ import java.util.HashMap;
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.server.common.exception.LoadMetaException;
import
org.apache.inlong.tubemq.server.master.bdbstore.bdbentitys.BdbGroupFlowCtrlEntity;
@@ -42,9 +41,9 @@ public class BdbGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
LoggerFactory.getLogger(BdbGroupResCtrlMapperImpl.class);
// consumer group configure store
private EntityStore groupConfStore;
- private PrimaryIndex<String/* groupName */, BdbGroupFlowCtrlEntity>
groupBaseCtrlIndex;
- private ConcurrentHashMap<String/* groupName */, GroupResCtrlEntity>
groupBaseCtrlCache =
- new ConcurrentHashMap<>();
+ private final PrimaryIndex<String/* groupName */, BdbGroupFlowCtrlEntity>
groupBaseCtrlIndex;
+ private final ConcurrentHashMap<String/* groupName */, GroupResCtrlEntity>
+ groupBaseCtrlCache = new ConcurrentHashMap<>();
public BdbGroupResCtrlMapperImpl(ReplicatedEnvironment repEnv, StoreConfig
storeConfig) {
groupConfStore = new EntityStore(repEnv,
@@ -97,44 +96,46 @@ public class BdbGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
}
@Override
- public boolean addGroupResCtrlConf(GroupResCtrlEntity memEntity,
ProcessResult result) {
+ public boolean addGroupResCtrlConf(GroupResCtrlEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
GroupResCtrlEntity curEntity =
groupBaseCtrlCache.get(memEntity.getGroupName());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group
").append(memEntity.getGroupName())
+ strBuff.append("The group
").append(memEntity.getGroupName())
.append("'s resource control already exists,
please delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putGroupConfigConfig2Bdb(memEntity, result)) {
+ if (putGroupConfigConfig2Bdb(memEntity, strBuff, result)) {
groupBaseCtrlCache.put(memEntity.getGroupName(), memEntity);
}
return result.isSuccess();
}
@Override
- public boolean updGroupResCtrlConf(GroupResCtrlEntity memEntity,
ProcessResult result) {
+ public boolean updGroupResCtrlConf(GroupResCtrlEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
GroupResCtrlEntity curEntity =
groupBaseCtrlCache.get(memEntity.getGroupName());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group
").append(memEntity.getGroupName())
+ strBuff.append("The group
").append(memEntity.getGroupName())
.append("'s resource control is not exists, please
add record first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
if (curEntity.equals(memEntity)) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The group
").append(memEntity.getGroupName())
+ strBuff.append("The group
").append(memEntity.getGroupName())
.append("'s resource control have not changed,
please delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putGroupConfigConfig2Bdb(memEntity, result)) {
+ if (putGroupConfigConfig2Bdb(memEntity, strBuff, result)) {
groupBaseCtrlCache.put(memEntity.getGroupName(), memEntity);
result.setSuccResult(curEntity);
}
@@ -188,21 +189,22 @@ public class BdbGroupResCtrlMapperImpl implements
GroupResCtrlMapper {
* Put Group configure info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
- private boolean putGroupConfigConfig2Bdb(GroupResCtrlEntity memEntity,
ProcessResult result) {
- BdbGroupFlowCtrlEntity retData = null;
+ private boolean putGroupConfigConfig2Bdb(GroupResCtrlEntity memEntity,
+ StringBuilder strBuff,
ProcessResult result) {
BdbGroupFlowCtrlEntity bdbEntity =
memEntity.buildBdbGroupFlowCtrlEntity();
try {
- retData = groupBaseCtrlIndex.put(bdbEntity);
+ groupBaseCtrlIndex.put(bdbEntity);
} catch (Throwable e) {
logger.error("[BDB Impl] put group resource control failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put group resource control failure: ")
+ strBuff.append("Put group resource control failure: ")
.append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
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 a0601b1..17c6b14 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
@@ -29,7 +29,6 @@ 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.server.common.exception.LoadMetaException;
import
org.apache.inlong.tubemq.server.master.bdbstore.bdbentitys.BdbTopicAuthControlEntity;
@@ -46,10 +45,10 @@ public class BdbTopicCtrlMapperImpl implements
TopicCtrlMapper {
// Topic control store
private EntityStore topicCtrlStore;
- private PrimaryIndex<String/* topicName */, BdbTopicAuthControlEntity>
topicCtrlIndex;
+ private final PrimaryIndex<String/* topicName */,
BdbTopicAuthControlEntity> topicCtrlIndex;
// data cache
- private ConcurrentHashMap<String/* topicName */, TopicCtrlEntity>
topicCtrlCache =
- new ConcurrentHashMap<>();
+ private final ConcurrentHashMap<String/* topicName */, TopicCtrlEntity>
+ topicCtrlCache = new ConcurrentHashMap<>();
public BdbTopicCtrlMapperImpl(ReplicatedEnvironment repEnv, StoreConfig
storeConfig) {
topicCtrlStore = new EntityStore(repEnv,
@@ -101,44 +100,46 @@ public class BdbTopicCtrlMapperImpl implements
TopicCtrlMapper {
}
@Override
- public boolean addTopicCtrlConf(TopicCtrlEntity memEntity, ProcessResult
result) {
+ public boolean addTopicCtrlConf(TopicCtrlEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
TopicCtrlEntity curEntity =
topicCtrlCache.get(memEntity.getTopicName());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic control
").append(memEntity.getTopicName())
+ strBuff.append("The topic control
").append(memEntity.getTopicName())
.append("'s configure already exists, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putTopicCtrlConfig2Bdb(memEntity, result)) {
+ if (putTopicCtrlConfig2Bdb(memEntity, strBuff, result)) {
topicCtrlCache.put(memEntity.getTopicName(), memEntity);
}
return result.isSuccess();
}
@Override
- public boolean updTopicCtrlConf(TopicCtrlEntity memEntity, ProcessResult
result) {
+ public boolean updTopicCtrlConf(TopicCtrlEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
TopicCtrlEntity curEntity =
topicCtrlCache.get(memEntity.getTopicName());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic control
").append(memEntity.getTopicName())
+ strBuff.append("The topic control
").append(memEntity.getTopicName())
.append("'s configure is not exists, please add
record first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
if (curEntity.equals(memEntity)) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic control
").append(memEntity.getTopicName())
+ strBuff.append("The topic control
").append(memEntity.getTopicName())
.append("'s configure have not changed, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putTopicCtrlConfig2Bdb(memEntity, result)) {
+ if (putTopicCtrlConfig2Bdb(memEntity, strBuff, result)) {
topicCtrlCache.put(memEntity.getTopicName(), memEntity);
result.setSuccResult(curEntity);
}
@@ -166,17 +167,17 @@ public class BdbTopicCtrlMapperImpl implements
TopicCtrlMapper {
@Override
public List<TopicCtrlEntity> getTopicCtrlConf(TopicCtrlEntity qryEntity) {
- List<TopicCtrlEntity> retEntitys = new ArrayList<>();
+ List<TopicCtrlEntity> retEntities = new ArrayList<>();
if (qryEntity == null) {
- retEntitys.addAll(topicCtrlCache.values());
+ retEntities.addAll(topicCtrlCache.values());
} else {
for (TopicCtrlEntity entity : topicCtrlCache.values()) {
if (entity != null && entity.isMatched(qryEntity)) {
- retEntitys.add(entity);
+ retEntities.add(entity);
}
}
}
- return retEntitys;
+ return retEntities;
}
@Override
@@ -203,20 +204,22 @@ public class BdbTopicCtrlMapperImpl implements
TopicCtrlMapper {
* Put topic control configure info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
- private boolean putTopicCtrlConfig2Bdb(TopicCtrlEntity memEntity,
ProcessResult result) {
- BdbTopicAuthControlEntity retData = null;
+ private boolean putTopicCtrlConfig2Bdb(TopicCtrlEntity memEntity,
+ StringBuilder strBuff,
ProcessResult result) {
BdbTopicAuthControlEntity bdbEntity =
memEntity.buildBdbTopicAuthControlEntity();
try {
- retData = topicCtrlIndex.put(bdbEntity);
+ topicCtrlIndex.put(bdbEntity);
} catch (Throwable e) {
logger.error("[BDB Impl] put topic control failure ", e);
- result.setFailResult(new
StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put topic control failure: ")
- .append(e.getMessage()).toString());
+ result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
+ strBuff.append("Put topic control failure: ")
+ .append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
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 94da730..6160707 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
@@ -30,7 +30,6 @@ 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;
@@ -49,15 +48,15 @@ public class BdbTopicDeployMapperImpl implements
TopicDeployMapper {
// Topic configure store
private EntityStore topicConfStore;
- private PrimaryIndex<String/* recordKey */, BdbTopicConfEntity>
topicConfIndex;
+ private final PrimaryIndex<String/* recordKey */, BdbTopicConfEntity>
topicConfIndex;
// data cache
- private ConcurrentHashMap<String/* recordKey */, TopicDeployEntity>
topicConfCache =
- new ConcurrentHashMap<>();
- private ConcurrentHashMap<Integer/* brokerId */, ConcurrentHashSet<String>>
+ private final ConcurrentHashMap<String/* recordKey */, TopicDeployEntity>
+ topicConfCache = new ConcurrentHashMap<>();
+ private final ConcurrentHashMap<Integer/* brokerId */,
ConcurrentHashSet<String>>
brokerIdCacheIndex = new ConcurrentHashMap<>();
- private ConcurrentHashMap<String/* topicName */, ConcurrentHashSet<String>>
+ private final ConcurrentHashMap<String/* topicName */,
ConcurrentHashSet<String>>
topicNameCacheIndex = new ConcurrentHashMap<>();
- private ConcurrentHashMap<Integer/* brokerId */, ConcurrentHashSet<String>>
+ private final ConcurrentHashMap<Integer/* brokerId */,
ConcurrentHashSet<String>>
brokerId2TopicCacheIndex = new ConcurrentHashMap<>();
public BdbTopicDeployMapperImpl(ReplicatedEnvironment repEnv, StoreConfig
storeConfig) {
@@ -108,44 +107,46 @@ public class BdbTopicDeployMapperImpl implements
TopicDeployMapper {
}
@Override
- public boolean addTopicConf(TopicDeployEntity memEntity, ProcessResult
result) {
+ public boolean addTopicConf(TopicDeployEntity memEntity,
+ StringBuilder strBuff, ProcessResult result) {
TopicDeployEntity curEntity =
topicConfCache.get(memEntity.getRecordKey());
if (curEntity != null) {
result.setFailResult(DataOpErrCode.DERR_EXISTED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic configure
").append(memEntity.getRecordKey())
+ strBuff.append("The topic configure
").append(memEntity.getRecordKey())
.append("'s configure already exists, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putTopicConfig2Bdb(memEntity, result)) {
+ if (putTopicConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
}
return result.isSuccess();
}
@Override
- public boolean updTopicConf(TopicDeployEntity memEntity, ProcessResult
result) {
+ public boolean updTopicConf(TopicDeployEntity memEntity,
+ StringBuilder strBuff, ProcessResult result) {
TopicDeployEntity curEntity =
topicConfCache.get(memEntity.getRecordKey());
if (curEntity == null) {
result.setFailResult(DataOpErrCode.DERR_NOT_EXIST.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic configure
").append(memEntity.getRecordKey())
+ strBuff.append("The topic configure
").append(memEntity.getRecordKey())
.append("'s configure is not exists, please add
record first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
if (curEntity.equals(memEntity)) {
result.setFailResult(DataOpErrCode.DERR_UNCHANGED.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("The topic configure
").append(memEntity.getRecordKey())
+ strBuff.append("The topic configure
").append(memEntity.getRecordKey())
.append("'s configure have not changed, please
delete it first!")
.toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
- if (putTopicConfig2Bdb(memEntity, result)) {
+ if (putTopicConfig2Bdb(memEntity, strBuff, result)) {
addOrUpdCacheRecord(memEntity);
result.setSuccResult(curEntity);
}
@@ -196,17 +197,17 @@ public class BdbTopicDeployMapperImpl implements
TopicDeployMapper {
@Override
public List<TopicDeployEntity> getTopicConf(TopicDeployEntity qryEntity) {
- List<TopicDeployEntity> retEntitys = new ArrayList<>();
+ List<TopicDeployEntity> retEntities = new ArrayList<>();
if (qryEntity == null) {
- retEntitys.addAll(topicConfCache.values());
+ retEntities.addAll(topicConfCache.values());
} else {
for (TopicDeployEntity entity : topicConfCache.values()) {
if (entity != null && entity.isMatched(qryEntity)) {
- retEntitys.add(entity);
+ retEntities.add(entity);
}
}
}
- return retEntitys;
+ return retEntities;
}
@Override
@@ -409,21 +410,22 @@ public class BdbTopicDeployMapperImpl implements
TopicDeployMapper {
* Put topic configure info into bdb store
*
* @param memEntity need add record
+ * @param strBuff the string buffer
* @param result process result with old value
- * @return
+ * @return the process result
*/
- private boolean putTopicConfig2Bdb(TopicDeployEntity memEntity,
ProcessResult result) {
- BdbTopicConfEntity retData = null;
+ private boolean putTopicConfig2Bdb(TopicDeployEntity memEntity,
+ StringBuilder strBuff, ProcessResult
result) {
BdbTopicConfEntity bdbEntity =
memEntity.buildBdbTopicConfEntity();
try {
- retData = topicConfIndex.put(bdbEntity);
+ topicConfIndex.put(bdbEntity);
} catch (Throwable e) {
logger.error("[BDB Impl] put topic configure failure ", e);
result.setFailResult(DataOpErrCode.DERR_STORE_ABNORMAL.getCode(),
- new StringBuilder(TBaseConstants.BUILDER_DEFAULT_SIZE)
- .append("Put topic configure failure: ")
+ strBuff.append("Put topic configure failure: ")
.append(e.getMessage()).toString());
+ strBuff.delete(0, strBuff.length());
return result.isSuccess();
}
result.setSuccResult(null);
diff --git
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
index d0466a4..186920e 100644
---
a/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
+++
b/inlong-tubemq/tubemq-server/src/main/java/org/apache/inlong/tubemq/server/master/nodemanage/nodebroker/DefBrokerRunManager.java
@@ -227,8 +227,9 @@ public class DefBrokerRunManager implements
BrokerRunManager, AliveObserver {
sBuffer.delete(0, sBuffer.length());
return result.isSuccess();
}
- String brokerConfInfo =
- brokerEntry.getBrokerDefaultConfInfo();
+ brokerEntry.getBrokerDefaultConfInfo(sBuffer);
+ String brokerConfInfo = sBuffer.toString();
+ sBuffer.delete(0, sBuffer.length());
Map<String, String> topicConfInfoMap =
metaDataManager.getBrokerTopicStrConfigInfo(brokerEntry,
sBuffer);
//
@@ -330,7 +331,9 @@ public class DefBrokerRunManager implements
BrokerRunManager, AliveObserver {
BrokerConfEntity brokerConfEntity =
metaDataManager.getBrokerConfByBrokerId(brokerId);
if (brokerConfEntity != null) {
- brokerConfInfo = brokerConfEntity.getBrokerDefaultConfInfo();
+ brokerConfEntity.getBrokerDefaultConfInfo(sBuffer);
+ brokerConfInfo = sBuffer.toString();
+ sBuffer.delete(0, sBuffer.length());
manageStatus = brokerConfEntity.getManageStatus();
}
Map<String, String> brokerTopicSetConfInfo =