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;


Reply via email to