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 bcbea4f [INLONG-3002] Tubemq new cluster supports multi-master port
and master webport configuration (#3006)
bcbea4f is described below
commit bcbea4fc5d41d0b9b01301ae7f4768630e33b0e4
Author: bluewang <[email protected]>
AuthorDate: Tue Mar 8 20:23:34 2022 +0800
[INLONG-3002] Tubemq new cluster supports multi-master port and master
webport configuration (#3006)
Co-authored-by: v_lizhwang <[email protected]>
---
.../manager/controller/cluster/ClusterController.java | 7 ++++---
.../manager/controller/cluster/request/AddClusterReq.java | 8 ++++----
.../apache/inlong/tubemq/manager/entry/ClusterEntry.java | 3 ---
.../inlong/tubemq/manager/service/ClusterServiceImpl.java | 14 ++++++++------
.../tubemq/manager/controller/TestClusterController.java | 14 ++++++++++----
5 files changed, 26 insertions(+), 20 deletions(-)
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
index e15cc1e..92ff640 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
@@ -98,9 +98,10 @@ public class ClusterController {
if (!req.legal()) {
return TubeMQResult.errorResult(TubeMQErrorConst.PARAM_ILLEGAL);
}
- List<String> masterIps = req.getMasterIps();
- for (String masterIp : masterIps) {
- TubeMQResult checkResult =
masterService.checkMasterNodeStatus(masterIp, req.getMasterWebPort());
+ List<MasterEntry> masterEntries = req.getMasterEntries();
+ for (MasterEntry masterEntry : masterEntries) {
+ TubeMQResult checkResult =
masterService.checkMasterNodeStatus(masterEntry.getIp(),
+ masterEntry.getWebPort());
if (checkResult.getErrCode() != SUCCESS_CODE) {
return TubeMQResult.errorResult("please check master ip and
webPort");
}
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/AddClusterReq.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/AddClusterReq.java
index 2eafbac..4db0bf2 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/AddClusterReq.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/AddClusterReq.java
@@ -20,19 +20,19 @@ package
org.apache.inlong.tubemq.manager.controller.cluster.request;
import lombok.Data;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
+import org.apache.inlong.tubemq.manager.entry.MasterEntry;
import java.util.List;
@Data
public class AddClusterReq {
- private List<String> masterIps;
+ private Integer id;
private String clusterName;
- private Integer masterPort;
- private Integer masterWebPort;
+ private List<MasterEntry> masterEntries;
private String createUser;
private String token;
public boolean legal() {
- return CollectionUtils.isNotEmpty(masterIps) && masterPort != null &&
StringUtils.isNotBlank(token);
+ return CollectionUtils.isNotEmpty(masterEntries) &&
StringUtils.isNotBlank(token);
}
}
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/entry/ClusterEntry.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/entry/ClusterEntry.java
index 11117f1..5ee2476 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/entry/ClusterEntry.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/entry/ClusterEntry.java
@@ -19,8 +19,6 @@ package org.apache.inlong.tubemq.manager.entry;
import java.util.Date;
import javax.persistence.Entity;
-import javax.persistence.GeneratedValue;
-import javax.persistence.GenerationType;
import javax.persistence.Id;
import javax.persistence.Table;
import javax.persistence.UniqueConstraint;
@@ -36,7 +34,6 @@ import lombok.Data;
@Data
public class ClusterEntry {
@Id
- @GeneratedValue(strategy = GenerationType.IDENTITY)
private long clusterId;
private String clusterName;
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
index 83f2209..7f09de3 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
@@ -49,6 +49,9 @@ public class ClusterServiceImpl implements ClusterService {
@Transactional(rollbackOn = Exception.class)
public void addClusterAndMasterNode(AddClusterReq req) {
ClusterEntry entry = new ClusterEntry();
+ if (req.getId() != null) {
+ entry.setClusterId(req.getId());
+ }
entry.setCreateTime(new Date());
entry.setCreateUser(req.getCreateUser());
entry.setClusterName(req.getClusterName());
@@ -95,12 +98,11 @@ public class ClusterServiceImpl implements ClusterService {
if (clusterEntry == null) {
return;
}
- for (String masterIp : req.getMasterIps()) {
- MasterEntry masterEntry = new MasterEntry();
- masterEntry.setPort(req.getMasterPort());
- masterEntry.setClusterId(clusterEntry.getClusterId());
- masterEntry.setWebPort(req.getMasterWebPort());
- masterEntry.setIp(masterIp);
+ for (MasterEntry masterEntry : req.getMasterEntries()) {
+ masterEntry.setPort(masterEntry.getPort());
+ masterEntry.setClusterId(req.getId());
+ masterEntry.setWebPort(masterEntry.getWebPort());
+ masterEntry.setIp(masterEntry.getIp());
masterEntry.setToken(req.getToken());
nodeService.addNode(masterEntry);
}
diff --git
a/inlong-tubemq/tubemq-manager/src/test/java/org/apache/inlong/tubemq/manager/controller/TestClusterController.java
b/inlong-tubemq/tubemq-manager/src/test/java/org/apache/inlong/tubemq/manager/controller/TestClusterController.java
index a813167..8548c74 100644
---
a/inlong-tubemq/tubemq-manager/src/test/java/org/apache/inlong/tubemq/manager/controller/TestClusterController.java
+++
b/inlong-tubemq/tubemq-manager/src/test/java/org/apache/inlong/tubemq/manager/controller/TestClusterController.java
@@ -46,7 +46,8 @@ import org.springframework.test.web.servlet.MockMvc;
import org.springframework.test.web.servlet.MvcResult;
import org.springframework.test.web.servlet.RequestBuilder;
-import java.util.Collections;
+import java.util.ArrayList;
+import java.util.List;
@Slf4j
@RunWith(SpringRunner.class)
@@ -158,10 +159,15 @@ public class TestClusterController {
public void testAddCluster() throws Exception {
AddClusterReq req = new AddClusterReq();
+ req.setId(4);
req.setClusterName("test");
- req.setMasterIps(Collections.singletonList("127.0.0.1"));
- req.setMasterWebPort(8080);
- req.setMasterPort(8089);
+ MasterEntry masterEntry = new MasterEntry();
+ masterEntry.setIp("127.0.0.1");
+ masterEntry.setPort(8089);
+ masterEntry.setWebPort(8080);
+ List<MasterEntry> masterEntries = new ArrayList<>();
+ masterEntries.add(masterEntry);
+ req.setMasterEntries(masterEntries);
req.setToken("abc");
ClusterEntry entry = getOneClusterEntry();