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";