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 e673c2565 [INLONG-6104][Manager] Support getting backup info in 
getSortSource (#6275)
e673c2565 is described below

commit e673c2565bf24e2550444942bab5e165bea60d81
Author: vernedeng <[email protected]>
AuthorDate: Sat Oct 29 19:51:13 2022 +0800

    [INLONG-6104][Manager] Support getting backup info in getSortSource (#6275)
---
 .../dao/mapper/InlongClusterEntityMapper.java      |   6 +-
 .../dao/mapper/InlongGroupEntityMapper.java        |   6 +-
 .../dao/mapper/InlongGroupExtEntityMapper.java     |  10 +-
 .../dao/mapper/InlongStreamEntityMapper.java       |   8 +-
 .../dao/mapper/InlongStreamExtEntityMapper.java    |   6 +
 .../manager/dao/mapper/StreamSinkEntityMapper.java |  10 +-
 .../mappers/InlongClusterEntityMapper.xml          |   3 +-
 .../resources/mappers/InlongGroupEntityMapper.xml  |   2 +-
 .../mappers/InlongGroupExtEntityMapper.xml         |   7 +
 .../resources/mappers/InlongStreamEntityMapper.xml |   8 +
 .../mappers/InlongStreamExtEntityMapper.xml        |  20 +-
 .../resources/mappers/StreamSinkEntityMapper.xml   |   8 +-
 .../sort/standalone/SortSourceClusterInfo.java     |   5 +-
 .../pojo/sort/standalone/SortSourceGroupInfo.java  |   6 +-
 .../pojo/sort/standalone/SortSourceStreamInfo.java |  34 +--
 ...reamInfo.java => SortSourceStreamSinkInfo.java} |   8 +-
 .../manager/service/core/SortConfigLoader.java     |  70 ++++++
 .../service/core/impl/SortConfigLoaderImpl.java    | 109 ++++++++
 .../service/core/impl/SortSourceServiceImpl.java   | 280 ++++++++++-----------
 .../service/group/InlongGroupOperator4Pulsar.java  |  17 +-
 .../service/group/InlongGroupServiceImpl.java      |   2 +-
 .../manager/service/sort/SortServiceImplTest.java  | 113 ++++++---
 22 files changed, 487 insertions(+), 251 deletions(-)

diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
index db7ae4a70..1f993527e 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
@@ -17,7 +17,10 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.InlongClusterEntity;
 import org.apache.inlong.manager.pojo.cluster.ClusterPageRequest;
 import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
@@ -49,7 +52,8 @@ public interface InlongClusterEntityMapper {
      *
      * @return All cluster info.
      */
-    List<SortSourceClusterInfo> selectAllClusters();
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<SortSourceClusterInfo> selectAllClusters();
 
     List<InlongClusterEntity> selectByClusterTag(String clusterTag);
 
diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
index b9ba70058..7dc92543d 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
@@ -17,7 +17,10 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.InlongGroupEntity;
 import org.apache.inlong.manager.pojo.group.InlongGroupBriefInfo;
 import org.apache.inlong.manager.pojo.group.InlongGroupPageRequest;
@@ -52,7 +55,8 @@ public interface InlongGroupEntityMapper {
      *
      * @return All inlong group info.
      */
-    List<SortSourceGroupInfo> selectAllGroups();
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<SortSourceGroupInfo> selectAllGroups();
 
     /**
      * Select all groups which are logical deleted before the specified last 
modify time
diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
index 4e7c1d4ab..38b694f1c 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
@@ -17,7 +17,10 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
 import org.springframework.stereotype.Repository;
 
@@ -38,7 +41,12 @@ public interface InlongGroupExtEntityMapper {
 
     int updateByPrimaryKey(InlongGroupExtEntity record);
 
-    InlongGroupExtEntity selectByUniqueKey(@Param("groupId") String groupId, 
@Param("keyName") String keyName);
+    InlongGroupExtEntity selectByUniqueKey(
+            @Param("inlongGroupId") String inlongGroupId,
+            @Param("keyName") String keyName);
+
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<InlongGroupExtEntity> selectByKeyName(@Param("keyName") String 
keyName);
 
     /**
      * Insert data in batches
diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
index 66f101a27..b676dcfeb 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
@@ -17,8 +17,12 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.InlongStreamEntity;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
 import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
 import org.apache.inlong.manager.pojo.stream.InlongStreamPageRequest;
 import org.springframework.stereotype.Repository;
@@ -50,6 +54,9 @@ public interface InlongStreamEntityMapper {
 
     int selectCountByGroupId(@Param("groupId") String groupId);
 
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<SortSourceStreamInfo> selectAllStreams();
+
     int updateByPrimaryKey(InlongStreamEntity record);
 
     int updateByIdentifierSelective(InlongStreamEntity streamEntity);
@@ -70,5 +77,4 @@ public interface InlongStreamEntityMapper {
      * @return rows deleted
      */
     int deleteByInlongGroupIds(@Param("groupIdList") List<String> groupIdList);
-
 }
diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
index 2c7316bf6..b9457210b 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
@@ -17,7 +17,10 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
 import org.springframework.stereotype.Repository;
 
@@ -47,6 +50,9 @@ public interface InlongStreamExtEntityMapper {
     InlongStreamExtEntity selectByKey(@Param("groupId") String groupId, 
@Param("streamId") String streamId,
             @Param("keyName") String keyName);
 
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<InlongStreamExtEntity> selectByKeyName(@Param("keyName") String 
keyName);
+
     int updateByPrimaryKey(InlongStreamExtEntity record);
 
     /**
diff --git 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
index 7142976dc..20fbd9e55 100644
--- 
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
+++ 
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
@@ -17,13 +17,16 @@
 
 package org.apache.inlong.manager.dao.mapper;
 
+import org.apache.ibatis.annotations.Options;
 import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
 import org.apache.inlong.manager.pojo.sink.SinkBriefInfo;
 import org.apache.inlong.manager.pojo.sink.SinkInfo;
 import org.apache.inlong.manager.pojo.sink.SinkPageRequest;
 import org.apache.inlong.manager.pojo.sort.standalone.SortIdInfo;
-import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
 import org.apache.inlong.manager.pojo.sort.standalone.SortTaskInfo;
 import org.springframework.stereotype.Repository;
 
@@ -117,6 +120,8 @@ public interface StreamSinkEntityMapper {
      */
     List<SinkInfo> selectAllConfig(@Param("groupId") String groupId, 
@Param("idList") List<String> streamIdList);
 
+    List<StreamSinkEntity> selectAllStreamSinks();
+
     /**
      * Select all tasks for sort-standalone
      *
@@ -136,7 +141,8 @@ public interface StreamSinkEntityMapper {
      *
      * @return All stream info
      */
-    List<SortSourceStreamInfo> selectAllStreams();
+    @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 
Integer.MIN_VALUE)
+    Cursor<SortSourceStreamSinkInfo> selectAllStreams();
 
     int updateByIdSelective(StreamSinkEntity record);
 
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
index d9d7fef62..931ca5f83 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
@@ -164,9 +164,10 @@
         </where>
         order by modify_time desc
     </select>
-    <select id="selectAllClusters" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo">
+    <select id="selectAllClusters" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo"
 >
         select name,
                type,
+               url,
                cluster_tags,
                ext_tag,
                ext_params
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
index e3b832339..47a97759f 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
@@ -202,7 +202,7 @@
     <select id="selectAllGroups" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo">
         select inlong_group_id    as groupId,
                inlong_cluster_tag as clusterTag,
-               mq_resource        as topic,
+               mq_resource        as mqResource,
                ext_params         as extParams,
                mq_type            as mqType
         from inlong_group
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
index dba70b312..9c16c25d7 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
@@ -53,6 +53,13 @@
         and key_name = #{keyName, jdbcType=VARCHAR}
         and is_deleted = 0
     </select>
+    <select id="selectByKeyName" 
resultType="org.apache.inlong.manager.dao.entity.InlongGroupExtEntity">
+        select
+        <include refid="Base_Column_List"/>
+        from inlong_group_ext
+        where key_name = #{keyName, jdbcType=VARCHAR}
+        and is_deleted = 0
+    </select>
 
     <delete id="deleteByPrimaryKey" parameterType="java.lang.Integer">
         delete
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
index 136c33f12..231769ccc 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
@@ -285,6 +285,14 @@
         where inlong_group_id = #{groupId, jdbcType=VARCHAR}
         and is_deleted = 0
     </select>
+    <select id="selectAllStreams" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo">
+        select
+               inlong_group_id,
+               inlong_stream_id,
+               mq_resource
+        from inlong_stream
+        where is_deleted = 0
+    </select>
     <select id="selectCountByGroupId" resultType="java.lang.Integer">
         select count(1)
         from inlong_stream
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
index 25e4e1a7a..03c785f58 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
@@ -35,30 +35,29 @@
 
     <insert id="insert" 
parameterType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
         insert into inlong_stream_ext (id, inlong_group_id, inlong_stream_id,
-                                       key_name, is_deleted, modify_time,
+                                       key_name, modify_time,
                                        key_value)
         values (#{id,jdbcType=INTEGER}, #{inlongGroupId,jdbcType=VARCHAR}, 
#{inlongStreamId,jdbcType=VARCHAR},
-                #{keyName,jdbcType=VARCHAR}, #{isDeleted,jdbcType=INTEGER}, 
#{modifyTime,jdbcType=TIMESTAMP},
+                #{keyName,jdbcType=VARCHAR}, #{modifyTime,jdbcType=TIMESTAMP},
                 #{keyValue,jdbcType=LONGVARCHAR})
     </insert>
     <!-- Bulk insert-->
     <insert id="insertAll" parameterType="java.util.List">
         insert into inlong_stream_ext
-        (id, inlong_group_id, inlong_stream_id, key_name, key_value, 
is_deleted)
+        (id, inlong_group_id, inlong_stream_id, key_name, key_value)
         values
         <foreach collection="extList" separator="," index="index" item="item">
-            (#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId}, 
#{item.keyName}, #{item.keyValue},
-            #{item.isDeleted})
+            (#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId}, 
#{item.keyName}, #{item.keyValue})
         </foreach>
     </insert>
     <!-- Bulk insert, update if exists-->
     <insert id="insertOnDuplicateKeyUpdate" parameterType="java.util.List">
         insert into inlong_stream_ext
-        (id, inlong_group_id, inlong_stream_id, key_name, key_value, 
is_deleted)
+        (id, inlong_group_id, inlong_stream_id, key_name, key_value)
         values
         <foreach collection="extList" separator="," index="index" item="item">
             (#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId},
-            #{item.keyName}, #{item.keyValue}, #{item.isDeleted})
+            #{item.keyName}, #{item.keyValue})
         </foreach>
         ON DUPLICATE KEY UPDATE
         inlong_group_id = VALUES(inlong_group_id),
@@ -87,6 +86,13 @@
         and key_name = #{keyName, jdbcType=VARCHAR}
         and is_deleted = 0
     </select>
+    <select id="selectByKeyName" 
resultType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
+        select
+        <include refid="Base_Column_List"/>
+        from inlong_stream_ext
+        where key_name = #{keyName, jdbcType=VARCHAR}
+        and is_deleted = 0
+    </select>
 
     <update id="updateByPrimaryKey" 
parameterType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
         update inlong_stream_ext
diff --git 
a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
 
b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
index c4af01696..aac142d59 100644
--- 
a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
+++ 
b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
@@ -350,6 +350,12 @@
             </if>
         </where>
     </select>
+    <select id="selectAllStreamSinks" 
resultType="org.apache.inlong.manager.dao.entity.StreamSinkEntity">
+        select
+        <include refid="Base_Column_List"/>
+        from stream_sink
+        where is_deleted = 0
+    </select>
     <select id="selectAllTasks" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortTaskInfo">
         select inlong_cluster_name as sortClusterName,
                sort_task_name,
@@ -368,7 +374,7 @@
         from stream_sink
         where is_deleted = 0
     </select>
-    <select id="selectAllStreams" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo">
+    <select id="selectAllStreams" 
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo">
         select inlong_cluster_name as sortClusterName,
                sort_task_name,
                inlong_group_id     as groupId,
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
index 148320178..933982729 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
@@ -31,6 +31,7 @@ import java.util.concurrent.ConcurrentHashMap;
 public class SortSourceClusterInfo {
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(SortSourceClusterInfo.class);
+    private static final Gson GSON = new Gson();
     private static final Splitter.MapSplitter MAP_SPLITTER = 
Splitter.on("&").trimResults()
             .withKeyValueSeparator("=");
     private static final String KEY_IS_CONSUMABLE = "consumer";
@@ -38,6 +39,7 @@ public class SortSourceClusterInfo {
     private static final long serialVersionUID = 1L;
     String name;
     String type;
+    String url;
     String clusterTags;
     String extTag;
     String extParams;
@@ -47,8 +49,7 @@ public class SortSourceClusterInfo {
     public Map<String, String> getExtParamsMap() {
         if (extParamsMap.isEmpty() && extParams != null) {
             try {
-                Gson gson = new Gson();
-                extParamsMap = gson.fromJson(extParams, Map.class);
+                extParamsMap = GSON.fromJson(extParams, Map.class);
             } catch (Throwable t) {
                 LOGGER.error("fail to parse cluster ext params", t);
             }
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
index 8d887d609..b7b951652 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
@@ -30,13 +30,14 @@ import java.util.concurrent.ConcurrentHashMap;
 public class SortSourceGroupInfo {
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(SortSourceGroupInfo.class);
+    private static final Gson GSON = new Gson();
     private static final String KEY_BACKUP_CLUSTER_TAG = "backup_cluster_tag";
     private static final String KEY_BACKUP_TOPIC = "backup_topic";
 
     private static final long serialVersionUID = 1L;
     String groupId;
     String clusterTag;
-    String topic;
+    String mqResource;
     String extParams;
     String mqType;
     Map<String, String> extParamsMap = new ConcurrentHashMap<>();
@@ -44,8 +45,7 @@ public class SortSourceGroupInfo {
     public Map<String, String> getExtParamsMap() {
         if (extParamsMap.isEmpty() && StringUtils.isNotBlank(extParams)) {
             try {
-                Gson gson = new Gson();
-                extParamsMap = gson.fromJson(extParams, Map.class);
+                extParamsMap = GSON.fromJson(extParams, Map.class);
             } catch (Throwable t) {
                 LOGGER.error("fail to parse group ext params", t);
             }
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
index 6462e3e77..bd89a573c 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
@@ -17,39 +17,11 @@
 
 package org.apache.inlong.manager.pojo.sort.standalone;
 
-import com.google.gson.Gson;
 import lombok.Data;
-import org.apache.commons.lang3.StringUtils;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
 
 @Data
 public class SortSourceStreamInfo {
-
-    private static final long serialVersionUID = 1L;
-    private static final Logger LOGGER = 
LoggerFactory.getLogger(SortSourceStreamInfo.class);
-    String sortClusterName;
-    String sortTaskName;
-    String groupId;
-    String extParams;
-    Map<String, String> extParamsMap;
-
-    public Map<String, String> getExtParamsMap() {
-        if (extParamsMap != null) {
-            return extParamsMap;
-        }
-        if (StringUtils.isNotBlank(extParams)) {
-            try {
-                Gson gson = new Gson();
-                extParamsMap = gson.fromJson(extParams, Map.class);
-            } catch (Throwable t) {
-                LOGGER.error("fail to parse source stream ext params", t);
-                extParamsMap = new ConcurrentHashMap<>();
-            }
-        }
-        return extParamsMap;
-    }
+    private String inlongGroupId;
+    private String inlongStreamId;
+    private String mqResource;
 }
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
similarity index 90%
copy from 
inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
copy to 
inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
index 6462e3e77..4aa2e1168 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
@@ -27,10 +27,11 @@ import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
 
 @Data
-public class SortSourceStreamInfo {
+public class SortSourceStreamSinkInfo {
 
     private static final long serialVersionUID = 1L;
-    private static final Logger LOGGER = 
LoggerFactory.getLogger(SortSourceStreamInfo.class);
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(SortSourceStreamSinkInfo.class);
+    private static final Gson GSON = new Gson();
     String sortClusterName;
     String sortTaskName;
     String groupId;
@@ -43,8 +44,7 @@ public class SortSourceStreamInfo {
         }
         if (StringUtils.isNotBlank(extParams)) {
             try {
-                Gson gson = new Gson();
-                extParamsMap = gson.fromJson(extParams, Map.class);
+                extParamsMap = GSON.fromJson(extParams, Map.class);
             } catch (Throwable t) {
                 LOGGER.error("fail to parse source stream ext params", t);
                 extParamsMap = new ConcurrentHashMap<>();
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
new file mode 100644
index 000000000..7bfa2a8ec
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
@@ -0,0 +1,70 @@
+/*
+ * 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.core;
+
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+
+import java.util.List;
+
+/**
+ * Loader for sort service to load configs thought Cursor
+ */
+public interface SortConfigLoader {
+    /**
+     * Load all clusters by cursor
+     * @return List of clusters, including MQ cluster and DataNode cluster.
+     */
+    List<SortSourceClusterInfo> loadAllClusters();
+
+    /**
+     * Load stream sinks by cursor
+     * @return List of Stream sinks.
+     */
+    List<SortSourceStreamSinkInfo> loadAllStreamSinks();
+
+    /**
+     * Load groups by cursor
+     * @return List of group info
+     */
+    List<SortSourceGroupInfo> loadAllGroup();
+
+    /**
+     * Load group backup info by cursor
+     * @param keyName Key name
+     * @return List of group backup info
+     */
+    List<InlongGroupExtEntity> loadGroupBackupInfo(String keyName);
+
+    /**
+     * Load stream backup info by cursor
+     * @param keyName Key name
+     * @return List of stream backup info
+     */
+    List<InlongStreamExtEntity> loadStreamBackupInfo(String keyName);
+
+    /**
+     * Load all inlong stream info by cursor
+     * @return List of stream info
+     */
+    List<SortSourceStreamInfo> loadAllStreams();
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
new file mode 100644
index 000000000..f442f9a46
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
@@ -0,0 +1,109 @@
+/*
+ * 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.core.impl;
+
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
+import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongGroupExtEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongStreamEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongStreamExtEntityMapper;
+import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+import org.apache.inlong.manager.service.core.SortConfigLoader;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.ArrayList;
+import java.util.List;
+
+@Service
+public class SortConfigLoaderImpl implements SortConfigLoader {
+    @Autowired
+    private InlongClusterEntityMapper clusterEntityMapper;
+    @Autowired
+    private StreamSinkEntityMapper streamSinkEntityMapper;
+    @Autowired
+    private InlongGroupEntityMapper inlongGroupEntityMapper;
+    @Autowired
+    private InlongGroupExtEntityMapper inlongGroupExtEntityMapper;
+    @Autowired
+    private InlongStreamExtEntityMapper inlongStreamExtEntityMapper;
+    @Autowired
+    private InlongStreamEntityMapper inlongStreamEntityMapper;
+
+    @Transactional
+    @Override
+    public List<SortSourceClusterInfo> loadAllClusters() {
+        Cursor<SortSourceClusterInfo> cursor = 
clusterEntityMapper.selectAllClusters();
+        List<SortSourceClusterInfo> allClusters = new ArrayList<>();
+        cursor.forEach(allClusters::add);
+        return allClusters;
+    }
+
+    @Transactional
+    @Override
+    public List<SortSourceStreamSinkInfo> loadAllStreamSinks() {
+        Cursor<SortSourceStreamSinkInfo> cursor = 
streamSinkEntityMapper.selectAllStreams();
+        List<SortSourceStreamSinkInfo> allStreamSinks = new ArrayList<>();
+        cursor.forEach(allStreamSinks::add);
+        return allStreamSinks;
+    }
+
+    @Transactional
+    @Override
+    public List<SortSourceGroupInfo> loadAllGroup() {
+        Cursor<SortSourceGroupInfo> cursor = 
inlongGroupEntityMapper.selectAllGroups();
+        List<SortSourceGroupInfo> allGroups = new ArrayList<>();
+        cursor.forEach(allGroups::add);
+        return allGroups;
+    }
+
+    @Transactional
+    @Override
+    public List<InlongGroupExtEntity> loadGroupBackupInfo(String keyName) {
+        Cursor<InlongGroupExtEntity> cursor = 
inlongGroupExtEntityMapper.selectByKeyName(keyName);
+        List<InlongGroupExtEntity> groupBackupInfo = new ArrayList<>();
+        cursor.forEach(groupBackupInfo::add);
+        return groupBackupInfo;
+    }
+
+    @Transactional
+    @Override
+    public List<InlongStreamExtEntity> loadStreamBackupInfo(String keyName) {
+        Cursor<InlongStreamExtEntity> cursor = 
inlongStreamExtEntityMapper.selectByKeyName(keyName);
+        List<InlongStreamExtEntity> streamBackupInfo = new ArrayList<>();
+        cursor.forEach(streamBackupInfo::add);
+        return streamBackupInfo;
+    }
+
+    @Transactional
+    @Override
+    public List<SortSourceStreamInfo> loadAllStreams() {
+        Cursor<SortSourceStreamInfo> cursor = 
inlongStreamEntityMapper.selectAllStreams();
+        List<SortSourceStreamInfo> allStreams = new ArrayList<>();
+        cursor.forEach(allStreams::add);
+        return allStreams;
+    }
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
index 611e8d02e..cadf713c5 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
@@ -20,17 +20,21 @@ package org.apache.inlong.manager.service.core.impl;
 import com.google.gson.Gson;
 import org.apache.commons.codec.digest.DigestUtils;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.inlong.common.constant.ClusterSwitch;
 import org.apache.inlong.common.pojo.sdk.CacheZone;
 import org.apache.inlong.common.pojo.sdk.CacheZoneConfig;
 import org.apache.inlong.common.pojo.sdk.SortSourceConfigResponse;
 import org.apache.inlong.common.pojo.sdk.Topic;
 import org.apache.inlong.manager.common.consts.MQType;
-import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
-import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
-import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.common.enums.ClusterType;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
 import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
 import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
 import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+import org.apache.inlong.manager.service.core.SortConfigLoader;
 import org.apache.inlong.manager.service.core.SortSourceService;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -45,9 +49,7 @@ import java.util.Collection;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
-import java.util.Locale;
 import java.util.Map;
-import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
@@ -74,10 +76,8 @@ public class SortSourceServiceImpl implements 
SortSourceService {
             add(MQType.PULSAR);
         }
     };
-    private static final String KEY_SERVICE_URL = "serviceUrl";
     private static final String KEY_AUTH = "authentication";
     private static final String KEY_TENANT = "tenant";
-    private static final String KEY_NAME_SPACE = "namespace";
 
     private static final int RESPONSE_CODE_SUCCESS = 0;
     private static final int RESPONSE_CODE_NO_UPDATE = 1;
@@ -93,12 +93,16 @@ public class SortSourceServiceImpl implements 
SortSourceService {
      */
     private Map<String, Map<String, CacheZoneConfig>> sortSourceConfigMap = 
new ConcurrentHashMap<>();
 
+    private Map<String, List<SortSourceClusterInfo>> mqClusters;
+    private Map<String, SortSourceGroupInfo> groupInfos;
+    private Map<String, SortSourceStreamInfo> allStreams;
+    private Map<String, String> backupClusterTag;
+    private Map<String, String> backupGroupMqResource;
+    private Map<String, String> backupStreamMqResource;
+    private Map<String, Map<String, List<String>>> groupMap;
+
     @Autowired
-    private InlongClusterEntityMapper clusterEntityMapper;
-    @Autowired
-    private StreamSinkEntityMapper streamSinkEntityMapper;
-    @Autowired
-    private InlongGroupEntityMapper inlongGroupEntityMapper;
+    private SortConfigLoader configLoader;
 
     @PostConstruct
     public void initialize() {
@@ -115,7 +119,8 @@ public class SortSourceServiceImpl implements 
SortSourceService {
     public void reload() {
         LOGGER.debug("start to reload sort config.");
         try {
-            reloadAllSourceConfig();
+            reloadAllConfigs();
+            parseAll();
         } catch (Throwable t) {
             LOGGER.error("fail to reload all source config", t);
         }
@@ -177,16 +182,21 @@ public class SortSourceServiceImpl implements 
SortSourceService {
 
     }
 
-    private void reloadAllSourceConfig() {
+    private void reloadAllConfigs() {
 
-        // get all streams.
-        List<SortSourceStreamInfo> allStreamInfos = 
streamSinkEntityMapper.selectAllStreams().stream()
-                .filter(dto -> dto.getSortClusterName() != null && 
dto.getSortTaskName() != null)
-                .collect(Collectors.toList());
+        // reload mq cluster and sort cluster
+        List<SortSourceClusterInfo> allClusters = 
configLoader.loadAllClusters();
 
-        // convert to Map<clusterName, Map<taskName, List<groupId>>> format.
-        Map<String, Map<String, List<String>>> groupMap = new 
ConcurrentHashMap<>();
-        allStreamInfos.forEach(stream -> {
+        // group mq clusters by cluster tag
+        mqClusters = allClusters.stream()
+                .filter(cluster -> 
SUPPORTED_MQ_TYPE.contains(cluster.getType()))
+                .filter(SortSourceClusterInfo::isConsumable)
+                
.collect(Collectors.groupingBy(SortSourceClusterInfo::getClusterTags));
+
+        // reload all stream sinks, to Map<clusterName, Map<taskName, 
List<groupId>>> format
+        List<SortSourceStreamSinkInfo> allStreamSinks = 
configLoader.loadAllStreamSinks();
+        groupMap = new HashMap<>();
+        allStreamSinks.forEach(stream -> {
             Map<String, List<String>> task2groupsMap =
                     groupMap.computeIfAbsent(stream.getSortClusterName(), k -> 
new ConcurrentHashMap<>());
             List<String> groupIdList =
@@ -194,109 +204,96 @@ public class SortSourceServiceImpl implements 
SortSourceService {
             groupIdList.add(stream.getGroupId());
         });
 
-        // get all groups. group by group id.
-        List<SortSourceGroupInfo> groupInfos = 
inlongGroupEntityMapper.selectAllGroups();
-        Map<String, SortSourceGroupInfo> allId2GroupInfos = groupInfos.stream()
-                .filter(dto -> dto.getGroupId() != null)
-                .collect(Collectors.toMap(SortSourceGroupInfo::getGroupId, dto 
-> dto, (g1, g2) -> g1));
-
-        // get all clusters. filter by type and check if consumable, then 
group by cluster tag.
-        List<SortSourceClusterInfo> clusterInfos = 
clusterEntityMapper.selectAllClusters();
-        Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos = 
clusterInfos.stream()
-                .filter(dto -> dto.getClusterTags() != null)
-                .filter(SortSourceClusterInfo::isConsumable)
-                .filter(cluster -> 
SUPPORTED_MQ_TYPE.contains(cluster.getType().toUpperCase(Locale.ROOT)))
-                
.collect(Collectors.groupingBy(SortSourceClusterInfo::getClusterTags));
+        // reload all groups
+        groupInfos = configLoader.loadAllGroup()
+                .stream()
+                .collect(Collectors.toMap(SortSourceGroupInfo::getGroupId, 
info -> info));
+
+        // reload all back up cluster
+        backupClusterTag = 
configLoader.loadGroupBackupInfo(ClusterSwitch.BACKUP_CLUSTER_TAG)
+                .stream()
+                
.collect(Collectors.toMap(InlongGroupExtEntity::getInlongGroupId, 
InlongGroupExtEntity::getKeyValue));
+
+        // reload all back up group mq resource
+        backupGroupMqResource = 
configLoader.loadGroupBackupInfo(ClusterSwitch.BACKUP_MQ_RESOURCE)
+                .stream()
+                
.collect(Collectors.toMap(InlongGroupExtEntity::getInlongGroupId, 
InlongGroupExtEntity::getKeyValue));
+
+        // reload all streams
+        allStreams = configLoader.loadAllStreams()
+                .stream()
+                
.collect(Collectors.toMap(SortSourceStreamInfo::getInlongGroupId, stream -> 
stream));
+
+        // reload all back up stream mq resource
+        backupStreamMqResource = 
configLoader.loadStreamBackupInfo(ClusterSwitch.BACKUP_MQ_RESOURCE)
+                .stream()
+                
.collect(Collectors.toMap(InlongStreamExtEntity::getInlongGroupId, 
InlongStreamExtEntity::getKeyValue));
+    }
 
-        // group clusters by name.
-        Map<String, SortSourceClusterInfo> name2ClusterInfos = 
clusterInfos.stream()
-                .collect(Collectors.toMap(SortSourceClusterInfo::getName, info 
-> info, (g1, g2) -> g1));
+    private void parseAll() {
 
         // Prepare CacheZones for each cluster and task
         Map<String, Map<String, String>> newMd5Map = new ConcurrentHashMap<>();
         Map<String, Map<String, CacheZoneConfig>> newConfigMap = new 
ConcurrentHashMap<>();
-        groupMap.forEach((clusterName, task2Group) -> {
-
-            // if there is no matched cluster name, just skip
-            if (!name2ClusterInfos.containsKey(clusterName)) {
-                return;
-            }
-            // find valid mq cluster list
-            String clusterTag = 
name2ClusterInfos.get(clusterName).getClusterTags();
-            final Map<String, List<SortSourceClusterInfo>> validClusterInfos = 
new ConcurrentHashMap<>();
-            if (allTag2ClusterInfos.containsKey(clusterTag)) {
-                validClusterInfos.put(clusterTag, 
allTag2ClusterInfos.get(clusterTag));
-            } else {
-                validClusterInfos.putAll(allTag2ClusterInfos);
-            }
 
+        groupMap.forEach((sortClusterName, task2GroupList) -> {
             // prepare the new config and md5
             Map<String, CacheZoneConfig> task2Config = new 
ConcurrentHashMap<>();
             Map<String, String> task2Md5 = new ConcurrentHashMap<>();
 
-            task2Group.forEach((task, groupList) -> {
-                // get topic properties under this cluster and task, group 
them by group id.
-                Map<String, Map<String, String>> group2topicProp = new 
HashMap<>();
-                allStreamInfos.stream().filter(stream -> 
stream.getSortTaskName().equals(task)
-                        && 
stream.getSortClusterName().equals(clusterName)).forEach(
-                        sortSourceStreamInfo -> 
group2topicProp.put(sortSourceStreamInfo.getGroupId(),
-                                sortSourceStreamInfo.getExtParamsMap()));
-
-                Map<String, CacheZone> cacheZones;
+            task2GroupList.forEach((taskName, groupList) -> {
                 try {
-                    cacheZones = this.getCacheZones(groupList, 
allId2GroupInfos, validClusterInfos, group2topicProp);
+                    CacheZoneConfig cacheZoneConfig =
+                            CacheZoneConfig.builder()
+                                    .sortClusterName(sortClusterName)
+                                    .sortTaskId(taskName)
+                                    .build();
+                    Map<String, CacheZone> cacheZoneMap =
+                            this.parseCacheZones(sortClusterName, taskName, 
groupList);
+                    cacheZoneConfig.setCacheZones(cacheZoneMap);
+
+                    // prepare md5
+                    String jsonStr = GSON.toJson(cacheZoneConfig);
+                    String md5 = DigestUtils.md5Hex(jsonStr);
+                    task2Config.put(taskName, cacheZoneConfig);
+                    task2Md5.put(taskName, md5);
                 } catch (Throwable t) {
-                    LOGGER.error("fail to get cacheZones of clusterName {}, 
task {}", clusterName, task);
-                    return;
+                    LOGGER.error("failed to parse sort source config of 
sortCluster={}, task={}",
+                            sortClusterName, taskName, t);
                 }
-                CacheZoneConfig config = CacheZoneConfig.builder()
-                        .cacheZones(cacheZones)
-                        .sortClusterName(clusterName)
-                        .sortTaskId(task)
-                        .build();
-                String jsonStr = GSON.toJson(config);
-                String md5 = DigestUtils.md5Hex(jsonStr);
-                task2Config.put(task, config);
-                task2Md5.put(task, md5);
             });
+            newConfigMap.put(sortClusterName, task2Config);
+            newMd5Map.put(sortClusterName, task2Md5);
 
-            newConfigMap.put(clusterName, task2Config);
-            newMd5Map.put(clusterName, task2Md5);
         });
-
         sortSourceConfigMap = newConfigMap;
         sortSourceMd5Map = newMd5Map;
     }
 
-    private Map<String, CacheZone> getCacheZones(
-            List<String> groupIdList,
-            Map<String, SortSourceGroupInfo> allId2GroupInfos,
-            Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos,
-            Map<String, Map<String, String>> group2topicProp) {
+    private Map<String, CacheZone> parseCacheZones(
+            String sortClusterName,
+            String taskName,
+            List<String> groupIdList) {
 
-        // stream of group info if group id exists.
-        List<SortSourceGroupInfo> groupInfoStream = groupIdList.stream()
-                .filter(allId2GroupInfos::containsKey)
-                .map(allId2GroupInfos::get)
+        // get group infos
+        List<SortSourceGroupInfo> groupInfoList = groupIdList.stream()
+                .filter(groupInfos::containsKey)
+                .map(groupInfos::get)
                 .collect(Collectors.toList());
 
-        // Group them by cluster tag.
-        Map<String, List<SortSourceGroupInfo>> tag2GroupInfos = 
groupInfoStream.stream()
+        // group them by cluster tag.
+        Map<String, List<SortSourceGroupInfo>> tag2GroupInfos = 
groupInfoList.stream()
                 
.collect(Collectors.groupingBy(SortSourceGroupInfo::getClusterTag));
 
-        // Group them by back up cluster tag if both 2nd tag and 2nd topic 
exist.
-        Map<String, List<SortSourceGroupInfo>> backupTag2GroupInfos = 
groupInfoStream.stream()
-                .filter(group -> group.getBackupClusterTag() != null && 
group.getBackupTopic() != null)
-                
.collect(Collectors.groupingBy(SortSourceGroupInfo::getBackupClusterTag));
+        // group them by second cluster tag.
+        Map<String, List<SortSourceGroupInfo>> backupTag2GroupInfos = 
groupInfoList.stream()
+                .filter(info -> 
backupClusterTag.containsKey(info.getGroupId()))
+                .collect(Collectors.groupingBy(info -> 
backupClusterTag.get(info.getGroupId())));
 
-        // get cache zone list.
-        List<CacheZone> firstTagCacheZoneList =
-                this.getCacheZoneListByTag(tag2GroupInfos, 
allTag2ClusterInfos, group2topicProp, false);
-        List<CacheZone> backupTagCacheZoneList =
-                this.getCacheZoneListByTag(backupTag2GroupInfos, 
allTag2ClusterInfos, group2topicProp, true);
+        List<CacheZone> cacheZones = this.parseCacheZonesByTag(tag2GroupInfos, 
false);
+        List<CacheZone> backupCacheZones = 
this.parseCacheZonesByTag(backupTag2GroupInfos, true);
 
-        // combine two cache zone list, and group by cache zone name.
-        return Stream.of(firstTagCacheZoneList, backupTagCacheZoneList)
+        return Stream.of(cacheZones, backupCacheZones)
                 .flatMap(Collection::stream)
                 .collect(Collectors.toMap(
                         CacheZone::getZoneName,
@@ -308,86 +305,63 @@ public class SortSourceServiceImpl implements 
SortSourceService {
                 );
     }
 
-    private List<CacheZone> getCacheZoneListByTag(
-            Map<String, List<SortSourceGroupInfo>> tag2GroupInfos,
-            Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos,
-            Map<String, Map<String, String>> group2topicProp,
-            boolean isBackupTag) {
-
-        // Tags of groups
-        List<String> tags = new ArrayList<>(tag2GroupInfos.keySet());
+    private List<CacheZone> parseCacheZonesByTag(Map<String, 
List<SortSourceGroupInfo>> tag2Groups, boolean isBackup) {
 
-        // Clusters that related to these tags
-        Map<String, List<SortSourceClusterInfo>> tag2ClusterInfos = new 
HashMap<>();
-        allTag2ClusterInfos.entrySet().stream().filter(entry -> 
tag2GroupInfos.containsKey(entry.getKey()))
-                .forEach(entry -> tag2ClusterInfos.put(entry.getKey(), 
entry.getValue()));
-
-        // get CacheZone list
-        return tags.stream()
-                .filter(tag2ClusterInfos::containsKey)
+        return tag2Groups.keySet().stream()
+                .filter(mqClusters::containsKey)
                 .flatMap(tag -> {
-                    List<SortSourceGroupInfo> groups = tag2GroupInfos.get(tag);
-                    List<SortSourceClusterInfo> clusters = 
tag2ClusterInfos.get(tag);
+                    List<SortSourceGroupInfo> groups = tag2Groups.get(tag);
+                    List<SortSourceClusterInfo> clusters = mqClusters.get(tag);
                     return clusters.stream()
                             .map(cluster -> {
                                 CacheZone zone = null;
                                 try {
-                                    zone = this.getCacheZone(groups, cluster, 
group2topicProp, isBackupTag);
+                                    zone = this.parseCacheZone(groups, 
cluster, isBackup);
                                 } catch (IllegalStateException e) {
                                     LOGGER.error("fail to init cache zone for 
cluster " + cluster, e);
                                 }
                                 return zone;
                             });
                 })
-                .filter(Objects::nonNull)
                 .collect(Collectors.toList());
     }
 
-    private CacheZone getCacheZone(
+    private CacheZone parseCacheZone(
             List<SortSourceGroupInfo> groups,
             SortSourceClusterInfo cluster,
-            Map<String, Map<String, String>> group2topicProp,
             boolean isBackupTag) {
+        switch (cluster.getType()) {
+            case ClusterType.PULSAR: return parsePulsarZone(groups, cluster, 
isBackupTag);
+            default:
+                throw new BusinessException(String.format("do not support 
cluster type=%s of cluster=%s",
+                        cluster.getType(), cluster));
+        }
+    }
 
-        // get basic Cache zone fields
+    private CacheZone parsePulsarZone(
+            List<SortSourceGroupInfo> groups,
+            SortSourceClusterInfo cluster,
+            boolean isBackupTag) {
         Map<String, String> param = cluster.getExtParamsMap();
-        String serviceUrl = Optional.ofNullable(param.get(KEY_SERVICE_URL))
-                .orElseThrow(
-                        () -> new IllegalStateException(("there is no 
serviceUrl for cluster " + cluster.getName())));
         String tenant = param.get(KEY_TENANT);
-        String namespace = param.get(KEY_NAME_SPACE);
-        String authentication = 
Optional.ofNullable(param.get(KEY_AUTH)).orElse("");
-
-        List<Topic> topics = groups.stream()
-                .map(groupInfo -> getTopic(groupInfo, tenant, namespace,
-                        group2topicProp.get(groupInfo.getGroupId()), 
isBackupTag))
+        String auth = param.get(KEY_AUTH);
+        List<Topic> sdkTopics = groups.stream()
+                .map(info -> {
+                    String namespace = info.getMqResource();
+                    String topic = 
allStreams.get(info.getGroupId()).getMqResource();
+                    if (isBackupTag) {
+                        namespace = 
Optional.ofNullable(backupGroupMqResource.get(info.getGroupId())).orElse(namespace);
+                        topic = 
Optional.ofNullable(backupStreamMqResource.get(info.getGroupId())).orElse(topic);
+                    }
+                    String fullTopic = 
tenant.concat("/").concat(namespace).concat("/").concat(topic);
+                    return Topic.builder().topic(fullTopic).build();
+                })
                 .collect(Collectors.toList());
-
         return CacheZone.builder()
-                .serviceUrl(serviceUrl)
-                .authentication(authentication)
-                .cacheZoneProperties(param)
                 .zoneName(cluster.getName())
-                .zoneType(cluster.getType())
-                .topics(topics)
-                .build();
-    }
-
-    private Topic getTopic(
-            SortSourceGroupInfo groupInfo,
-            String tenant,
-            String namespace,
-            Map<String, String> topicProperties,
-            boolean isBackupTag) {
-
-        String topic = isBackupTag ? groupInfo.getBackupTopic() : 
groupInfo.getTopic();
-        StringBuilder fullTopic = new StringBuilder();
-        Optional.ofNullable(tenant).ifPresent(t -> 
fullTopic.append(t).append("/"));
-        Optional.ofNullable(namespace).ifPresent(n -> 
fullTopic.append(n).append("/"));
-        fullTopic.append(topic);
-        return Topic.builder()
-                .topic(fullTopic.toString())
-                .topicProperties(topicProperties)
+                .serviceUrl(cluster.getUrl())
+                .topics(sdkTopics)
+                .authentication(auth)
                 .build();
     }
 
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
index dc09f2be3..908e61220 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
@@ -135,16 +135,17 @@ public class InlongGroupOperator4Pulsar extends 
AbstractGroupOperator {
 
         // set backup topics, each inlong stream is associated with a Pulsar 
topic
         List<InlongStreamBriefInfo> streamTopics = 
streamService.getTopicList(groupId);
-        streamTopics.forEach(stream -> {
-            InlongStreamExtEntity streamExtEntity = 
streamExtMapper.selectByKey(groupId, stream.getInlongStreamId(),
-                    BACKUP_MQ_RESOURCE);
-            if (streamExtEntity != null && 
StringUtils.isNotBlank(streamExtEntity.getKeyValue())) {
-                stream.setMqResource(streamExtEntity.getKeyValue());
-            }
-        });
         List<String> topics = streamTopics.stream()
-                .map(InlongStreamBriefInfo::getMqResource)
+                .map(stream -> {
+                    InlongStreamExtEntity streamExtEntity = 
streamExtMapper.selectByKey(groupId,
+                            stream.getInlongStreamId(), BACKUP_MQ_RESOURCE);
+                    if (streamExtEntity != null && 
StringUtils.isNotBlank(streamExtEntity.getKeyValue())) {
+                        return streamExtEntity.getKeyValue();
+                    }
+                    return stream.getMqResource();
+                })
                 .collect(Collectors.toList());
+
         topicInfo.setTopics(topics);
 
         return topicInfo;
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
index ddcb6f0c5..f89fd0f70 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
@@ -389,7 +389,7 @@ public class InlongGroupServiceImpl implements 
InlongGroupService {
     public InlongGroupTopicInfo getBackupTopic(String groupId) {
         // backup topic info saved in the ext table
         InlongGroupExtEntity extEntity = 
groupExtMapper.selectByUniqueKey(groupId, BACKUP_CLUSTER_TAG);
-        if (StringUtils.isBlank(extEntity.getKeyValue())) {
+        if (extEntity == null || StringUtils.isBlank(extEntity.getKeyValue())) 
{
             LOGGER.warn("not found any backup topic for groupId={}", groupId);
             return null;
         }
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
index b65735236..f8eb9eff8 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
@@ -17,19 +17,28 @@
 
 package org.apache.inlong.manager.service.sort;
 
+import org.apache.inlong.common.constant.ClusterSwitch;
 import org.apache.inlong.common.pojo.sdk.SortSourceConfigResponse;
 import org.apache.inlong.common.pojo.sortstandalone.SortClusterResponse;
 import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.enums.ClusterType;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.dao.entity.InlongClusterEntity;
-import org.apache.inlong.manager.dao.entity.InlongGroupEntity;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
 import org.apache.inlong.manager.dao.mapper.DataNodeEntityMapper;
 import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
 import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
 import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.pojo.cluster.pulsar.PulsarClusterRequest;
+import org.apache.inlong.manager.pojo.group.InlongGroupExtInfo;
+import org.apache.inlong.manager.pojo.group.pulsar.InlongPulsarRequest;
+import org.apache.inlong.manager.pojo.stream.InlongStreamExtInfo;
+import org.apache.inlong.manager.pojo.stream.InlongStreamRequest;
 import org.apache.inlong.manager.service.ServiceBaseTest;
+import org.apache.inlong.manager.service.cluster.InlongClusterService;
 import org.apache.inlong.manager.service.core.SortService;
+import org.apache.inlong.manager.service.group.InlongGroupService;
+import org.apache.inlong.manager.service.stream.InlongStreamService;
 import org.json.JSONObject;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
@@ -40,7 +49,9 @@ import org.junit.jupiter.api.TestMethodOrder;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.transaction.annotation.Transactional;
 
+import java.util.ArrayList;
 import java.util.Date;
+import java.util.List;
 
 /**
  * Sort service test for {@link SortService}
@@ -51,7 +62,9 @@ public class SortServiceImplTest extends ServiceBaseTest {
     private static final String TEST_CLUSTER = "testCluster";
     private static final String TEST_TASK = "testTask";
     private static final String TEST_GROUP = "testGroup";
+    private static final String TEST_STREAM = "1";
     private static final String TEST_TAG = "testTag";
+    private static final String BACK_UP_TAG = "testBackupTag";
     private static final String TEST_TOPIC = "testTopic";
     private static final String TEST_SINK_TYPE = "testSinkType";
     private static final String TEST_CREATOR = "testUser";
@@ -66,6 +79,12 @@ public class SortServiceImplTest extends ServiceBaseTest {
     private DataNodeEntityMapper dataNodeEntityMapper;
     @Autowired
     private SortService sortService;
+    @Autowired
+    private InlongClusterService clusterService;
+    @Autowired
+    private InlongGroupService groupService;
+    @Autowired
+    private InlongStreamService streamService;
 
     @Test
     @Order(1)
@@ -173,10 +192,12 @@ public class SortServiceImplTest extends ServiceBaseTest {
     @BeforeEach
     private void prepareAll() {
         this.prepareCluster(TEST_CLUSTER);
-        this.prepareTask(TEST_TASK, TEST_GROUP, TEST_CLUSTER);
-        this.prepareGroupId(TEST_GROUP);
-        this.preparePulsar("testPulsar", true);
+        this.preparePulsar("testPulsar", true, TEST_TAG);
+        this.preparePulsar("testPulsar2", true, BACK_UP_TAG);
         this.prepareDataNode(TEST_TASK);
+        this.prepareGroupId(TEST_GROUP);
+        this.prepareStreamId(TEST_GROUP, TEST_STREAM);
+        this.prepareTask(TEST_TASK, TEST_GROUP, TEST_CLUSTER);
     }
 
     private void prepareDataNode(String taskName) {
@@ -195,20 +216,50 @@ public class SortServiceImplTest extends ServiceBaseTest {
     }
 
     private void prepareGroupId(String groupId) {
-        InlongGroupEntity entity = new InlongGroupEntity();
-        entity.setInlongGroupId(groupId);
-        entity.setInlongClusterTag(TEST_TAG);
-        entity.setMqResource(TEST_TOPIC);
-        entity.setMqType("PULSAR");
-        entity.setName("testName");
-        entity.setDescription("testDescription");
-        entity.setCreator(TEST_CREATOR);
-        entity.setInCharges(TEST_CREATOR);
-        entity.setCreateTime(new Date());
-        entity.setModifyTime(new Date());
-        entity.setIsDeleted(InlongConstants.UN_DELETED);
-        entity.setVersion(InlongConstants.INITIAL_VERSION);
-        inlongGroupEntityMapper.insert(entity);
+        InlongPulsarRequest request = new InlongPulsarRequest();
+        request.setInlongGroupId(groupId);
+        request.setMqResource("test_namespace");
+        request.setInlongClusterTag(TEST_TAG);
+        request.setVersion(InlongConstants.INITIAL_VERSION);
+        request.setName("test_group_name");
+        request.setMqType(ClusterType.PULSAR);
+        request.setInCharges(TEST_CREATOR);
+        List<InlongGroupExtInfo> extList = new ArrayList<>();
+        InlongGroupExtInfo ext1 = InlongGroupExtInfo
+                .builder()
+                .inlongGroupId(groupId)
+                .keyName(ClusterSwitch.BACKUP_CLUSTER_TAG)
+                .keyValue(BACK_UP_TAG)
+                .build();
+        InlongGroupExtInfo ext2 = InlongGroupExtInfo
+                .builder()
+                .inlongGroupId(groupId)
+                .keyName(ClusterSwitch.BACKUP_MQ_RESOURCE)
+                .keyValue("backup_name")
+                .build();
+
+        extList.add(ext1);
+        extList.add(ext2);
+        request.setExtList(extList);
+        groupService.save(request, "test operator");
+    }
+
+    private void prepareStreamId(String groupId, String streamId) {
+        InlongStreamRequest request = new InlongStreamRequest();
+        request.setInlongGroupId(groupId);
+        request.setInlongStreamId(streamId);
+        request.setName("test_stream_name");
+        request.setMqResource(TEST_TOPIC);
+        request.setVersion(InlongConstants.INITIAL_VERSION);
+        List<InlongStreamExtInfo> extInfos = new ArrayList<>();
+        InlongStreamExtInfo ext = new InlongStreamExtInfo();
+        extInfos.add(ext);
+        ext.setInlongStreamId(streamId);
+        ext.setInlongGroupId(groupId);
+        ext.setKeyName(ClusterSwitch.BACKUP_MQ_RESOURCE);
+        ext.setKeyValue("backup_topic");
+        request.setExtList(extInfos);
+        streamService.save(request, "test_operator");
     }
 
     private void prepareCluster(String clusterName) {
@@ -226,25 +277,21 @@ public class SortServiceImplTest extends ServiceBaseTest {
         clusterEntityMapper.insert(entity);
     }
 
-    private void preparePulsar(String pulsarName, boolean isConsumable) {
-        InlongClusterEntity entity = new InlongClusterEntity();
-        entity.setName(pulsarName);
-        entity.setType("PULSAR");
-        entity.setClusterTags(TEST_TAG);
-        entity.setCreator(TEST_CREATOR);
-        entity.setInCharges(TEST_CREATOR);
-        Date now = new Date();
-        entity.setCreateTime(now);
-        entity.setModifyTime(now);
-        entity.setIsDeleted(InlongConstants.UN_DELETED);
-        entity.setVersion(InlongConstants.INITIAL_VERSION);
+    private void preparePulsar(String pulsarName, boolean isConsumable, String 
tag) {
+        PulsarClusterRequest request = new PulsarClusterRequest();
+        request.setUrl("testServiceUrl");
+        request.setName(pulsarName);
+        request.setType(ClusterType.PULSAR);
+        request.setClusterTags(tag);
+        request.setInCharges(TEST_CREATOR);
+        request.setVersion(InlongConstants.INITIAL_VERSION);
         String extTag = "zone=" + TEST_TAG
                 + "&producer=true"
                 + "&consumer=" + (isConsumable ? "true" : "false");
-        entity.setExtTag(extTag);
-        
entity.setExtParams("{\"tenant\":\"testTenant\",\"namespace\":\"testNS\",\"serviceUrl\":\"testServiceUrl\","
+        request.setExtTag(extTag);
+        request.setExtParams("{\"tenant\":\"testTenant\","
                 + 
"\"authentication\":\"testAuth\",\"adminUrl\":\"testAdmin\"}");
-        clusterEntityMapper.insert(entity);
+        clusterService.save(request, "test operator");
     }
 
     private void prepareTask(String taskName, String groupId, String 
clusterName) {

Reply via email to