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 7e47c0592f [INLONG-8300][Manager] Support previewing data of Pulsar
(#8301)
7e47c0592f is described below
commit 7e47c0592fd0d9dc4b413c53061bbc6188d949d3
Author: fuweng11 <[email protected]>
AuthorDate: Wed Jun 28 17:45:48 2023 +0800
[INLONG-8300][Manager] Support previewing data of Pulsar (#8301)
---
.../org/apache/inlong/common}/util/StringUtil.java | 2 +-
.../java/org/apache/inlong/common}/util/Utils.java | 2 +-
.../api/inner/client/InlongStreamClient.java | 10 +++
.../client/api/service/InlongStreamApi.java | 5 ++
.../manager/pojo/consume/BriefMQMessage.java | 55 +++++++++++++
inlong-manager/manager-service/pom.xml | 6 +-
.../service/message/DeserializeOperator.java | 61 +++++++++++++++
.../message/DeserializeOperatorFactory.java | 49 ++++++++++++
.../message/InlongMsgDeserializeOperator.java | 82 ++++++++++++++++++++
.../service/message/PbMsgDeserializeOperator.java | 89 ++++++++++++++++++++++
.../service/message/RawMsgDeserializeOperator.java | 53 +++++++++++++
.../resource/queue/QueueResourceOperator.java | 18 +++++
.../resource/queue/pulsar/PulsarOperator.java | 52 ++++++++++++-
.../queue/pulsar/PulsarResourceOperator.java | 38 +++++++++
.../service/stream/InlongStreamService.java | 13 ++++
.../service/stream/InlongStreamServiceImpl.java | 30 ++++++++
.../web/controller/InlongStreamController.java | 14 ++++
.../sdk/sort/impl/decode/MessageDeserializer.java | 4 +-
.../sdk/sort/manager/InlongMultiTopicManager.java | 2 +-
.../sdk/sort/manager/InlongSingleTopicManager.java | 2 +-
.../sort/impl/decode/MessageDeserializerTest.java | 2 +-
21 files changed, 580 insertions(+), 9 deletions(-)
diff --git
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/StringUtil.java
b/inlong-common/src/main/java/org/apache/inlong/common/util/StringUtil.java
similarity index 99%
rename from
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/StringUtil.java
rename to
inlong-common/src/main/java/org/apache/inlong/common/util/StringUtil.java
index a1fc09ccf3..cddef9dc8f 100644
---
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/StringUtil.java
+++ b/inlong-common/src/main/java/org/apache/inlong/common/util/StringUtil.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.inlong.sdk.sort.util;
+package org.apache.inlong.common.util;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
diff --git
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/Utils.java
b/inlong-common/src/main/java/org/apache/inlong/common/util/Utils.java
similarity index 99%
rename from
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/Utils.java
rename to inlong-common/src/main/java/org/apache/inlong/common/util/Utils.java
index 01074a7cfa..a8009d8c52 100644
---
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/util/Utils.java
+++ b/inlong-common/src/main/java/org/apache/inlong/common/util/Utils.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.inlong.sdk.sort.util;
+package org.apache.inlong.common.util;
import org.xerial.snappy.Snappy;
diff --git
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/inner/client/InlongStreamClient.java
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/inner/client/InlongStreamClient.java
index 9a7fe113e5..02f54cb1f7 100644
---
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/inner/client/InlongStreamClient.java
+++
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/inner/client/InlongStreamClient.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.common.PageResult;
import org.apache.inlong.manager.pojo.common.Response;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.sink.ParseFieldRequest;
import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
@@ -246,4 +247,13 @@ public class InlongStreamClient {
ParseFieldRequest request =
ParseFieldRequest.builder().method(method).statement(statement).build();
return parseFields(request);
}
+
+ public List<BriefMQMessage> queryMessage(String groupId, String streamId,
Integer messageCount) {
+ Preconditions.expectNotBlank(groupId, ErrorCodeEnum.GROUP_ID_IS_EMPTY);
+ Preconditions.expectNotBlank(streamId,
ErrorCodeEnum.STREAM_ID_IS_EMPTY);
+ Response<List<BriefMQMessage>> response = ClientUtils.executeHttpCall(
+ inlongStreamApi.listMessages(groupId, streamId, messageCount));
+ ClientUtils.assertRespSuccess(response);
+ return response.getData();
+ }
}
diff --git
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/service/InlongStreamApi.java
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/service/InlongStreamApi.java
index 5aa8e9ed77..e23e6bc5bc 100644
---
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/service/InlongStreamApi.java
+++
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/service/InlongStreamApi.java
@@ -19,6 +19,7 @@ package org.apache.inlong.manager.client.api.service;
import org.apache.inlong.manager.pojo.common.PageResult;
import org.apache.inlong.manager.pojo.common.Response;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.sink.ParseFieldRequest;
import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
@@ -73,4 +74,8 @@ public interface InlongStreamApi {
@POST("stream/parseFields")
Call<Response<List<StreamField>>> parseFields(@Body ParseFieldRequest
parseFieldRequest);
+
+ @GET("stream/listMessages")
+ Call<Response<List<BriefMQMessage>>> listMessages(@Query("groupId") String
groupId,
+ @Query("streamId") String streamId, @Query("messageCount") Integer
messageCount);
}
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/consume/BriefMQMessage.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/consume/BriefMQMessage.java
new file mode 100644
index 0000000000..149d0a2353
--- /dev/null
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/consume/BriefMQMessage.java
@@ -0,0 +1,55 @@
+/*
+ * 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.consume;
+
+import io.swagger.annotations.ApiModel;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+/**
+ * Inlong display message info
+ */
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+@ApiModel("Inlong brief mq message info")
+public class BriefMQMessage {
+
+ @ApiModelProperty(value = "index id")
+ private Integer id;
+
+ @ApiModelProperty(value = "Inlong group id")
+ private String inlongGroupId;
+
+ @ApiModelProperty(value = "Inlong stream id")
+ private String inlongStreamId;
+
+ @ApiModelProperty(value = "Date")
+ private Long dt;
+
+ @ApiModelProperty(value = "Node ip")
+ private String nodeIp;
+
+ @ApiModelProperty(value = "Message body")
+ private String body;
+
+}
diff --git a/inlong-manager/manager-service/pom.xml
b/inlong-manager/manager-service/pom.xml
index 4a4375526c..71aabc248e 100644
--- a/inlong-manager/manager-service/pom.xml
+++ b/inlong-manager/manager-service/pom.xml
@@ -609,6 +609,10 @@
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
</dependency>
-
+ <dependency>
+ <groupId>org.apache.inlong</groupId>
+ <artifactId>sdk-common</artifactId>
+ <version>${project.version}</version>
+ </dependency>
</dependencies>
</project>
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperator.java
new file mode 100644
index 0000000000..876c828f12
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperator.java
@@ -0,0 +1,61 @@
+/*
+ * 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.service.message;
+
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Deserialize of message operator
+ */
+public interface DeserializeOperator {
+
+ String COMPRESS_TYPE_KEY = "compressType";
+ String NODE_IP = "NodeIP";
+ String MSG_TIME_KEY = "msgTime";
+ char INLONGMSG_ATTR_ENTRY_DELIMITER = '&';
+ char INLONGMSG_ATTR_KV_DELIMITER = '=';
+
+ // keys in attributes
+ String INLONGMSG_ATTR_TIME_T = "t";
+ String INLONGMSG_ATTR_TIME_DT = "dt";
+
+ /**
+ * Determines whether the current instance matches the specified type.
+ */
+ boolean accept(DataProxyMsgEncType type);
+
+ /**
+ * List brief mq message info
+ *
+ * @param streamInfo inlong stream info
+ * @param msgBytes messages
+ * @param headers message headers
+ * @param index message index
+ * @return list of brief mq message info
+ */
+ default List<BriefMQMessage> decodeMsg(InlongStreamInfo streamInfo,
+ byte[] msgBytes, Map<String, String> headers, int index) throws
Exception {
+ return null;
+ }
+
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperatorFactory.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperatorFactory.java
new file mode 100644
index 0000000000..6d7d43e82e
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/DeserializeOperatorFactory.java
@@ -0,0 +1,49 @@
+/*
+ * 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.service.message;
+
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.List;
+
+/**
+ * Factory for {@link DeserializeOperator}.
+ */
+@Service
+public class DeserializeOperatorFactory {
+
+ @Autowired
+ private List<DeserializeOperator> operatorList;
+
+ /**
+ * Get a message queue resource operator instance via the given mqType
+ */
+ public DeserializeOperator getInstance(DataProxyMsgEncType type) {
+ return operatorList.stream()
+ .filter(inst -> inst.accept(type))
+ .findFirst()
+ .orElseThrow(() -> new
BusinessException(ErrorCodeEnum.MQ_TYPE_NOT_SUPPORT,
+
String.format(ErrorCodeEnum.MQ_TYPE_NOT_SUPPORT.getMessage(), type)));
+ }
+
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/InlongMsgDeserializeOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/InlongMsgDeserializeOperator.java
new file mode 100644
index 0000000000..c72c6ea4d7
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/InlongMsgDeserializeOperator.java
@@ -0,0 +1,82 @@
+/*
+ * 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.service.message;
+
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.common.msg.AttributeConstants;
+import org.apache.inlong.common.msg.InLongMsg;
+import org.apache.inlong.common.util.StringUtil;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+import java.nio.charset.Charset;
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+
+@Slf4j
+@Service
+public class InlongMsgDeserializeOperator implements DeserializeOperator {
+
+ @Override
+ public boolean accept(DataProxyMsgEncType type) {
+ return DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.equals(type);
+ }
+
+ @Override
+ public List<BriefMQMessage> decodeMsg(InlongStreamInfo streamInfo,
+ byte[] msgBytes, Map<String, String> headers, int index) throws
Exception {
+ String groupId = headers.get(AttributeConstants.GROUP_ID);
+ String streamId = headers.get(AttributeConstants.STREAM_ID);
+ List<BriefMQMessage> messageList = new ArrayList<>();
+ InLongMsg inLongMsg = InLongMsg.parseFrom(msgBytes);
+ for (String attr : inLongMsg.getAttrs()) {
+ Map<String, String> attributes = StringUtil.splitKv(attr,
INLONGMSG_ATTR_ENTRY_DELIMITER,
+ INLONGMSG_ATTR_KV_DELIMITER, null, null);
+ // Extracts time from the attributes
+ long msgTime;
+ if (attributes.containsKey(INLONGMSG_ATTR_TIME_T)) {
+ String date = attributes.get(INLONGMSG_ATTR_TIME_T).trim();
+ msgTime = StringUtil.parseDateTime(date);
+ } else if (attributes.containsKey(INLONGMSG_ATTR_TIME_DT)) {
+ String epoch = attributes.get(INLONGMSG_ATTR_TIME_DT).trim();
+ msgTime = Long.parseLong(epoch);
+ } else {
+ throw new
IllegalArgumentException(String.format("PARSE_ATTR_ERROR_STRING%s",
+ INLONGMSG_ATTR_TIME_T + " or " +
INLONGMSG_ATTR_TIME_DT));
+ }
+ Iterator<byte[]> iterator = inLongMsg.getIterator(attr);
+ while (iterator.hasNext()) {
+ byte[] bodyBytes = iterator.next();
+ if (Objects.isNull(bodyBytes)) {
+ continue;
+ }
+ BriefMQMessage inLongMessage =
+ new BriefMQMessage(index, groupId, streamId, msgTime,
attributes.get(NODE_IP),
+ new String(bodyBytes,
Charset.forName(streamInfo.getDataEncoding())));
+ messageList.add(inLongMessage);
+ }
+ }
+ return messageList;
+ }
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/PbMsgDeserializeOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/PbMsgDeserializeOperator.java
new file mode 100644
index 0000000000..0ae13afb42
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/PbMsgDeserializeOperator.java
@@ -0,0 +1,89 @@
+/*
+ * 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.service.message;
+
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.common.msg.AttributeConstants;
+import org.apache.inlong.common.util.Utils;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
+import org.apache.inlong.sdk.commons.protocol.ProxySdk.INLONG_COMPRESSED_TYPE;
+import org.apache.inlong.sdk.commons.protocol.ProxySdk.MapFieldEntry;
+import org.apache.inlong.sdk.commons.protocol.ProxySdk.MessageObj;
+import org.apache.inlong.sdk.commons.protocol.ProxySdk.MessageObjs;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+import java.nio.charset.Charset;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+@Slf4j
+@Service
+public class PbMsgDeserializeOperator implements DeserializeOperator {
+
+ @Override
+ public boolean accept(DataProxyMsgEncType type) {
+ return DataProxyMsgEncType.MSG_ENCODE_TYPE_PB.equals(type);
+ }
+
+ @Override
+ public List<BriefMQMessage> decodeMsg(InlongStreamInfo streamInfo,
+ byte[] msgBytes, Map<String, String> headers, int index) throws
Exception {
+ List<BriefMQMessage> messageList = new ArrayList<>();
+ int compressType =
Integer.parseInt(headers.getOrDefault(COMPRESS_TYPE_KEY, "0"));
+ byte[] values = msgBytes;
+ switch (compressType) {
+ case INLONG_COMPRESSED_TYPE.INLONG_NO_COMPRESS_VALUE:
+ break;
+ case INLONG_COMPRESSED_TYPE.INLONG_GZ_VALUE:
+ values = Utils.gzipDecompress(msgBytes, 0, msgBytes.length);
+ break;
+ case INLONG_COMPRESSED_TYPE.INLONG_SNAPPY_VALUE:
+ values = Utils.snappyDecompress(msgBytes, 0, msgBytes.length);
+ break;
+ default:
+ throw new IllegalArgumentException("Unknown compress type:" +
compressType);
+ }
+ messageList = transformMessageObjs(MessageObjs.parseFrom(values),
streamInfo, index);
+ return messageList;
+ }
+
+ private List<BriefMQMessage> transformMessageObjs(MessageObjs messageObjs,
InlongStreamInfo streamInfo, int index) {
+ if (null == messageObjs) {
+ return null;
+ }
+ List<BriefMQMessage> briefMQMessages = new ArrayList<>();
+ for (MessageObj messageObj : messageObjs.getMsgsList()) {
+ List<MapFieldEntry> mapFieldEntries = messageObj.getParamsList();
+ Map<String, String> headers = new HashMap<>();
+ for (MapFieldEntry mapFieldEntry : mapFieldEntries) {
+ headers.put(mapFieldEntry.getKey(), mapFieldEntry.getValue());
+ }
+ BriefMQMessage briefMQMessage = new BriefMQMessage(index,
headers.get(AttributeConstants.GROUP_ID),
+ headers.get(AttributeConstants.STREAM_ID),
messageObj.getMsgTime(),
+ headers.get(NODE_IP),
+ new String(messageObj.getBody().toByteArray(),
Charset.forName(streamInfo.getDataEncoding())));
+ briefMQMessages.add(briefMQMessage);
+ }
+ return briefMQMessages;
+ }
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/RawMsgDeserializeOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/RawMsgDeserializeOperator.java
new file mode 100644
index 0000000000..e039959323
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/message/RawMsgDeserializeOperator.java
@@ -0,0 +1,53 @@
+/*
+ * 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.service.message;
+
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
+import org.apache.inlong.common.msg.AttributeConstants;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+
+import java.nio.charset.Charset;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+@Slf4j
+@Service
+public class RawMsgDeserializeOperator implements DeserializeOperator {
+
+ @Override
+ public boolean accept(DataProxyMsgEncType type) {
+ return DataProxyMsgEncType.MSG_ENCODE_TYPE_RAW.equals(type);
+ }
+
+ @Override
+ public List<BriefMQMessage> decodeMsg(InlongStreamInfo streamInfo,
+ byte[] msgBytes, Map<String, String> headers, int index) throws
Exception {
+ String groupId = headers.get(AttributeConstants.GROUP_ID);
+ String streamId = headers.get(AttributeConstants.STREAM_ID);
+ long msgTime = Long.parseLong(headers.getOrDefault(MSG_TIME_KEY, "0"));
+ return Collections
+ .singletonList(new BriefMQMessage(null, groupId, streamId,
msgTime, headers.get(NODE_IP),
+ new String(msgBytes,
Charset.forName(streamInfo.getDataEncoding()))));
+ }
+
+}
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 b620fab480..b53045e9cb 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
@@ -17,12 +17,17 @@
package org.apache.inlong.manager.service.resource.queue;
+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;
+import java.util.List;
+
/**
* Interface of the message queue resource operator
*/
@@ -71,4 +76,17 @@ public interface QueueResourceOperator {
default void deleteQueueForStream(InlongGroupInfo groupInfo,
InlongStreamInfo streamInfo, String operator) {
}
+ /**
+ * Query brief mq message info
+ *
+ * @param groupInfo inlong group info
+ * @param streamInfo inlong stream info
+ * @param messageCount Count of messages to query'
+ * @return query brief mq message info
+ */
+ default List<BriefMQMessage> queryLastestMessage(InlongGroupInfo
groupInfo, InlongStreamInfo streamInfo,
+ Integer messageCount) throws PulsarClientException {
+ 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 fe597f1e84..bf6342cdb7 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
@@ -17,20 +17,26 @@
package org.apache.inlong.manager.service.resource.queue.pulsar;
+import org.apache.inlong.common.enums.DataProxyMsgEncType;
import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.conversion.ConversionHandle;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.group.pulsar.InlongPulsarInfo;
import org.apache.inlong.manager.pojo.queue.pulsar.PulsarTopicInfo;
+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 com.google.common.collect.Sets;
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.client.admin.Namespaces;
import org.apache.pulsar.client.admin.PulsarAdmin;
import org.apache.pulsar.client.admin.PulsarAdminException;
+import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.policies.data.PersistencePolicies;
import org.apache.pulsar.common.policies.data.RetentionPolicies;
@@ -40,9 +46,12 @@ import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+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.
*/
@@ -57,7 +66,9 @@ public class PulsarOperator {
private static final int MAX_PARTITION = 1000;
private static final int RETRY_TIMES = 3;
private static final int DELAY_SECONDS = 5;
-
+ private static final String PARSE_ATTR_ERROR_STRING = "Could not find %s
in attributes!";
+ @Autowired
+ public DeserializeOperatorFactory deserializeOperatorFactory;
@Autowired
private ConversionHandle conversionHandle;
@@ -374,4 +385,43 @@ public class PulsarOperator {
return false;
}
+ /**
+ * Query topic message for the given pulsar cluster.
+ */
+ public List<BriefMQMessage> queryLastestMessage(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<>();
+
+ try {
+ messages = pulsarAdmin.topics().peekMessages(topicFullName,
subName, messageCount);
+ } catch (PulsarAdminException e) {
+ String errMsg = "failed to query peek messages";
+ LOGGER.error(errMsg, e);
+ throw new BusinessException(errMsg);
+ }
+
+ int index = 0;
+ for (Message<byte[]> pulsarMessage : messages) {
+ try {
+ Map<String, String> headers = pulsarMessage.getProperties();
+ int wrapTypeId =
Integer.parseInt(headers.getOrDefault(VERSION_KEY,
+
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);
+
+ }
+ }
+
+ LOGGER.info("success query message by subs={} for topic={}", subName,
topicFullName);
+ return messageList;
+ }
+
}
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 c6b6197ae3..d3a79c9714 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
@@ -27,6 +27,7 @@ import
org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.cluster.ClusterInfo;
import org.apache.inlong.manager.pojo.cluster.pulsar.PulsarClusterInfo;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.group.pulsar.InlongPulsarInfo;
import org.apache.inlong.manager.pojo.queue.pulsar.PulsarTopicInfo;
@@ -44,9 +45,11 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.client.admin.PulsarAdmin;
+import org.apache.pulsar.client.api.PulsarClientException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import java.util.ArrayList;
import java.util.List;
/**
@@ -61,6 +64,8 @@ public class PulsarResourceOperator implements
QueueResourceOperator {
*/
public static final String PULSAR_SUBSCRIPTION = "%s_%s_%s_consumer_group";
+ public static final String PULSAR_SUBSCRIPTION_REALTIME_REVIEW =
"%s_%s_consumer_group_realtime_review";
+
@Autowired
private InlongClusterService clusterService;
@Autowired
@@ -299,4 +304,37 @@ public class PulsarResourceOperator implements
QueueResourceOperator {
}
}
+ /**
+ * Query lastest message from pulsar
+ */
+ public List<BriefMQMessage> queryLastestMessage(InlongGroupInfo groupInfo,
+ InlongStreamInfo streamInfo, Integer messageCount) throws
PulsarClientException {
+ String groupId = streamInfo.getInlongGroupId();
+ InlongPulsarInfo inlongPulsarInfo = ((InlongPulsarInfo) groupInfo);
+ PulsarClusterInfo pulsarCluster = (PulsarClusterInfo)
clusterService.getOne(groupInfo.getInlongClusterTag(),
+ null, ClusterType.PULSAR);
+ List<BriefMQMessage> briefMQMessages = new ArrayList<>();
+
+ try (PulsarAdmin pulsarAdmin =
PulsarUtils.getPulsarAdmin(pulsarCluster)) {
+ String tenant = inlongPulsarInfo.getPulsarTenant();
+ if (StringUtils.isBlank(tenant)) {
+ tenant = pulsarCluster.getPulsarTenant();
+ }
+
+ String namespace = groupInfo.getMqResource();
+ String topicName = streamInfo.getMqResource();
+ String fullTopicName = tenant + "/" + namespace + "/" + topicName;
+ String clusterTag = inlongPulsarInfo.getInlongClusterTag();
+ String subs = String.format(PULSAR_SUBSCRIPTION_REALTIME_REVIEW,
clusterTag, topicName);
+ briefMQMessages =
+ pulsarOperator.queryLastestMessage(pulsarAdmin,
fullTopicName, subs, messageCount, streamInfo);
+
+ // insert the consumer group info into the inlong_consume table
+ Integer id = consumeService.saveBySystem(groupInfo, topicName,
subs);
+ log.info("success to save inlong consume [{}] for subs={},
groupId={}, topic={}",
+ id, subs, groupId, topicName);
+ }
+ return briefMQMessages;
+ }
+
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamService.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamService.java
index 9f83d7764a..f0f53cd409 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamService.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamService.java
@@ -18,6 +18,7 @@
package org.apache.inlong.manager.service.stream;
import org.apache.inlong.manager.pojo.common.PageResult;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.sink.ParseFieldRequest;
import org.apache.inlong.manager.pojo.stream.InlongStreamApproveRequest;
import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
@@ -256,4 +257,16 @@ public interface InlongStreamService {
* Converts an Excel file to a streamFields
*/
List<StreamField> parseFields(MultipartFile file);
+
+ /**
+ * List brief mq message info
+ *
+ * @param groupId inlong group id
+ * @param streamId inlong stream id
+ * @param messageCount Count of messages to query'
+ * @param operator operator
+ * @return list of brief mq message info
+ */
+ List<BriefMQMessage> listMessages(String groupId, String streamId, Integer
messageCount, String operator);
+
}
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 3963601a5d..0ae4f16b25 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
@@ -36,6 +36,8 @@ import
org.apache.inlong.manager.dao.mapper.InlongStreamFieldEntityMapper;
import org.apache.inlong.manager.pojo.common.OrderFieldEnum;
import org.apache.inlong.manager.pojo.common.OrderTypeEnum;
import org.apache.inlong.manager.pojo.common.PageResult;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
+import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.sink.ParseFieldRequest;
import org.apache.inlong.manager.pojo.sink.SinkBriefInfo;
import org.apache.inlong.manager.pojo.sink.StreamSink;
@@ -50,6 +52,10 @@ import
org.apache.inlong.manager.pojo.stream.InlongStreamRequest;
import org.apache.inlong.manager.pojo.stream.StreamField;
import org.apache.inlong.manager.pojo.user.UserInfo;
import org.apache.inlong.manager.pojo.user.UserRoleCode;
+import org.apache.inlong.manager.service.group.InlongGroupOperator;
+import org.apache.inlong.manager.service.group.InlongGroupOperatorFactory;
+import org.apache.inlong.manager.service.resource.queue.QueueResourceOperator;
+import
org.apache.inlong.manager.service.resource.queue.QueueResourceOperatorFactory;
import org.apache.inlong.manager.service.sink.StreamSinkService;
import org.apache.inlong.manager.service.source.StreamSourceService;
import org.apache.inlong.manager.service.user.UserService;
@@ -70,6 +76,7 @@ import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
@@ -124,6 +131,12 @@ public class InlongStreamServiceImpl implements
InlongStreamService {
private ObjectMapper objectMapper;
@Autowired
private UserService userService;
+ @Autowired
+ @Lazy
+ private QueueResourceOperatorFactory queueOperatorFactory;
+ @Autowired
+ @Lazy
+ private InlongGroupOperatorFactory groupOperatorFactory;
@Transactional(rollbackFor = Throwable.class)
@Override
@@ -992,4 +1005,21 @@ public class InlongStreamServiceImpl implements
InlongStreamService {
return entity;
}
+ @Override
+ public List<BriefMQMessage> listMessages(String groupId, String streamId,
Integer messageCount, String operator) {
+ InlongGroupEntity groupEntity = groupMapper.selectByGroupId(groupId);
+ // check user
+ userService.checkUser(groupEntity.getInCharges(), operator,
ErrorCodeEnum.GROUP_PERMISSION_DENIED.getMessage());
+ InlongGroupOperator instance =
groupOperatorFactory.getInstance(groupEntity.getMqType());
+ InlongGroupInfo groupInfo = instance.getFromEntity(groupEntity);
+ InlongStreamInfo inlongStreamInfo = get(groupId, streamId);
+ List<BriefMQMessage> messageList = new ArrayList<>();
+ QueueResourceOperator queueOperator =
queueOperatorFactory.getInstance(groupEntity.getMqType());
+ try {
+ messageList = queueOperator.queryLastestMessage(groupInfo,
inlongStreamInfo, messageCount);
+ } catch (Exception e) {
+ LOGGER.error("query message error ", e);
+ }
+ return messageList;
+ }
}
diff --git
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
index ab6baff32a..64bbeb4671 100644
---
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
+++
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.tool.excel.ExcelTool;
import org.apache.inlong.manager.common.validation.UpdateValidation;
import org.apache.inlong.manager.pojo.common.PageResult;
import org.apache.inlong.manager.pojo.common.Response;
+import org.apache.inlong.manager.pojo.consume.BriefMQMessage;
import org.apache.inlong.manager.pojo.sink.ParseFieldRequest;
import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
@@ -234,4 +235,17 @@ public class InlongStreamController {
}
}
+ @RequestMapping(value = "/stream/listMessages", method = RequestMethod.GET)
+ @ApiOperation(value = "Get inlong stream message")
+ @ApiImplicitParams({
+ @ApiImplicitParam(name = "groupId", dataTypeClass = String.class,
required = true),
+ @ApiImplicitParam(name = "streamId", dataTypeClass = String.class,
required = true),
+ @ApiImplicitParam(name = "messageCount", dataTypeClass =
String.class, required = true)
+ })
+ public Response<List<BriefMQMessage>> listMessages(@RequestParam String
groupId, @RequestParam String streamId,
+ @RequestParam Integer messageCount) {
+ String username = LoginUserUtils.getLoginUser().getName();
+ return Response.success(streamService.listMessages(groupId, streamId,
messageCount, username));
+ }
+
}
diff --git
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializer.java
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializer.java
index 99894456ca..7af7c06363 100644
---
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializer.java
+++
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializer.java
@@ -19,6 +19,8 @@ package org.apache.inlong.sdk.sort.impl.decode;
import org.apache.inlong.common.enums.DataProxyMsgEncType;
import org.apache.inlong.common.msg.InLongMsg;
+import org.apache.inlong.common.util.StringUtil;
+import org.apache.inlong.common.util.Utils;
import org.apache.inlong.sdk.commons.protocol.EventConstants;
import org.apache.inlong.sdk.commons.protocol.ProxySdk.MapFieldEntry;
import org.apache.inlong.sdk.commons.protocol.ProxySdk.MessageObj;
@@ -27,8 +29,6 @@ import org.apache.inlong.sdk.sort.api.ClientContext;
import org.apache.inlong.sdk.sort.api.Deserializer;
import org.apache.inlong.sdk.sort.entity.InLongMessage;
import org.apache.inlong.sdk.sort.entity.InLongTopic;
-import org.apache.inlong.sdk.sort.util.StringUtil;
-import org.apache.inlong.sdk.sort.util.Utils;
import java.io.IOException;
import java.util.ArrayList;
diff --git
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongMultiTopicManager.java
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongMultiTopicManager.java
index cb5cb474fc..de30a3f308 100644
---
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongMultiTopicManager.java
+++
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongMultiTopicManager.java
@@ -17,6 +17,7 @@
package org.apache.inlong.sdk.sort.manager;
+import org.apache.inlong.common.util.StringUtil;
import org.apache.inlong.sdk.sort.api.ClientContext;
import org.apache.inlong.sdk.sort.api.InlongTopicTypeEnum;
import org.apache.inlong.sdk.sort.api.QueryConsumeConfig;
@@ -27,7 +28,6 @@ import org.apache.inlong.sdk.sort.entity.ConsumeConfig;
import org.apache.inlong.sdk.sort.entity.InLongTopic;
import org.apache.inlong.sdk.sort.fetcher.tube.TubeConsumerCreator;
import org.apache.inlong.sdk.sort.util.PeriodicTask;
-import org.apache.inlong.sdk.sort.util.StringUtil;
import org.apache.inlong.tubemq.client.config.TubeClientConfig;
import org.apache.inlong.tubemq.client.exception.TubeClientException;
import org.apache.inlong.tubemq.client.factory.MessageSessionFactory;
diff --git
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongSingleTopicManager.java
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongSingleTopicManager.java
index b1af5a3126..a59e81d80d 100644
---
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongSingleTopicManager.java
+++
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/manager/InlongSingleTopicManager.java
@@ -17,6 +17,7 @@
package org.apache.inlong.sdk.sort.manager;
+import org.apache.inlong.common.util.StringUtil;
import org.apache.inlong.sdk.sort.api.ClientContext;
import org.apache.inlong.sdk.sort.api.InlongTopicTypeEnum;
import org.apache.inlong.sdk.sort.api.QueryConsumeConfig;
@@ -27,7 +28,6 @@ import org.apache.inlong.sdk.sort.entity.ConsumeConfig;
import org.apache.inlong.sdk.sort.entity.InLongTopic;
import org.apache.inlong.sdk.sort.fetcher.tube.TubeConsumerCreator;
import org.apache.inlong.sdk.sort.util.PeriodicTask;
-import org.apache.inlong.sdk.sort.util.StringUtil;
import org.apache.inlong.tubemq.client.config.TubeClientConfig;
import org.apache.inlong.tubemq.client.factory.MessageSessionFactory;
import org.apache.inlong.tubemq.client.factory.TubeSingleSessionFactory;
diff --git
a/inlong-sdk/sort-sdk/src/test/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializerTest.java
b/inlong-sdk/sort-sdk/src/test/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializerTest.java
index 1f0cd88621..88d8e04b79 100644
---
a/inlong-sdk/sort-sdk/src/test/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializerTest.java
+++
b/inlong-sdk/sort-sdk/src/test/java/org/apache/inlong/sdk/sort/impl/decode/MessageDeserializerTest.java
@@ -18,6 +18,7 @@
package org.apache.inlong.sdk.sort.impl.decode;
import org.apache.inlong.common.msg.InLongMsg;
+import org.apache.inlong.common.util.Utils;
import org.apache.inlong.sdk.commons.protocol.ProxySdk.MapFieldEntry;
import org.apache.inlong.sdk.commons.protocol.ProxySdk.MessageObj;
import org.apache.inlong.sdk.commons.protocol.ProxySdk.MessageObjs;
@@ -26,7 +27,6 @@ import org.apache.inlong.sdk.sort.entity.CacheZoneCluster;
import org.apache.inlong.sdk.sort.entity.InLongMessage;
import org.apache.inlong.sdk.sort.entity.InLongTopic;
import org.apache.inlong.sdk.sort.impl.ClientContextImpl;
-import org.apache.inlong.sdk.sort.util.Utils;
import com.google.protobuf.ByteString;
import org.junit.Assert;