This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new d9162c2d0 [INLONG-3621][TubeMQ] Added cluster switching method and
delete cluster support to delete master (#3630)
d9162c2d0 is described below
commit d9162c2d0f748128fe4fe91b4d221cd6e81aa255
Author: bluewang <[email protected]>
AuthorDate: Tue Apr 12 17:41:28 2022 +0800
[INLONG-3621][TubeMQ] Added cluster switching method and delete cluster
support to delete master (#3630)
---
.../controller/cluster/ClusterController.java | 3 +++
.../cluster/request/SwitchClusterReq.java | 30 ++++++++++++++++++++++
.../manager/repository/MasterRepository.java | 7 +++++
.../tubemq/manager/service/ClusterServiceImpl.java | 5 ++++
.../tubemq/manager/service/MasterServiceImpl.java | 9 +++++++
.../inlong/tubemq/manager/service/TubeConst.java | 1 +
.../manager/service/interfaces/MasterService.java | 8 ++++++
7 files changed, 63 insertions(+)
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 5530e4691..99c1903e0 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
@@ -29,6 +29,7 @@ import
org.apache.inlong.tubemq.manager.controller.TubeMQResult;
import org.apache.inlong.tubemq.manager.controller.cluster.dto.ClusterDto;
import
org.apache.inlong.tubemq.manager.controller.cluster.request.AddClusterReq;
import
org.apache.inlong.tubemq.manager.controller.cluster.request.DeleteClusterReq;
+import
org.apache.inlong.tubemq.manager.controller.cluster.request.SwitchClusterReq;
import org.apache.inlong.tubemq.manager.controller.cluster.vo.ClusterVo;
import
org.apache.inlong.tubemq.manager.controller.group.result.ConsumerGroupInfoRes;
import
org.apache.inlong.tubemq.manager.controller.group.result.ConsumerInfoRes;
@@ -81,6 +82,8 @@ public class ClusterController {
return deleteCluster(gson.fromJson(req,
DeleteClusterReq.class));
case TubeConst.MODIFY:
return changeCluster(gson.fromJson(req, ClusterDto.class));
+ case TubeConst.SWITCH:
+ return masterService.baseRequestMaster(gson.fromJson(req,
SwitchClusterReq.class));
default:
return
TubeMQResult.errorResult(TubeMQErrorConst.NO_SUCH_METHOD);
}
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
new file mode 100644
index 000000000..d03cc5ef6
--- /dev/null
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
@@ -0,0 +1,30 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.tubemq.manager.controller.cluster.request;
+
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.ToString;
+import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
+
+@Data
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+public class SwitchClusterReq extends BaseReq {
+ private String confModAuthToken;
+}
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
index e0b653664..54d538115 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
@@ -56,4 +56,11 @@ public interface MasterRepository extends
JpaRepository<MasterEntry, Long> {
* @return
*/
List<MasterEntry> findMasterEntryByIpEquals(String masterIp);
+
+ /**
+ * delete master by cluster id
+ *
+ * @return
+ */
+ Integer deleteByClusterId(Long clusterId);
}
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 4fcd2963b..151367c8e 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
@@ -31,6 +31,7 @@ import org.apache.inlong.tubemq.manager.entry.ClusterEntry;
import org.apache.inlong.tubemq.manager.entry.MasterEntry;
import org.apache.inlong.tubemq.manager.repository.ClusterRepository;
import org.apache.inlong.tubemq.manager.service.interfaces.ClusterService;
+import org.apache.inlong.tubemq.manager.service.interfaces.MasterService;
import org.apache.inlong.tubemq.manager.service.interfaces.NodeService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
@@ -45,6 +46,9 @@ public class ClusterServiceImpl implements ClusterService {
@Autowired
NodeService nodeService;
+ @Autowired
+ MasterService masterService;
+
@Override
@Transactional(rollbackOn = Exception.class)
public void addClusterAndMasterNode(AddClusterReq req) {
@@ -64,6 +68,7 @@ public class ClusterServiceImpl implements ClusterService {
@Override
@Transactional(rollbackOn = Exception.class)
public void deleteCluster(Long clusterId) {
+ masterService.deleteMaster(clusterId);
Integer successCode = clusterRepository.deleteByClusterId(clusterId);
if (successCode.equals(DELETE_FAIL)) {
throw new RuntimeException("no such cluster with clusterId = " +
clusterId);
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
index 3d6379544..4d24133b7 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
@@ -37,6 +37,7 @@ import
org.apache.inlong.tubemq.manager.controller.TubeMQResult;
import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
import org.apache.inlong.tubemq.manager.entry.MasterEntry;
import org.apache.inlong.tubemq.manager.repository.MasterRepository;
+import static org.apache.inlong.tubemq.manager.service.TubeConst.DELETE_FAIL;
import org.apache.inlong.tubemq.manager.service.interfaces.MasterService;
import org.apache.inlong.tubemq.manager.service.tube.TubeHttpResponse;
import org.apache.inlong.tubemq.manager.utils.ConvertUtils;
@@ -189,4 +190,12 @@ public class MasterServiceImpl implements MasterService {
+ method + "&" + "clusterId=" + clusterId;
}
+ @Override
+ public void deleteMaster(Long clusterId) {
+ Integer successCode = masterRepository.deleteByClusterId(clusterId);
+ if (successCode.equals(DELETE_FAIL)) {
+ throw new RuntimeException("no such master with clusterId = " +
clusterId);
+ }
+ }
+
}
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
index ef851f9b8..87a1205ac 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
@@ -51,6 +51,7 @@ public class TubeConst {
public static final String CLONE = "clone";
public static final String ADD = "add";
public static final String QUERY = "query";
+ public static final String SWITCH = "switch";
public static final String REBALANCE_CONSUMER_GROUP = "rebalanceGroup";
public static final String REBALANCE_CONSUMER = "rebalanceConsumer";
public static final String SET_READ_OR_WRITE = "setReadOrWrite";
diff --git
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
index 0ac583b24..3e062fc95 100644
---
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
+++
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
@@ -101,4 +101,12 @@ public interface MasterService {
TubeMQResult checkMasterNodeStatus(String masterIp, Integer masterPort);
String getQueryCountUrl(Integer clusterId, String method);
+
+ /**
+ * delete master by cluster id
+ *
+ * @param clusterId
+ */
+ void deleteMaster(Long clusterId);
+
}