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 9c80a9636 [INLONG-3633][TubeMQ] Add topic publish and subscribe 
methods, cluster modification add parameters,cluster multiple masterIP query 
topic failed (#3657)
9c80a9636 is described below

commit 9c80a963668c6b2ccec9842da76ea8d125e56ff2
Author: bluewang <[email protected]>
AuthorDate: Wed Apr 13 16:26:05 2022 +0800

    [INLONG-3633][TubeMQ] Add topic publish and subscribe methods, cluster 
modification add parameters,cluster multiple masterIP query topic failed (#3657)
---
 .../manager/controller/cluster/dto/ClusterDto.java  |  1 +
 .../controller/topic/TopicWebController.java        |  6 ++++++
 .../request/SetPublishReq.java}                     | 21 ++++++++++++---------
 .../request/SetSubscribeReq.java}                   | 21 ++++++++++++---------
 .../tubemq/manager/service/ClusterServiceImpl.java  |  1 +
 .../tubemq/manager/service/MasterServiceImpl.java   |  6 ++----
 .../inlong/tubemq/manager/service/TubeConst.java    |  2 ++
 7 files changed, 36 insertions(+), 22 deletions(-)

diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
index 741ff03d1..0ba7d1760 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
@@ -24,6 +24,7 @@ import org.apache.commons.lang3.StringUtils;
 public class ClusterDto {
     private Long clusterId;
     private String clusterName;
+    private int reloadBrokerSize;
 
     public boolean legal() {
         return (clusterId != null && StringUtils.isNotBlank(clusterName));
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/TopicWebController.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/TopicWebController.java
index 727e64488..5e9c7967a 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/TopicWebController.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/TopicWebController.java
@@ -29,6 +29,8 @@ import 
org.apache.inlong.tubemq.manager.controller.topic.request.DeleteTopicReq;
 import 
org.apache.inlong.tubemq.manager.controller.topic.request.ModifyTopicReq;
 import 
org.apache.inlong.tubemq.manager.controller.topic.request.QueryCanWriteReq;
 import 
org.apache.inlong.tubemq.manager.controller.topic.request.SetAuthControlReq;
+import org.apache.inlong.tubemq.manager.controller.topic.request.SetPublishReq;
+import 
org.apache.inlong.tubemq.manager.controller.topic.request.SetSubscribeReq;
 import org.apache.inlong.tubemq.manager.service.TubeConst;
 import org.apache.inlong.tubemq.manager.service.TubeMQErrorConst;
 import org.apache.inlong.tubemq.manager.service.interfaces.MasterService;
@@ -79,6 +81,10 @@ public class TopicWebController {
                 return masterService.baseRequestMaster(gson.fromJson(req, 
DeleteTopicReq.class));
             case TubeConst.QUERY_CAN_WRITE:
                 return queryCanWrite(gson.fromJson(req, 
QueryCanWriteReq.class));
+            case TubeConst.PUBLISH:
+                return masterService.baseRequestMaster(gson.fromJson(req, 
SetPublishReq.class));
+            case TubeConst.SUBSCRIBE:
+                return masterService.baseRequestMaster(gson.fromJson(req, 
SetSubscribeReq.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/dto/ClusterDto.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetPublishReq.java
similarity index 63%
copy from 
inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
copy to 
inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetPublishReq.java
index 741ff03d1..58b366556 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetPublishReq.java
@@ -15,17 +15,20 @@
  * limitations under the License.
  */
 
-package org.apache.inlong.tubemq.manager.controller.cluster.dto;
+package org.apache.inlong.tubemq.manager.controller.topic.request;
 
 import lombok.Data;
-import org.apache.commons.lang3.StringUtils;
+import lombok.EqualsAndHashCode;
+import lombok.ToString;
+import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
 
 @Data
-public class ClusterDto {
-    private Long clusterId;
-    private String clusterName;
-
-    public boolean legal() {
-        return (clusterId != null && StringUtils.isNotBlank(clusterName));
-    }
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+public class SetPublishReq extends BaseReq {
+    private String topicName;
+    private String brokerId;
+    private String confModAuthToken;
+    private String modifyUser;
+    private boolean acceptPublish;
 }
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetSubscribeReq.java
similarity index 63%
copy from 
inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
copy to 
inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetSubscribeReq.java
index 741ff03d1..27e6925f3 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/dto/ClusterDto.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/topic/request/SetSubscribeReq.java
@@ -15,17 +15,20 @@
  * limitations under the License.
  */
 
-package org.apache.inlong.tubemq.manager.controller.cluster.dto;
+package org.apache.inlong.tubemq.manager.controller.topic.request;
 
 import lombok.Data;
-import org.apache.commons.lang3.StringUtils;
+import lombok.EqualsAndHashCode;
+import lombok.ToString;
+import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
 
 @Data
-public class ClusterDto {
-    private Long clusterId;
-    private String clusterName;
-
-    public boolean legal() {
-        return (clusterId != null && StringUtils.isNotBlank(clusterName));
-    }
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+public class SetSubscribeReq extends BaseReq {
+    private String topicName;
+    private String brokerId;
+    private String confModAuthToken;
+    private String modifyUser;
+    private boolean acceptSubscribe;
 }
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 151367c8e..dd940a92b 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
@@ -98,6 +98,7 @@ public class ClusterServiceImpl implements ClusterService {
             ClusterEntry cluster = clusterRepository
                     .findClusterEntryByClusterId(clusterDto.getClusterId());
             cluster.setClusterName(clusterDto.getClusterName());
+            cluster.setReloadBrokerSize(clusterDto.getReloadBrokerSize());
             clusterRepository.save(cluster);
         } catch (Exception e) {
             return TubeMQResult.errorResult(e.getMessage());
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 4d24133b7..6df1ee34b 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
@@ -106,8 +106,7 @@ public class MasterServiceImpl implements MasterService {
         if (req.getClusterId() == null) {
             return TubeMQResult.errorResult("please input clusterId");
         }
-        MasterEntry masterEntry = 
masterRepository.findMasterEntryByClusterIdEquals(
-                req.getClusterId());
+        MasterEntry masterEntry = 
getMasterNode(Long.valueOf(req.getClusterId()));
         if (masterEntry == null) {
             return TubeMQResult.errorResult("no such cluster");
         }
@@ -171,8 +170,7 @@ public class MasterServiceImpl implements MasterService {
     public String getQueryUrl(Map<String, String> queryBody) throws Exception {
         int clusterId = Integer.parseInt(queryBody.get("clusterId"));
         queryBody.remove("clusterId");
-        MasterEntry masterEntry =
-                masterRepository.findMasterEntryByClusterIdEquals(clusterId);
+        MasterEntry masterEntry = getMasterNode(Long.valueOf(clusterId));
         return TubeConst.SCHEMA + masterEntry.getIp() + ":" + 
masterEntry.getWebPort()
                 + "/" + TubeConst.TUBE_REQUEST_PATH + "?" + 
ConvertUtils.covertMapToQueryString(queryBody);
     }
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 87a1205ac..57acf76ee 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
@@ -52,6 +52,8 @@ public class TubeConst {
     public static final String ADD = "add";
     public static final String QUERY = "query";
     public static final String SWITCH = "switch";
+    public static final String PUBLISH = "publish";
+    public static final String SUBSCRIBE = "subscribe";
     public static final String REBALANCE_CONSUMER_GROUP = "rebalanceGroup";
     public static final String REBALANCE_CONSUMER = "rebalanceConsumer";
     public static final String SET_READ_OR_WRITE = "setReadOrWrite";

Reply via email to