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

Reply via email to