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/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 19a19d4ed5 [INLONG-8354][Manager] Support previewing data of TubeMQ
(#8357)
19a19d4ed5 is described below
commit 19a19d4ed59c15bc5f92fef753fd0626148b2c9e
Author: fuweng11 <[email protected]>
AuthorDate: Fri Jun 30 10:34:47 2023 +0800
[INLONG-8354][Manager] Support previewing data of TubeMQ (#8357)
Co-authored-by: healchow <[email protected]>
---
.../manager/common/consts/InlongConstants.java | 11 ++++
.../manager/pojo/queue/tubemq/TubeBrokerInfo.java | 27 +++++++-
.../pojo/queue/tubemq/TubeMessageResponse.java | 48 +++++++++++++++
.../resource/queue/QueueResourceOperator.java | 11 ++--
.../resource/queue/pulsar/PulsarOperator.java | 20 +++---
.../queue/pulsar/PulsarResourceOperator.java | 6 +-
.../resource/queue/tubemq/TubeMQOperator.java | 71 ++++++++++++++++++++++
.../queue/tubemq/TubeMQResourceOperator.java | 14 +++++
.../service/stream/InlongStreamServiceImpl.java | 2 +-
9 files changed, 187 insertions(+), 23 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
index 3ce5d5e9e8..1f45d451d4 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
@@ -17,6 +17,8 @@
package org.apache.inlong.manager.common.consts;
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+
import com.google.common.collect.Sets;
import java.util.Set;
@@ -48,6 +50,8 @@ public class InlongConstants {
public static final String COLON = ":";
+ public static final String EQUAL = "=";
+
public static final String SEMICOLON = ";";
public static final String HYPHEN = "-";
@@ -187,4 +191,11 @@ public class InlongConstants {
*/
public static final String BATCH_PARSING_FILED_JSON_COMMENT_PROP = "desc";
+ /**
+ * Message compression type, 0: Raw message without any InLong format, 1:
InlongMsgPb, 2: InlongMsg
+ * <p/>
+ * See more: {@link DataProxyMsgEncType}
+ */
+ public static final String MSG_ENCODE_VER = "msgEnType";
+
}
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeBrokerInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeBrokerInfo.java
index 61658a0bb1..d1131b52c4 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeBrokerInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeBrokerInfo.java
@@ -17,7 +17,10 @@
package org.apache.inlong.manager.pojo.queue.tubemq;
+import org.apache.inlong.manager.common.consts.InlongConstants;
+
import lombok.Data;
+import org.apache.commons.collections.CollectionUtils;
import java.util.ArrayList;
import java.util.List;
@@ -113,6 +116,25 @@ public class TubeBrokerInfo {
return allIdList;
}
+ /**
+ * Get one online broker address of TubeMQ.
+ */
+ public String getOnlineBrokerAddress() {
+ if (CollectionUtils.isEmpty(data)) {
+ return null;
+ }
+
+ String brokerAddress = null;
+ for (BrokerInfo brokerInfo : data) {
+ if (brokerInfo.isBrokerOnline()) {
+ brokerAddress = brokerInfo.getBrokerIp() +
InlongConstants.COLON + brokerInfo.getBrokerWebPort();
+ break;
+ }
+ }
+
+ return brokerAddress;
+ }
+
/**
* Broker info
*/
@@ -122,6 +144,7 @@ public class TubeBrokerInfo {
private int brokerId;
private String brokerIp;
private int brokerPort;
+ private int brokerWebPort;
private String manageStatus;
private String runStatus;
private String subStatus;
@@ -138,8 +161,8 @@ public class TubeBrokerInfo {
}
private boolean isWorking() {
- return RUNNING.equals(runStatus) && (ONLINE.equals(manageStatus)
|| ONLY_READ.equals(manageStatus)
- || ONLY_WRITE.equals(manageStatus));
+ return RUNNING.equals(runStatus) &&
+ (ONLINE.equals(manageStatus) ||
ONLY_READ.equals(manageStatus) || ONLY_WRITE.equals(manageStatus));
}
private boolean isConfigurable() {
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeMessageResponse.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeMessageResponse.java
new file mode 100644
index 0000000000..723c0d5469
--- /dev/null
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/queue/tubemq/TubeMessageResponse.java
@@ -0,0 +1,48 @@
+/*
+ * 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.manager.pojo.queue.tubemq;
+
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * The message info from TubeMQ.
+ */
+@Data
+public class TubeMessageResponse {
+
+ // true, or false
+ private boolean result;
+
+ // 0 is success, other is failed
+ private int errCode;
+
+ // OK, or err msg
+ private String errMsg;
+
+ private List<TubeDataInfo> dataSet;
+
+ @Data
+ public static class TubeDataInfo {
+
+ private int index;
+ private String data;
+ private String attr;
+ }
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/QueueResourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/QueueResourceOperator.java
index b53045e9cb..805e92b7fc 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/QueueResourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/QueueResourceOperator.java
@@ -21,8 +21,6 @@ import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
-import org.apache.pulsar.client.api.PulsarClientException;
-
import javax.validation.constraints.NotBlank;
import javax.validation.constraints.NotNull;
@@ -77,15 +75,16 @@ public interface QueueResourceOperator {
}
/**
- * Query brief mq message info
+ * Query latest messages from MQ.
*
* @param groupInfo inlong group info
* @param streamInfo inlong stream info
- * @param messageCount Count of messages to query'
+ * @param messageCount count of messages to query
+ * @throws Exception any exception if occurred
* @return query brief mq message info
*/
- default List<BriefMQMessage> queryLastestMessage(InlongGroupInfo
groupInfo, InlongStreamInfo streamInfo,
- Integer messageCount) throws PulsarClientException {
+ default List<BriefMQMessage> queryLatestMessages(InlongGroupInfo
groupInfo, InlongStreamInfo streamInfo,
+ Integer messageCount) throws Exception {
return null;
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarOperator.java
index bf6342cdb7..e81cdd56cf 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarOperator.java
@@ -50,8 +50,6 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Map;
-import static org.elasticsearch.ingest.Pipeline.VERSION_KEY;
-
/**
* Pulsar operator, supports creating topics and creating subscription.
*/
@@ -388,34 +386,34 @@ public class PulsarOperator {
/**
* Query topic message for the given pulsar cluster.
*/
- public List<BriefMQMessage> queryLastestMessage(PulsarAdmin pulsarAdmin,
String topicFullName, String subName,
+ public List<BriefMQMessage> queryLatestMessage(PulsarAdmin pulsarAdmin,
String topicFullName, String subName,
Integer messageCount, InlongStreamInfo streamInfo) {
LOGGER.info("begin to query message for topic {}, subName={}",
topicFullName, subName);
- List<Message<byte[]>> messages = new ArrayList<>();
- List<BriefMQMessage> messageList = new ArrayList<>();
+ List<Message<byte[]>> messages;
try {
messages = pulsarAdmin.topics().peekMessages(topicFullName,
subName, messageCount);
} catch (PulsarAdminException e) {
- String errMsg = "failed to query peek messages";
+ String errMsg = "failed to query peek messages: ";
LOGGER.error(errMsg, e);
- throw new BusinessException(errMsg);
+ throw new BusinessException(errMsg + e.getMessage());
}
int index = 0;
+ List<BriefMQMessage> messageList = new ArrayList<>();
for (Message<byte[]> pulsarMessage : messages) {
try {
Map<String, String> headers = pulsarMessage.getProperties();
- int wrapTypeId =
Integer.parseInt(headers.getOrDefault(VERSION_KEY,
+ int wrapTypeId =
Integer.parseInt(headers.getOrDefault(InlongConstants.MSG_ENCODE_VER,
Integer.toString(DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getId())));
DeserializeOperator deserializeOperator =
deserializeOperatorFactory.getInstance(
DataProxyMsgEncType.valueOf(wrapTypeId));
messageList.addAll(
deserializeOperator.decodeMsg(streamInfo,
pulsarMessage.getData(), headers, ++index));
} catch (Exception e) {
- String errMsg = "decode msg error";
- LOGGER.error("decode msg error", e);
- throw new BusinessException(errMsg);
+ String errMsg = "decode msg error: ";
+ LOGGER.error(errMsg, e);
+ throw new BusinessException(errMsg + e.getMessage());
}
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarResourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarResourceOperator.java
index d3a79c9714..f08faeaced 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarResourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/pulsar/PulsarResourceOperator.java
@@ -305,9 +305,9 @@ public class PulsarResourceOperator implements
QueueResourceOperator {
}
/**
- * Query lastest message from pulsar
+ * Query latest message from pulsar
*/
- public List<BriefMQMessage> queryLastestMessage(InlongGroupInfo groupInfo,
+ public List<BriefMQMessage> queryLatestMessages(InlongGroupInfo groupInfo,
InlongStreamInfo streamInfo, Integer messageCount) throws
PulsarClientException {
String groupId = streamInfo.getInlongGroupId();
InlongPulsarInfo inlongPulsarInfo = ((InlongPulsarInfo) groupInfo);
@@ -327,7 +327,7 @@ public class PulsarResourceOperator implements
QueueResourceOperator {
String clusterTag = inlongPulsarInfo.getInlongClusterTag();
String subs = String.format(PULSAR_SUBSCRIPTION_REALTIME_REVIEW,
clusterTag, topicName);
briefMQMessages =
- pulsarOperator.queryLastestMessage(pulsarAdmin,
fullTopicName, subs, messageCount, streamInfo);
+ pulsarOperator.queryLatestMessage(pulsarAdmin,
fullTopicName, subs, messageCount, streamInfo);
// insert the consumer group info into the inlong_consume table
Integer id = consumeService.saveBySystem(groupInfo, topicName,
subs);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQOperator.java
index da2948cdec..06193a49bd 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQOperator.java
@@ -17,14 +17,22 @@
package org.apache.inlong.manager.service.resource.queue.tubemq;
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.util.HttpUtils;
import org.apache.inlong.manager.pojo.cluster.tubemq.TubeClusterInfo;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.queue.tubemq.ConsumerGroupResponse;
import org.apache.inlong.manager.pojo.queue.tubemq.TopicResponse;
import org.apache.inlong.manager.pojo.queue.tubemq.TubeBrokerInfo;
import org.apache.inlong.manager.pojo.queue.tubemq.TubeHttpResponse;
+import org.apache.inlong.manager.pojo.queue.tubemq.TubeMessageResponse;
+import
org.apache.inlong.manager.pojo.queue.tubemq.TubeMessageResponse.TubeDataInfo;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import org.apache.inlong.manager.service.cluster.InlongClusterServiceImpl;
+import org.apache.inlong.manager.service.message.DeserializeOperator;
+import org.apache.inlong.manager.service.message.DeserializeOperatorFactory;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
@@ -38,7 +46,11 @@ import org.springframework.web.client.RestTemplate;
import javax.annotation.Nonnull;
+import java.util.ArrayList;
+import java.util.Base64;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
/**
* TubeMQ operator, supports creating topics and creating consumer groups.
@@ -58,15 +70,19 @@ public class TubeMQOperator {
private static final String BROKER_ID = "&brokerId=";
private static final String CREATE_USER = "&createUser=";
private static final String CONF_MOD_AUTH_TOKEN = "&confModAuthToken=";
+ private static final String MSG_COUNT = "&msgCount=";
private static final String QUERY_TOPIC_PATH =
"/webapi.htm?method=admin_query_cluster_topic_view";
private static final String QUERY_BROKER_PATH =
"/webapi.htm?method=admin_query_broker_run_status";
private static final String ADD_TOPIC_PATH =
"/webapi.htm?method=admin_add_new_topic_record";
private static final String QUERY_CONSUMER_PATH =
"/webapi.htm?method=admin_query_allowed_consumer_group_info";
private static final String ADD_CONSUMER_PATH =
"/webapi.htm?method=admin_add_authorized_consumergroup_info";
+ private static final String QUERY_MESSAGE_PATH =
"/broker.htm?method=admin_snapshot_message";
@Autowired
private RestTemplate restTemplate;
+ @Autowired
+ public DeserializeOperatorFactory deserializeOperatorFactory;
/**
* Create topic for the given tubemq cluster.
@@ -248,4 +264,59 @@ public class TubeMQOperator {
}
}
+ /**
+ * Query topic message for the given tubemq cluster.
+ */
+ public List<BriefMQMessage> queryLastMessage(TubeClusterInfo tubeCluster,
String topicName,
+ Integer msgCount, InlongStreamInfo streamInfo) {
+ LOGGER.info("begin to query message for topic {} in cluster: {}",
topicName, tubeCluster);
+ String masterUrl = tubeCluster.getMasterWebUrl();
+ TubeBrokerInfo brokerView = this.getBrokerInfo(masterUrl);
+ String brokerUrl = brokerView.getOnlineBrokerAddress();
+
+ List<BriefMQMessage> messageList = new ArrayList<>();
+ try {
+ if (StringUtils.isEmpty(brokerUrl) ||
StringUtils.isEmpty(topicName)) {
+ throw new BusinessException("tubemq master url or tubemq topic
cannot be null");
+ }
+
+ if (!this.isTopicExist(masterUrl, topicName)) {
+ LOGGER.error("tubemq topic {} not exists in {}, skip to
query", topicName, masterUrl);
+ throw new BusinessException("TubeMQ master url or TubeMQ topic
cannot be null");
+ }
+
+ String url = "http://" + brokerUrl + QUERY_MESSAGE_PATH +
TOPIC_NAME + topicName + MSG_COUNT + msgCount;
+ TubeMessageResponse response = HttpUtils.request(restTemplate,
url, HttpMethod.GET,
+ null, new HttpHeaders(), TubeMessageResponse.class);
+ if (response.getErrCode() != SUCCESS_CODE) {
+ String msg = String.format("failed to query message for topic
%s, error: %s",
+ topicName, response.getErrMsg());
+ LOGGER.error(msg + " in {} for broker {}", masterUrl,
brokerUrl);
+ throw new BusinessException(msg);
+ }
+
+ int index = 0;
+ for (TubeDataInfo tubeDataInfo : response.getDataSet()) {
+ Map<String, String> map = new HashMap<>();
+ for (String kv :
tubeDataInfo.getAttr().split(InlongConstants.COMMA)) {
+ map.put(kv.split(InlongConstants.EQUAL)[0],
kv.split(InlongConstants.EQUAL)[1]);
+ }
+
+ int wrapTypeId =
Integer.parseInt(map.getOrDefault(InlongConstants.MSG_ENCODE_VER,
+
Integer.toString(DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getId())));
+ byte[] messageData =
Base64.getDecoder().decode(tubeDataInfo.getData());
+ DeserializeOperator deserializeOperator =
deserializeOperatorFactory.getInstance(
+ DataProxyMsgEncType.valueOf(wrapTypeId));
+ messageList.addAll(deserializeOperator.decodeMsg(streamInfo,
messageData, map, index));
+ }
+
+ LOGGER.info("success query messages for topic={}", topicName);
+ } catch (Exception e) {
+ String errMsg = "failed to query messages: ";
+ LOGGER.error(errMsg, e);
+ throw new BusinessException(errMsg + e.getMessage());
+ }
+
+ return messageList;
+ }
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQResourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQResourceOperator.java
index 114f2fd979..1baa970179 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQResourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/queue/tubemq/TubeMQResourceOperator.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.cluster.tubemq.TubeClusterInfo;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import org.apache.inlong.manager.service.cluster.InlongClusterService;
@@ -35,6 +36,8 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import java.util.List;
+
/**
* Operator for create TubeMQ Topic and ConsumerGroup
*/
@@ -108,4 +111,15 @@ public class TubeMQResourceOperator implements
QueueResourceOperator {
// currently, not support delete tubemq resource for stream
}
+ public List<BriefMQMessage> queryLatestMessages(InlongGroupInfo groupInfo,
InlongStreamInfo streamInfo,
+ Integer messageCount) {
+ Preconditions.expectNotNull(groupInfo, "inlong group info cannot be
null");
+
+ String clusterTag = groupInfo.getInlongClusterTag();
+ TubeClusterInfo tubeCluster = (TubeClusterInfo)
clusterService.getOne(clusterTag, null, ClusterType.TUBEMQ);
+ String topicName = groupInfo.getMqResource();
+
+ return tubeMQOperator.queryLastMessage(tubeCluster, topicName,
messageCount, streamInfo);
+ }
+
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
index 0ae4f16b25..fcceb0c0b2 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
@@ -1016,7 +1016,7 @@ public class InlongStreamServiceImpl implements
InlongStreamService {
List<BriefMQMessage> messageList = new ArrayList<>();
QueueResourceOperator queueOperator =
queueOperatorFactory.getInstance(groupEntity.getMqType());
try {
- messageList = queueOperator.queryLastestMessage(groupInfo,
inlongStreamInfo, messageCount);
+ messageList = queueOperator.queryLatestMessages(groupInfo,
inlongStreamInfo, messageCount);
} catch (Exception e) {
LOGGER.error("query message error ", e);
}