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 b65788b6f [INLONG-4949][Manager] Support to extend different types of 
Sort protocol (#4954)
b65788b6f is described below

commit b65788b6fb6f6f04e106fbdcfc424774f200cec4
Author: healchow <[email protected]>
AuthorDate: Mon Jul 11 21:41:56 2022 +0800

    [INLONG-4949][Manager] Support to extend different types of Sort protocol 
(#4954)
---
 .../apache/inlong/manager/client/BaseExample.java  |   6 +-
 .../apache/inlong/manager/client/ut/BaseTest.java  |   6 +-
 .../manager/common/consts/InlongConstants.java     |  13 +-
 .../inlong/manager/common/enums/GroupMode.java     |  13 +-
 .../form/process/GroupResourceProcessForm.java     |  13 --
 .../plugin/listener/RestartSortListener.java       |   6 +-
 .../plugin/listener/RestartStreamListener.java     |   6 +-
 .../plugin/listener/StartupSortListener.java       |   6 +-
 .../plugin/listener/StartupStreamListener.java     |   6 +-
 .../plugin/listener/RestartSortListenerTest.java   |   2 +-
 .../plugin/listener/StartupSortListenerTest.java   |   2 +-
 .../operation/InlongGroupProcessOperation.java     |  10 +-
 .../service/resource/SinkResourceListener.java     |   3 +-
 .../resource/StreamSinkResourceListener.java       |   3 +-
 .../service/sort/CreateSortConfigListener.java     | 114 ----------
 .../sort/CreateStreamSortConfigListener.java       | 184 ---------------
 .../service/sort/PushSortConfigListener.java       | 109 ---------
 ...V2.java => SortConfig4NormalGroupOperator.java} | 158 +++++++------
 .../manager/service/sort/SortConfigListener.java   |  88 ++++++++
 .../manager/service/sort/SortConfigOperator.java   |  47 ++++
 .../service/sort/SortConfigOperatorFactory.java    |  49 ++++
 .../service/sort/StreamSortConfigListener.java     |  95 ++++++++
 .../service/sort/ZookeeperDisabledSelector.java    |   5 +-
 .../service/sort/ZookeeperEnabledSelector.java     |   5 +-
 .../service/sort/light/LightGroupSortListener.java |   8 +-
 .../manager/service/sort/util/DataFlowUtils.java   |  94 --------
 .../manager/service/sort/util/SinkInfoUtils.java   | 246 ---------------------
 .../manager/service/sort/util/SourceInfoUtils.java | 112 ----------
 .../listener/AbstractSourceOperateListener.java    |   6 +-
 .../approve/GroupApproveProcessListener.java       |   2 +-
 .../listener/GroupTaskListenerFactory.java         |   6 +-
 .../listener/StreamTaskListenerFactory.java        |   6 +-
 .../manager/service/sort/DisableZkForSortTest.java |   6 +-
 33 files changed, 440 insertions(+), 995 deletions(-)

diff --git 
a/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/BaseExample.java
 
b/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/BaseExample.java
index 268001f1d..7f35072a3 100644
--- 
a/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/BaseExample.java
+++ 
b/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/BaseExample.java
@@ -81,9 +81,9 @@ public class BaseExample {
         pulsarInfo.setMqResource(namespace);
 
         // set enable zk, create resource, lightweight mode, and cluster tag
-        pulsarInfo.setEnableZookeeper(0);
-        pulsarInfo.setEnableCreateResource(1);
-        pulsarInfo.setLightweight(0);
+        pulsarInfo.setEnableZookeeper(InlongConstants.DISABLE_ZK);
+        
pulsarInfo.setEnableCreateResource(InlongConstants.ENABLE_CREATE_RESOURCE);
+        pulsarInfo.setLightweight(InlongConstants.NORMAL_MODE);
         pulsarInfo.setInlongClusterTag("default_cluster");
 
         pulsarInfo.setDailyRecords(10000000);
diff --git 
a/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/ut/BaseTest.java
 
b/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/ut/BaseTest.java
index a5f3182f8..ec06cbe44 100644
--- 
a/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/ut/BaseTest.java
+++ 
b/inlong-manager/manager-client-examples/src/test/java/org/apache/inlong/manager/client/ut/BaseTest.java
@@ -112,9 +112,9 @@ public class BaseTest {
         pulsarInfo.setMqResource(NAMESPACE);
 
         // set enable zk, create resource, lightweight mode, and cluster tag
-        pulsarInfo.setEnableZookeeper(0);
-        pulsarInfo.setEnableCreateResource(1);
-        pulsarInfo.setLightweight(1);
+        pulsarInfo.setEnableZookeeper(InlongConstants.DISABLE_ZK);
+        
pulsarInfo.setEnableCreateResource(InlongConstants.ENABLE_CREATE_RESOURCE);
+        pulsarInfo.setLightweight(InlongConstants.LIGHTWEIGHT_MODE);
         pulsarInfo.setInlongClusterTag("default_cluster");
 
         pulsarInfo.setDailyRecords(10000000);
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
index 2054de5f0..df368b0bb 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/consts/InlongConstants.java
@@ -23,18 +23,21 @@ package org.apache.inlong.manager.common.consts;
 public class InlongConstants {
 
     public static final Integer UN_DELETED = 0;
-
     public static final Integer IS_DELETED = 1;
 
     public static final Integer DELETED_STATUS = 10;
 
-    public static final Integer DISABLE_CREATE_RESOURCE = 0;
+    public static final Integer NORMAL_MODE = 0;
+    public static final Integer LIGHTWEIGHT_MODE = 1;
 
-    public static final Integer ENABLE_CREATE_RESOURCE = 1;
+    public static final Integer DISABLE_ZK = 0;
+    public static final Integer ENABLE_ZK = 1;
 
-    public static final Integer SYNC_SEND = 1;
+    public static final Integer DISABLE_CREATE_RESOURCE = 0;
+    public static final Integer ENABLE_CREATE_RESOURCE = 1;
 
     public static final Integer UN_SYNC_SEND = 0;
+    public static final Integer SYNC_SEND = 1;
 
     /**
      * Pulsar config
@@ -55,7 +58,7 @@ public class InlongConstants {
     /**
      * Sort config
      */
-    public static final String DATA_FLOW = "dataFlow";
+    public static final String DATAFLOW = "dataflow";
 
     public static final String STREAMS = "streams";
 
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupMode.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupMode.java
index 1cd2cd414..43e7606d4 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupMode.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupMode.java
@@ -18,8 +18,11 @@
 package org.apache.inlong.manager.common.enums;
 
 import lombok.Getter;
+import org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
 
+import java.util.Objects;
+
 /**
  * Mode of inlong group
  */
@@ -31,10 +34,10 @@ public enum GroupMode {
     NORMAL("normal"),
 
     /**
-     * Light group init with sort in Inlong Cluster
+     * Lightweight group init with sort in Inlong Cluster
      * StreamSource -> Sort -> StreamSink
      */
-    LIGHT("light");
+    LIGHTWEIGHT("lightweight");
 
     @Getter
     private final String mode;
@@ -49,12 +52,12 @@ public enum GroupMode {
                 return groupMode;
             }
         }
-        throw new IllegalArgumentException(String.format("Unsupported group 
mode=%s", mode));
+        throw new IllegalArgumentException(String.format("Unsupported group 
mode for %s", mode));
     }
 
     public static GroupMode parseGroupMode(InlongGroupInfo groupInfo) {
-        if (groupInfo.getLightweight() != null && groupInfo.getLightweight() 
== 1) {
-            return GroupMode.LIGHT;
+        if (Objects.equals(groupInfo.getLightweight(), 
InlongConstants.LIGHTWEIGHT_MODE)) {
+            return GroupMode.LIGHTWEIGHT;
         }
         return GroupMode.NORMAL;
     }
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/process/GroupResourceProcessForm.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/process/GroupResourceProcessForm.java
index e28191d04..03c9e15cb 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/process/GroupResourceProcessForm.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/process/GroupResourceProcessForm.java
@@ -45,9 +45,6 @@ public class GroupResourceProcessForm extends BaseProcessForm 
{
     @Setter
     private GroupOperateType groupOperateType = GroupOperateType.INIT;
 
-    @Deprecated
-    private String streamId;
-
     private List<InlongStreamInfo> streamInfos;
 
     public InlongGroupInfo getGroupInfo() {
@@ -72,16 +69,6 @@ public class GroupResourceProcessForm extends 
BaseProcessForm {
         return groupInfo.getInlongGroupId();
     }
 
-    @Deprecated
-    public String getInlongStreamId() {
-        return streamId;
-    }
-
-    @Deprecated
-    public void setInlongStreamId(String streamId) {
-        this.streamId = streamId;
-    }
-
     @Override
     public Map<String, Object> showInList() {
         Map<String, Object> show = Maps.newHashMap();
diff --git 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartSortListener.java
 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartSortListener.java
index f845bcb47..6ef6f2a6c 100644
--- 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartSortListener.java
+++ 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartSortListener.java
@@ -88,8 +88,8 @@ public class RestartSortListener implements 
SortOperateListener {
             String message = String.format("sort job id is empty for groupId 
[%s]", groupId);
             return ListenerResult.fail(message);
         }
-        String dataFlows = kvConf.get(InlongConstants.DATA_FLOW);
-        if (StringUtils.isEmpty(dataFlows)) {
+        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
+        if (StringUtils.isEmpty(dataflow)) {
             String message = String.format("dataflow is empty for groupId 
[%s]", groupId);
             log.error(message);
             return ListenerResult.fail(message);
@@ -105,7 +105,7 @@ public class RestartSortListener implements 
SortOperateListener {
         FlinkService flinkService = new FlinkService(flinkInfo.getEndpoint());
         FlinkOperation flinkOperation = new FlinkOperation(flinkService);
         try {
-            flinkOperation.genPath(flinkInfo, dataFlows);
+            flinkOperation.genPath(flinkInfo, dataflow);
             flinkOperation.restart(flinkInfo);
             log.info("job restart success for [{}]", jobId);
             return ListenerResult.success();
diff --git 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartStreamListener.java
 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartStreamListener.java
index 1c59755ec..157153e32 100644
--- 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartStreamListener.java
+++ 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/RestartStreamListener.java
@@ -91,8 +91,8 @@ public class RestartStreamListener implements 
SortOperateListener {
             String message = String.format("sort job id is empty for groupId 
[%s] streamId [%s]", groupId, streamId);
             return ListenerResult.fail(message);
         }
-        String dataFlows = kvConf.get(InlongConstants.DATA_FLOW);
-        if (StringUtils.isEmpty(dataFlows)) {
+        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
+        if (StringUtils.isEmpty(dataflow)) {
             String message = String.format("dataflow is empty for groupId [%s] 
streamId [%s]", groupId, streamId);
             log.error(message);
             return ListenerResult.fail(message);
@@ -108,7 +108,7 @@ public class RestartStreamListener implements 
SortOperateListener {
         FlinkService flinkService = new FlinkService(flinkInfo.getEndpoint());
         FlinkOperation flinkOperation = new FlinkOperation(flinkService);
         try {
-            flinkOperation.genPath(flinkInfo, dataFlows);
+            flinkOperation.genPath(flinkInfo, dataflow);
             flinkOperation.restart(flinkInfo);
             log.info("job restart success for [{}]", jobId);
             return ListenerResult.success();
diff --git 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
index db1e8cf1c..5c80896a7 100644
--- 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
+++ 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
@@ -91,8 +91,8 @@ public class StartupSortListener implements 
SortOperateListener {
             kvConf.putAll(result);
         }
 
-        String dataFlows = kvConf.get(InlongConstants.DATA_FLOW);
-        if (StringUtils.isEmpty(dataFlows)) {
+        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
+        if (StringUtils.isEmpty(dataflow)) {
             String message = String.format("dataflow is empty for groupId 
[%s]", groupId);
             log.error(message);
             return ListenerResult.fail(message);
@@ -109,7 +109,7 @@ public class StartupSortListener implements 
SortOperateListener {
         FlinkOperation flinkOperation = new FlinkOperation(flinkService);
 
         try {
-            flinkOperation.genPath(flinkInfo, dataFlows);
+            flinkOperation.genPath(flinkInfo, dataflow);
             flinkOperation.start(flinkInfo);
             log.info("job submit success, jobId is [{}]", 
flinkInfo.getJobId());
         } catch (Exception e) {
diff --git 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
index f880a9151..adbbf11ef 100644
--- 
a/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
+++ 
b/inlong-manager/manager-plugins/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
@@ -88,8 +88,8 @@ public class StartupStreamListener implements 
SortOperateListener {
             kvConf.putAll(result);
         }
 
-        String dataFlows = kvConf.get(InlongConstants.DATA_FLOW);
-        if (StringUtils.isEmpty(dataFlows)) {
+        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
+        if (StringUtils.isEmpty(dataflow)) {
             String message = String.format("dataflow is empty for groupId [%s] 
and streamId [%s]", groupId, streamId);
             log.error(message);
             return ListenerResult.fail(message);
@@ -105,7 +105,7 @@ public class StartupStreamListener implements 
SortOperateListener {
         FlinkOperation flinkOperation = new FlinkOperation(flinkService);
 
         try {
-            flinkOperation.genPath(flinkInfo, dataFlows);
+            flinkOperation.genPath(flinkInfo, dataflow);
             flinkOperation.start(flinkInfo);
             log.info("job submit success, jobId is [{}]", 
flinkInfo.getJobId());
         } catch (Exception e) {
diff --git 
a/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/RestartSortListenerTest.java
 
b/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/RestartSortListenerTest.java
index a9990d33b..80d36359e 100644
--- 
a/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/RestartSortListenerTest.java
+++ 
b/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/RestartSortListenerTest.java
@@ -64,7 +64,7 @@ public class RestartSortListenerTest {
         inlongGroupExtInfoList.add(inlongGroupExtInfo5);
 
         InlongGroupExtInfo inlongGroupExtInfo6 = new InlongGroupExtInfo();
-        inlongGroupExtInfo6.setKeyName(InlongConstants.DATA_FLOW);
+        inlongGroupExtInfo6.setKeyName(InlongConstants.DATAFLOW);
         inlongGroupExtInfo6.setKeyValue("{\"streamId\":{\n"
                 + "    \"id\":1,\n"
                 + "    \"source_info\":{\n"
diff --git 
a/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/StartupSortListenerTest.java
 
b/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/StartupSortListenerTest.java
index 7c4b963b0..c47f5f71c 100644
--- 
a/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/StartupSortListenerTest.java
+++ 
b/inlong-manager/manager-plugins/src/test/java/org/apache/inlong/manager/plugin/listener/StartupSortListenerTest.java
@@ -61,7 +61,7 @@ public class StartupSortListenerTest {
         inlongGroupExtInfos.add(inlongGroupExtInfo2);
 
         InlongGroupExtInfo inlongGroupExtInfo5 = new InlongGroupExtInfo();
-        inlongGroupExtInfo5.setKeyName(InlongConstants.DATA_FLOW);
+        inlongGroupExtInfo5.setKeyName(InlongConstants.DATAFLOW);
         inlongGroupExtInfo5.setKeyValue("{\"streamId\":{\n"
                 + "    \"id\": 1,\n"
                 + "    \"source_info\":\n"
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
index e8b632665..dd6ac6724 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
@@ -110,7 +110,7 @@ public class InlongGroupProcessOperation {
                 GroupResourceProcessForm form = genGroupProcessForm(groupInfo, 
GroupOperateType.SUSPEND);
                 executorService.execute(() -> 
workflowService.start(ProcessName.SUSPEND_GROUP_PROCESS, operator, form));
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 LightGroupResourceProcessForm lightForm = 
genLightGroupProcessForm(groupInfo, GroupOperateType.SUSPEND);
                 executorService.execute(
                         () -> 
workflowService.start(ProcessName.SUSPEND_LIGHT_GROUP_PROCESS, operator, 
lightForm));
@@ -139,7 +139,7 @@ public class InlongGroupProcessOperation {
                 GroupResourceProcessForm form = genGroupProcessForm(groupInfo, 
GroupOperateType.SUSPEND);
                 result = 
workflowService.start(ProcessName.SUSPEND_GROUP_PROCESS, operator, form);
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 LightGroupResourceProcessForm lightForm = 
genLightGroupProcessForm(groupInfo, GroupOperateType.SUSPEND);
                 result = 
workflowService.start(ProcessName.SUSPEND_LIGHT_GROUP_PROCESS, operator, 
lightForm);
                 break;
@@ -165,7 +165,7 @@ public class InlongGroupProcessOperation {
                 GroupResourceProcessForm form = genGroupProcessForm(groupInfo, 
GroupOperateType.RESTART);
                 executorService.execute(() -> 
workflowService.start(ProcessName.RESTART_GROUP_PROCESS, operator, form));
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 LightGroupResourceProcessForm lightForm = 
genLightGroupProcessForm(groupInfo, GroupOperateType.RESTART);
                 executorService.execute(
                         () -> 
workflowService.start(ProcessName.RESTART_LIGHT_GROUP_PROCESS, operator, 
lightForm));
@@ -193,7 +193,7 @@ public class InlongGroupProcessOperation {
                 GroupResourceProcessForm form = genGroupProcessForm(groupInfo, 
GroupOperateType.RESTART);
                 result = 
workflowService.start(ProcessName.RESTART_GROUP_PROCESS, operator, form);
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 LightGroupResourceProcessForm lightForm = 
genLightGroupProcessForm(groupInfo, GroupOperateType.RESTART);
                 result = 
workflowService.start(ProcessName.RESTART_LIGHT_GROUP_PROCESS, operator, 
lightForm);
                 break;
@@ -306,7 +306,7 @@ public class InlongGroupProcessOperation {
                 GroupResourceProcessForm form = genGroupProcessForm(groupInfo, 
GroupOperateType.DELETE);
                 workflowService.start(ProcessName.DELETE_GROUP_PROCESS, 
operator, form);
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 LightGroupResourceProcessForm lightForm = 
genLightGroupProcessForm(groupInfo,
                         GroupOperateType.DELETE);
                 workflowService.start(ProcessName.DELETE_LIGHT_GROUP_PROCESS, 
operator, lightForm);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/SinkResourceListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/SinkResourceListener.java
index b46174b5c..8f9f4fffd 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/SinkResourceListener.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/SinkResourceListener.java
@@ -37,7 +37,8 @@ import java.util.List;
 import java.util.stream.Collectors;
 
 /**
- * Event listener of operate sink resources, such as create or update Hive 
table, Kafka topics, ES indices, etc.
+ * Event listener of operate sink resources,
+ * such as create or update Hive table, Kafka topics, ES indices, etc.
  */
 @Slf4j
 @Service
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/StreamSinkResourceListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/StreamSinkResourceListener.java
index cdaf78f47..44ce6bdcf 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/StreamSinkResourceListener.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/StreamSinkResourceListener.java
@@ -37,7 +37,8 @@ import java.util.List;
 import java.util.stream.Collectors;
 
 /**
- * Event listener of create hive table for one inlong stream
+ * Event listener of operate sink resources for one inlong stream,
+ * such as create or update Hive table, Kafka topics, ES indices, etc.
  */
 @Service
 @Slf4j
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListener.java
deleted file mode 100644
index 422a90ec0..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListener.java
+++ /dev/null
@@ -1,114 +0,0 @@
-/*
- * 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.sort;
-
-import com.google.common.collect.Lists;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.commons.lang3.StringUtils;
-import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.inlong.manager.common.consts.InlongConstants;
-import org.apache.inlong.manager.common.enums.GroupOperateType;
-import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
-import org.apache.inlong.manager.common.pojo.group.InlongGroupExtInfo;
-import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
-import org.apache.inlong.manager.common.pojo.sink.StreamSink;
-import 
org.apache.inlong.manager.common.pojo.workflow.form.process.GroupResourceProcessForm;
-import org.apache.inlong.manager.service.sink.StreamSinkService;
-import org.apache.inlong.manager.service.sort.util.DataFlowUtils;
-import org.apache.inlong.manager.workflow.WorkflowContext;
-import org.apache.inlong.manager.workflow.event.ListenerResult;
-import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
-import org.apache.inlong.manager.workflow.event.task.TaskEvent;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Component;
-
-import java.util.List;
-
-/**
- * Create sort config when disable the ZooKeeper
- */
-@Deprecated
-@Component
-public class CreateSortConfigListener implements SortOperateListener {
-
-    private static final Logger LOGGER = 
LoggerFactory.getLogger(CreateSortConfigListener.class);
-    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); // 
thread safe
-
-    @Autowired
-    private StreamSinkService streamSinkService;
-    @Autowired
-    private DataFlowUtils dataFlowUtils;
-
-    @Override
-    public TaskEvent event() {
-        return TaskEvent.COMPLETE;
-    }
-
-    @Override
-    public ListenerResult listen(WorkflowContext context) throws 
WorkflowListenerException {
-        LOGGER.info("Create sort config for groupId={}", 
context.getProcessForm().getInlongGroupId());
-        GroupResourceProcessForm form = (GroupResourceProcessForm) 
context.getProcessForm();
-        GroupOperateType groupOperateType = form.getGroupOperateType();
-        if (groupOperateType == GroupOperateType.SUSPEND || groupOperateType 
== GroupOperateType.DELETE) {
-            return ListenerResult.success();
-        }
-        InlongGroupInfo groupInfo = form.getGroupInfo();
-        String groupId = groupInfo.getInlongGroupId();
-        if (StringUtils.isEmpty(groupId)) {
-            LOGGER.warn("GroupId is null for context={}", context);
-            return ListenerResult.success();
-        }
-
-        List<StreamSink> streamSinks = streamSinkService.listSink(groupId, 
null);
-        if (CollectionUtils.isEmpty(streamSinks)) {
-            LOGGER.warn("Sink not found by groupId={}", groupId);
-            return ListenerResult.success();
-        }
-
-        try {
-            // TODO Support more than one sinks under a stream
-            // Map<String, DataFlowInfo> dataFlowInfoMap = 
streamSinks.stream().map(sink -> {
-            //             DataFlowInfo flowInfo = 
dataFlowUtils.createDataFlow(groupInfo, sink);
-            //             return Pair.of(sink.getInlongStreamId(), flowInfo);
-            //         }
-            // ).collect(Collectors.toMap(Pair::getKey, Pair::getValue));
-
-            // String dataFlows = 
OBJECT_MAPPER.writeValueAsString(dataFlowInfoMap);
-            InlongGroupExtInfo extInfo = new InlongGroupExtInfo();
-            extInfo.setInlongGroupId(groupId);
-            extInfo.setKeyName(InlongConstants.DATA_FLOW);
-            // extInfo.setKeyValue(dataFlows);
-            if (groupInfo.getExtList() == null) {
-                groupInfo.setExtList(Lists.newArrayList());
-            }
-            upsertDataFlow(groupInfo, extInfo);
-        } catch (Exception e) {
-            LOGGER.error("create sort config failed for sink list={} ", 
streamSinks, e);
-            throw new WorkflowListenerException("create sort config failed: " 
+ e.getMessage());
-        }
-        return ListenerResult.success();
-    }
-
-    private void upsertDataFlow(InlongGroupInfo groupInfo, InlongGroupExtInfo 
extInfo) {
-        groupInfo.getExtList().removeIf(ext -> 
InlongConstants.DATA_FLOW.equals(ext.getKeyName()));
-        groupInfo.getExtList().add(extInfo);
-    }
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateStreamSortConfigListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateStreamSortConfigListener.java
deleted file mode 100644
index e32e7f9b3..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateStreamSortConfigListener.java
+++ /dev/null
@@ -1,184 +0,0 @@
-/*
- * 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.sort;
-
-import com.google.common.collect.Lists;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.commons.lang3.StringUtils;
-import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.inlong.manager.common.consts.InlongConstants;
-import org.apache.inlong.manager.common.enums.ClusterType;
-import org.apache.inlong.manager.common.enums.GroupOperateType;
-import org.apache.inlong.manager.common.enums.MQType;
-import org.apache.inlong.manager.common.enums.SourceType;
-import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
-import org.apache.inlong.manager.common.pojo.cluster.ClusterInfo;
-import org.apache.inlong.manager.common.pojo.cluster.pulsar.PulsarClusterInfo;
-import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
-import org.apache.inlong.manager.common.pojo.sink.StreamSink;
-import org.apache.inlong.manager.common.pojo.source.StreamSource;
-import org.apache.inlong.manager.common.pojo.source.kafka.KafkaSource;
-import org.apache.inlong.manager.common.pojo.source.pulsar.PulsarSource;
-import org.apache.inlong.manager.common.pojo.stream.InlongStreamExtInfo;
-import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
-import 
org.apache.inlong.manager.common.pojo.workflow.form.process.StreamResourceProcessForm;
-import org.apache.inlong.manager.service.cluster.InlongClusterService;
-import org.apache.inlong.manager.service.sink.StreamSinkService;
-import org.apache.inlong.manager.service.sort.util.ExtractNodeUtils;
-import org.apache.inlong.manager.service.sort.util.LoadNodeUtils;
-import org.apache.inlong.manager.service.source.StreamSourceService;
-import org.apache.inlong.manager.workflow.WorkflowContext;
-import org.apache.inlong.manager.workflow.event.ListenerResult;
-import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
-import org.apache.inlong.manager.workflow.event.task.TaskEvent;
-import org.apache.inlong.sort.protocol.GroupInfo;
-import org.apache.inlong.sort.protocol.StreamInfo;
-import org.apache.inlong.sort.protocol.enums.PulsarScanStartupMode;
-import org.apache.inlong.sort.protocol.node.Node;
-import org.apache.inlong.sort.protocol.transformation.relation.NodeRelation;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Component;
-
-import java.util.List;
-import java.util.stream.Collectors;
-
-/**
- * Create sort config for one stream if Zookeeper is disabled.
- */
-@Slf4j
-@Component
-public class CreateStreamSortConfigListener implements SortOperateListener {
-
-    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); // 
thread safe
-
-    @Autowired
-    private StreamSourceService streamSourceService;
-    @Autowired
-    private StreamSinkService streamSinkService;
-    @Autowired
-    private InlongClusterService clusterService;
-
-    @Override
-    public TaskEvent event() {
-        return TaskEvent.COMPLETE;
-    }
-
-    @Override
-    public ListenerResult listen(WorkflowContext context) throws Exception {
-        StreamResourceProcessForm form = (StreamResourceProcessForm) 
context.getProcessForm();
-        GroupOperateType groupOperateType = form.getGroupOperateType();
-        if (groupOperateType == GroupOperateType.SUSPEND || groupOperateType 
== GroupOperateType.DELETE) {
-            return ListenerResult.success();
-        }
-        InlongGroupInfo groupInfo = form.getGroupInfo();
-        InlongStreamInfo streamInfo = form.getStreamInfo();
-        final String groupId = streamInfo.getInlongGroupId();
-        final String streamId = streamInfo.getInlongStreamId();
-        List<StreamSink> streamSinks = streamInfo.getSinkList();
-        if (CollectionUtils.isEmpty(streamSinks)) {
-            log.warn("Sink not found by groupId={}", groupId);
-            return ListenerResult.success();
-        }
-        try {
-            List<StreamSource> sources = createPulsarSources(groupInfo, 
streamInfo);
-            List<Node> nodes = createNodesForStream(sources, streamSinks);
-            List<NodeRelation> nodeRelations = 
createNodeRelationsForStream(sources, streamSinks);
-            StreamInfo sortStreamInfo = new StreamInfo(streamId, nodes, 
nodeRelations);
-            GroupInfo sortGroupInfo = new GroupInfo(groupId, 
Lists.newArrayList(sortStreamInfo));
-            String dataFlows = OBJECT_MAPPER.writeValueAsString(sortGroupInfo);
-            addExtInfo(groupInfo, streamInfo, InlongConstants.DATA_FLOW, 
dataFlows);
-        } catch (Exception e) {
-            log.error("create sort config failed for sink list={} of 
groupId={}, streamId={}", streamSinks, groupId,
-                    streamId, e);
-            throw new WorkflowListenerException("create sort config failed: " 
+ e.getMessage());
-        }
-        return ListenerResult.success();
-    }
-
-    private void addExtInfo(InlongGroupInfo groupInfo, InlongStreamInfo 
streamInfo, String key, String value) {
-        InlongStreamExtInfo extInfo = new InlongStreamExtInfo();
-        extInfo.setInlongGroupId(groupInfo.getInlongGroupId());
-        extInfo.setInlongStreamId(streamInfo.getInlongStreamId());
-        extInfo.setKeyName(key);
-        extInfo.setKeyValue(value);
-        if (streamInfo.getExtList() == null) {
-            streamInfo.setExtList(Lists.newArrayList());
-        }
-        upsertExtInfo(streamInfo, extInfo, key);
-    }
-
-    private List<StreamSource> createPulsarSources(InlongGroupInfo groupInfo, 
InlongStreamInfo streamInfo) {
-        if (!MQType.MQ_PULSAR.equals(groupInfo.getMqType())) {
-            String errMsg = String.format("Unsupported MQ type %s", 
groupInfo.getMqType());
-            log.error(errMsg);
-            throw new WorkflowListenerException(errMsg);
-        }
-
-        PulsarSource pulsarSource = new PulsarSource();
-        String streamId = streamInfo.getInlongStreamId();
-        pulsarSource.setSourceName(streamId);
-        pulsarSource.setNamespace(groupInfo.getMqResource());
-        pulsarSource.setTopic(streamInfo.getMqResource());
-
-        ClusterInfo clusterInfo = 
clusterService.getOne(groupInfo.getInlongClusterTag(), null,
-                ClusterType.PULSAR);
-        PulsarClusterInfo pulsarCluster = (PulsarClusterInfo) clusterInfo;
-        String adminUrl = pulsarCluster.getAdminUrl();
-        String serviceUrl = pulsarCluster.getUrl();
-
-        pulsarSource.setAdminUrl(adminUrl);
-        pulsarSource.setServiceUrl(serviceUrl);
-        pulsarSource.setInlongComponent(true);
-        List<StreamSource> sources = 
streamSourceService.listSource(groupInfo.getInlongGroupId(), streamId);
-        for (StreamSource source : sources) {
-            if (StringUtils.isEmpty(pulsarSource.getSerializationType())
-                    && StringUtils.isNotEmpty(source.getSerializationType())) {
-                
pulsarSource.setSerializationType(source.getSerializationType());
-            }
-            if (SourceType.forType(source.getSourceType()) == 
SourceType.KAFKA) {
-                pulsarSource.setPrimaryKey(((KafkaSource) 
source).getPrimaryKey());
-            }
-        }
-        
pulsarSource.setScanStartupMode(PulsarScanStartupMode.EARLIEST.getValue());
-        pulsarSource.setFieldList(streamInfo.getFieldList());
-        return Lists.newArrayList(pulsarSource);
-    }
-
-    private List<Node> createNodesForStream(List<StreamSource> sources, 
List<StreamSink> streamSinks) {
-        List<Node> nodes = Lists.newArrayList();
-        nodes.addAll(ExtractNodeUtils.createExtractNodes(sources));
-        nodes.addAll(LoadNodeUtils.createLoadNodes(streamSinks));
-        return nodes;
-    }
-
-    private List<NodeRelation> createNodeRelationsForStream(List<StreamSource> 
sources, List<StreamSink> streamSinks) {
-        NodeRelation relation = new NodeRelation();
-        List<String> inputs = 
sources.stream().map(StreamSource::getSourceName).collect(Collectors.toList());
-        List<String> outputs = 
streamSinks.stream().map(StreamSink::getSinkName).collect(Collectors.toList());
-        relation.setInputs(inputs);
-        relation.setOutputs(outputs);
-        return Lists.newArrayList(relation);
-    }
-
-    private void upsertExtInfo(InlongStreamInfo streamInfo, 
InlongStreamExtInfo extInfo, String keyName) {
-        streamInfo.getExtList().removeIf(ext -> 
keyName.equals(ext.getKeyName()));
-        streamInfo.getExtList().add(extInfo);
-    }
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/PushSortConfigListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/PushSortConfigListener.java
deleted file mode 100644
index c37dc658e..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/PushSortConfigListener.java
+++ /dev/null
@@ -1,109 +0,0 @@
-/*
- * 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.sort;
-
-import com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
-import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
-import org.apache.inlong.manager.common.pojo.sink.StreamSink;
-import 
org.apache.inlong.manager.common.pojo.workflow.form.process.GroupResourceProcessForm;
-import org.apache.inlong.manager.service.group.InlongGroupService;
-import org.apache.inlong.manager.service.sink.StreamSinkService;
-import org.apache.inlong.manager.service.sort.util.DataFlowUtils;
-import org.apache.inlong.manager.workflow.WorkflowContext;
-import org.apache.inlong.manager.workflow.event.ListenerResult;
-import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
-import org.apache.inlong.manager.workflow.event.task.TaskEvent;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Component;
-
-import java.util.List;
-
-/**
- * Push sort config when enable the ZooKeeper
- */
-@Deprecated
-@Component
-public class PushSortConfigListener implements SortOperateListener {
-
-    private static final Logger LOGGER = 
LoggerFactory.getLogger(PushSortConfigListener.class);
-
-    @Autowired
-    private InlongGroupService groupService;
-    @Autowired
-    private StreamSinkService streamSinkService;
-    @Autowired
-    private DataFlowUtils dataFlowUtils;
-    @Autowired
-    private ObjectMapper objectMapper;
-
-    @Override
-    public TaskEvent event() {
-        return TaskEvent.COMPLETE;
-    }
-
-    @Override
-    public ListenerResult listen(WorkflowContext context) throws 
WorkflowListenerException {
-        if (LOGGER.isDebugEnabled()) {
-            LOGGER.debug("begin to push sort config by context={}", context);
-        }
-
-        GroupResourceProcessForm form = (GroupResourceProcessForm) 
context.getProcessForm();
-        String groupId = form.getGroupInfo().getInlongGroupId();
-        InlongGroupInfo groupInfo = groupService.get(groupId);
-
-        // if streamId not null, just push the config belongs to the groupId 
and the streamId
-        String streamId = form.getInlongStreamId();
-        List<StreamSink> streamSinks = streamSinkService.listSink(groupId, 
streamId);
-        if (CollectionUtils.isEmpty(streamSinks)) {
-            LOGGER.warn("Sink not found by groupId={}", groupId);
-            return ListenerResult.success();
-        }
-
-        for (StreamSink streamSink : streamSinks) {
-            if (LOGGER.isDebugEnabled()) {
-                LOGGER.debug("sink info: {}", streamSink);
-            }
-
-            Integer sinkId = streamSink.getId();
-            try {
-                // DataFlowInfo dataFlowInfo = 
dataFlowUtils.createDataFlow(groupInfo, streamSink);
-                // String zkUrl = clusterBean.getZkUrl();
-                // String zkRoot = clusterBean.getZkRoot();
-                // push data flow info to zk
-                // String sortClusterName = clusterBean.getAppName();
-                // ZkTools.updateDataFlowInfo(dataFlowInfo, sortClusterName, 
sinkId, zkUrl, zkRoot);
-                // add sink id to zk
-                // ZkTools.addDataFlowToCluster(sortClusterName, sinkId, 
zkUrl, zkRoot);
-
-                if (LOGGER.isDebugEnabled()) {
-                    // LOGGER.debug("success to push config to sort:{}", 
objectMapper.writeValueAsString(dataFlowInfo));
-                }
-            } catch (Exception e) {
-                LOGGER.error("push sort config to zookeeper failed, sinkId={} 
", sinkId, e);
-                throw new WorkflowListenerException("push sort config to 
zookeeper failed: " + e.getMessage());
-            }
-        }
-
-        return ListenerResult.success();
-    }
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListenerV2.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfig4NormalGroupOperator.java
similarity index 66%
rename from 
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListenerV2.java
rename to 
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfig4NormalGroupOperator.java
index b4f2f68be..dc4df0c09 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/CreateSortConfigListenerV2.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfig4NormalGroupOperator.java
@@ -19,14 +19,12 @@ package org.apache.inlong.manager.service.sort;
 
 import com.google.common.collect.Lists;
 import com.google.common.collect.Maps;
-import lombok.SneakyThrows;
-import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.collections.CollectionUtils;
 import org.apache.commons.lang3.StringUtils;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
 import org.apache.inlong.common.enums.DataTypeEnum;
 import org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.enums.ClusterType;
-import org.apache.inlong.manager.common.enums.GroupOperateType;
 import org.apache.inlong.manager.common.enums.MQType;
 import org.apache.inlong.manager.common.enums.SourceType;
 import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
@@ -38,24 +36,22 @@ import 
org.apache.inlong.manager.common.pojo.sink.StreamSink;
 import org.apache.inlong.manager.common.pojo.source.StreamSource;
 import org.apache.inlong.manager.common.pojo.source.kafka.KafkaSource;
 import org.apache.inlong.manager.common.pojo.source.pulsar.PulsarSource;
+import org.apache.inlong.manager.common.pojo.stream.InlongStreamExtInfo;
 import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
-import 
org.apache.inlong.manager.common.pojo.workflow.form.process.GroupResourceProcessForm;
 import org.apache.inlong.manager.service.cluster.InlongClusterService;
 import org.apache.inlong.manager.service.sink.StreamSinkService;
 import org.apache.inlong.manager.service.sort.util.ExtractNodeUtils;
 import org.apache.inlong.manager.service.sort.util.LoadNodeUtils;
 import org.apache.inlong.manager.service.source.StreamSourceService;
-import org.apache.inlong.manager.workflow.WorkflowContext;
-import org.apache.inlong.manager.workflow.event.ListenerResult;
-import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
-import org.apache.inlong.manager.workflow.event.task.TaskEvent;
 import org.apache.inlong.sort.protocol.GroupInfo;
 import org.apache.inlong.sort.protocol.StreamInfo;
 import org.apache.inlong.sort.protocol.enums.PulsarScanStartupMode;
 import org.apache.inlong.sort.protocol.node.Node;
 import org.apache.inlong.sort.protocol.transformation.relation.NodeRelation;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Component;
+import org.springframework.stereotype.Service;
 
 import java.util.ArrayList;
 import java.util.HashMap;
@@ -63,10 +59,13 @@ import java.util.List;
 import java.util.Map;
 import java.util.stream.Collectors;
 
-@Component
-@Slf4j
-public class CreateSortConfigListenerV2 implements SortOperateListener {
+/**
+ * Sort config operator, used to create a Sort config for the InlongGroup in 
normal mode with ZK disabled.
+ */
+@Service
+public class SortConfig4NormalGroupOperator implements SortConfigOperator {
 
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(SortConfig4NormalGroupOperator.class);
     private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
 
     @Autowired
@@ -77,92 +76,86 @@ public class CreateSortConfigListenerV2 implements 
SortOperateListener {
     private InlongClusterService clusterService;
 
     @Override
-    public TaskEvent event() {
-        return TaskEvent.COMPLETE;
+    public Boolean accept(Integer isNormal, Integer enableZk) {
+        return InlongConstants.NORMAL_MODE.equals(isNormal)
+                && InlongConstants.DISABLE_ZK.equals(enableZk);
     }
 
-    @SneakyThrows
     @Override
-    public ListenerResult listen(WorkflowContext context) {
-        log.info("Create sort config V2 for groupId={}", 
context.getProcessForm().getInlongGroupId());
-        GroupResourceProcessForm form = (GroupResourceProcessForm) 
context.getProcessForm();
-        GroupOperateType groupOperateType = form.getGroupOperateType();
-        if (groupOperateType == GroupOperateType.SUSPEND || groupOperateType 
== GroupOperateType.DELETE) {
-            return ListenerResult.success();
-        }
-        InlongGroupInfo groupInfo = form.getGroupInfo();
-        List<InlongStreamInfo> streamInfos = form.getStreamInfos();
-        int sinkCount = streamInfos.stream()
-                .map(s -> s.getSinkList() == null ? 0 : s.getSinkList().size())
-                .reduce(0, Integer::sum);
-        if (sinkCount == 0) {
-            log.warn("not any sink for group {} found, skip creating sort 
config", groupInfo.getInlongGroupId());
-            return ListenerResult.success();
+    public void buildConfig(InlongGroupInfo groupInfo, List<InlongStreamInfo> 
streamInfos, boolean isStream)
+            throws Exception {
+        if (groupInfo == null || CollectionUtils.isEmpty(streamInfos)) {
+            LOGGER.warn("group info is null or stream infos is empty, no need 
to build sort config for disable zk");
+            return;
         }
 
-        GroupInfo configInfo = createGroupInfo(groupInfo, streamInfos);
-        String dataFlows = OBJECT_MAPPER.writeValueAsString(configInfo);
-        addExtInfo(groupInfo, InlongConstants.DATA_FLOW, dataFlows);
-        return ListenerResult.success();
-    }
+        GroupInfo configInfo = this.createSortGroupInfo(groupInfo, 
streamInfos);
+        String dataflow = OBJECT_MAPPER.writeValueAsString(configInfo);
 
-    private void addExtInfo(InlongGroupInfo groupInfo, String key, String 
value) {
-        if (groupInfo.getExtList() == null) {
-            groupInfo.setExtList(Lists.newArrayList());
+        if (isStream) {
+            this.addToStreamExt(streamInfos, dataflow);
+        } else {
+            this.addToGroupExt(groupInfo, dataflow);
         }
-        InlongGroupExtInfo extInfo = new InlongGroupExtInfo();
-        extInfo.setInlongGroupId(groupInfo.getInlongGroupId());
-        extInfo.setKeyName(key);
-        extInfo.setKeyValue(value);
-        upsertExtInfo(groupInfo, extInfo);
     }
 
     /**
-     * TODO need support TubeMQ
+     * Create GroupInfo for Sort protocol.
+     *
+     * @see org.apache.inlong.sort.protocol.GroupInfo
      */
-    private GroupInfo createGroupInfo(InlongGroupInfo groupInfo, 
List<InlongStreamInfo> streamInfoList) {
+    private GroupInfo createSortGroupInfo(InlongGroupInfo groupInfo, 
List<InlongStreamInfo> streamInfoList) {
         String groupId = groupInfo.getInlongGroupId();
         List<StreamSink> streamSinks = sinkService.listSink(groupId, null);
         Map<String, List<StreamSink>> sinkMap = streamSinks.stream()
                 .collect(Collectors.groupingBy(StreamSink::getInlongStreamId, 
HashMap::new,
                         Collectors.toCollection(ArrayList::new)));
-        Map<String, List<StreamSource>> sourceMap = 
createPulsarSources(groupInfo, streamInfoList);
 
-        List<StreamInfo> streamInfos = new ArrayList<>();
+        // get source info
+        Map<String, List<StreamSource>> sourceMap;
+        if (MQType.MQ_PULSAR.equals(groupInfo.getMqType())) {
+            sourceMap = this.createPulsarSources(groupInfo, streamInfoList);
+        } else {
+            // TODO need to support TubeMQ
+            String errMsg = String.format("Unsupported MQ type: %s", 
groupInfo.getMqType());
+            LOGGER.error(errMsg);
+            throw new WorkflowListenerException(errMsg);
+        }
+
+        // create StreamInfo for Sort protocol
+        List<StreamInfo> sortStreamInfos = new ArrayList<>();
         for (InlongStreamInfo inlongStream : streamInfoList) {
             String streamId = inlongStream.getInlongStreamId();
-            StreamInfo streamInfo = new StreamInfo(streamId,
-                    createNodesForStream(sourceMap.get(streamId), 
sinkMap.get(streamId)),
-                    createNodeRelationsForStream(sourceMap.get(streamId), 
sinkMap.get(streamId)));
-            streamInfos.add(streamInfo);
+            List<StreamSource> sources = sourceMap.get(streamId);
+            List<StreamSink> sinks = sinkMap.get(streamId);
+            StreamInfo sortStream = new StreamInfo(streamId,
+                    this.createNodesForStream(sources, sinks),
+                    this.createNodeRelationsForStream(sources, sinks));
+            sortStreamInfos.add(sortStream);
         }
 
-        return new GroupInfo(groupId, streamInfos);
+        return new GroupInfo(groupId, sortStreamInfos);
     }
 
+    /**
+     * Create Pulsar sources for Sort.
+     */
     private Map<String, List<StreamSource>> createPulsarSources(
             InlongGroupInfo groupInfo, List<InlongStreamInfo> streamInfoList) {
-
-        if (!MQType.MQ_PULSAR.equals(groupInfo.getMqType())) {
-            String errMsg = String.format("Unsupported MQ type %s", 
groupInfo.getMqType());
-            log.error(errMsg);
-            throw new WorkflowListenerException(errMsg);
-        }
-
-        Map<String, List<StreamSource>> sourceMap = Maps.newHashMap();
         ClusterInfo clusterInfo = 
clusterService.getOne(groupInfo.getInlongClusterTag(), null,
                 ClusterType.PULSAR);
-
         PulsarClusterInfo pulsarCluster = (PulsarClusterInfo) clusterInfo;
         String adminUrl = pulsarCluster.getAdminUrl();
         String serviceUrl = pulsarCluster.getUrl();
-        String tenant = StringUtils.isEmpty(pulsarCluster.getTenant()) ? 
InlongConstants.DEFAULT_PULSAR_TENANT
-                : pulsarCluster.getTenant();
+        String tenant = StringUtils.isEmpty(pulsarCluster.getTenant())
+                ? InlongConstants.DEFAULT_PULSAR_TENANT : 
pulsarCluster.getTenant();
+
+        Map<String, List<StreamSource>> sourceMap = Maps.newHashMap();
         streamInfoList.forEach(streamInfo -> {
             PulsarSource pulsarSource = new PulsarSource();
             String streamId = streamInfo.getInlongStreamId();
-            pulsarSource.setTenant(tenant);
             pulsarSource.setSourceName(streamId);
+            pulsarSource.setTenant(tenant);
             pulsarSource.setNamespace(groupInfo.getMqResource());
             pulsarSource.setTopic(streamInfo.getMqResource());
             pulsarSource.setAdminUrl(adminUrl);
@@ -179,6 +172,8 @@ public class CreateSortConfigListenerV2 implements 
SortOperateListener {
                     pulsarSource.setPrimaryKey(((KafkaSource) 
sourceInfo).getPrimaryKey());
                 }
             }
+
+            // if the SerializationType is still null, set it to the CSV
             if (StringUtils.isEmpty(pulsarSource.getSerializationType())) {
                 pulsarSource.setSerializationType(DataTypeEnum.CSV.getName());
             }
@@ -186,6 +181,7 @@ public class CreateSortConfigListenerV2 implements 
SortOperateListener {
             pulsarSource.setFieldList(streamInfo.getFieldList());
             sourceMap.computeIfAbsent(streamId, key -> 
Lists.newArrayList()).add(pulsarSource);
         });
+
         return sourceMap;
     }
 
@@ -205,9 +201,41 @@ public class CreateSortConfigListenerV2 implements 
SortOperateListener {
         return Lists.newArrayList(relation);
     }
 
-    private void upsertExtInfo(InlongGroupInfo groupInfo, InlongGroupExtInfo 
extInfo) {
+    /**
+     * Add config into inlong group ext info
+     */
+    private void addToGroupExt(InlongGroupInfo groupInfo, String value) {
+        if (groupInfo.getExtList() == null) {
+            groupInfo.setExtList(Lists.newArrayList());
+        }
+
+        InlongGroupExtInfo extInfo = new InlongGroupExtInfo();
+        extInfo.setInlongGroupId(groupInfo.getInlongGroupId());
+        extInfo.setKeyName(InlongConstants.DATAFLOW);
+        extInfo.setKeyValue(value);
+
         groupInfo.getExtList().removeIf(ext -> 
extInfo.getKeyName().equals(ext.getKeyName()));
         groupInfo.getExtList().add(extInfo);
     }
 
+    /**
+     * Add config into inlong stream ext info
+     */
+    private void addToStreamExt(List<InlongStreamInfo> streamInfos, String 
value) {
+        streamInfos.forEach(streamInfo -> {
+            if (streamInfo.getExtList() == null) {
+                streamInfo.setExtList(Lists.newArrayList());
+            }
+
+            InlongStreamExtInfo extInfo = new InlongStreamExtInfo();
+            extInfo.setInlongGroupId(streamInfo.getInlongGroupId());
+            extInfo.setInlongStreamId(streamInfo.getInlongStreamId());
+            extInfo.setKeyName(InlongConstants.DATAFLOW);
+            extInfo.setKeyValue(value);
+
+            streamInfo.getExtList().removeIf(ext -> 
extInfo.getKeyName().equals(ext.getKeyName()));
+            streamInfo.getExtList().add(extInfo);
+        });
+    }
+
 }
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigListener.java
new file mode 100644
index 000000000..ffaeec151
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigListener.java
@@ -0,0 +1,88 @@
+/*
+ * 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.sort;
+
+import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
+import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
+import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
+import 
org.apache.inlong.manager.common.pojo.workflow.form.process.GroupResourceProcessForm;
+import org.apache.inlong.manager.workflow.WorkflowContext;
+import org.apache.inlong.manager.workflow.event.ListenerResult;
+import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
+import org.apache.inlong.manager.workflow.event.task.TaskEvent;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import java.util.List;
+
+/**
+ * Event listener of build the Sort config,
+ * such as update the form config, or build and push config to ZK, etc.
+ */
+@Component
+public class SortConfigListener implements SortOperateListener {
+
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(SortConfigListener.class);
+
+    @Autowired
+    private SortConfigOperatorFactory operatorFactory;
+
+    @Override
+    public TaskEvent event() {
+        return TaskEvent.COMPLETE;
+    }
+
+    @Override
+    public ListenerResult listen(WorkflowContext context) throws 
WorkflowListenerException {
+        GroupResourceProcessForm form = (GroupResourceProcessForm) 
context.getProcessForm();
+        String groupId = form.getInlongGroupId();
+        LOGGER.info("begin to build sort config for groupId={}", groupId);
+
+        GroupOperateType operateType = form.getGroupOperateType();
+        if (operateType == GroupOperateType.SUSPEND || operateType == 
GroupOperateType.DELETE) {
+            LOGGER.info("not build sort config for groupId={}, as the group 
operate type={}", groupId, operateType);
+            return ListenerResult.success();
+        }
+        InlongGroupInfo groupInfo = form.getGroupInfo();
+        List<InlongStreamInfo> streamInfos = form.getStreamInfos();
+        int sinkCount = streamInfos.stream()
+                .map(stream -> stream.getSinkList() == null ? 0 : 
stream.getSinkList().size())
+                .reduce(0, Integer::sum);
+        if (sinkCount == 0) {
+            LOGGER.warn("not build sort config for groupId={}, as not found 
any sink", groupId);
+            return ListenerResult.success();
+        }
+
+        try {
+            SortConfigOperator operator = 
operatorFactory.getInstance(groupInfo.getLightweight(),
+                    groupInfo.getEnableZookeeper());
+            operator.buildConfig(groupInfo, streamInfos, false);
+        } catch (Exception e) {
+            String msg = String.format("failed to build sort config for 
groupId=%s, ", groupId);
+            LOGGER.error(msg + "streamInfos=" + streamInfos, e);
+            throw new WorkflowListenerException(msg + e.getMessage());
+        }
+
+        LOGGER.info("success to build sort config for groupId={}", groupId);
+        return ListenerResult.success();
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperator.java
new file mode 100644
index 000000000..e6c35253f
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperator.java
@@ -0,0 +1,47 @@
+/*
+ * 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.sort;
+
+import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
+import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
+
+import java.util.List;
+
+/**
+ * Interface of the Sort config operator
+ */
+public interface SortConfigOperator {
+
+    /**
+     * Determines whether the current instance matches the specified type.
+     *
+     * @param isNormal is the inlong group is normal mode, 0: normal mode, 1: 
lightweight mode
+     * @param enableZk is the inlong group enable the ZooKeeper, 1: enable, 0: 
disable
+     */
+    Boolean accept(Integer isNormal, Integer enableZk);
+
+    /**
+     * Build Sort config.
+     *
+     * @param groupInfo inlong group info
+     * @param streamInfos inlong stream info list
+     * @param isStream is the config built for inlong stream
+     */
+    void buildConfig(InlongGroupInfo groupInfo, List<InlongStreamInfo> 
streamInfos, boolean isStream) throws Exception;
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperatorFactory.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperatorFactory.java
new file mode 100644
index 000000000..4c6c547b4
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/SortConfigOperatorFactory.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.sort;
+
+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 SortConfigOperator}.
+ */
+@Service
+public class SortConfigOperatorFactory {
+
+    @Autowired
+    private List<SortConfigOperator> operatorList;
+
+    /**
+     * Get a Sort config operator instance.
+     *
+     * @param isNormal is the inlong group is normal mode, 0: normal mode, 1: 
lightweight mode
+     * @param enableZk is the inlong group enable the ZooKeeper, 1: enable, 0: 
disable
+     */
+    public SortConfigOperator getInstance(Integer isNormal, Integer enableZk) {
+        return operatorList.stream()
+                .filter(inst -> inst.accept(isNormal, enableZk))
+                .findFirst()
+                .orElseThrow(() -> new BusinessException("not found any 
instance of SortConfigOperator when isNormal="
+                        + isNormal + ", enableZk=" + enableZk));
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/StreamSortConfigListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/StreamSortConfigListener.java
new file mode 100644
index 000000000..17ddacf4d
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/StreamSortConfigListener.java
@@ -0,0 +1,95 @@
+/*
+ * 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.sort;
+
+import org.apache.commons.collections.CollectionUtils;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
+import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
+import org.apache.inlong.manager.common.pojo.sink.StreamSink;
+import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
+import 
org.apache.inlong.manager.common.pojo.workflow.form.process.StreamResourceProcessForm;
+import org.apache.inlong.manager.workflow.WorkflowContext;
+import org.apache.inlong.manager.workflow.event.ListenerResult;
+import org.apache.inlong.manager.workflow.event.task.SortOperateListener;
+import org.apache.inlong.manager.workflow.event.task.TaskEvent;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import java.util.Collections;
+import java.util.List;
+
+/**
+ * Event listener of build the Sort config for one inlong stream,
+ * such as update the form config, or build and push config to ZK, etc.
+ */
+@Component
+public class StreamSortConfigListener implements SortOperateListener {
+
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(StreamSortConfigListener.class);
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); // 
thread safe
+
+    @Autowired
+    private SortConfigOperatorFactory operatorFactory;
+
+    @Override
+    public TaskEvent event() {
+        return TaskEvent.COMPLETE;
+    }
+
+    @Override
+    public ListenerResult listen(WorkflowContext context) throws 
WorkflowListenerException {
+        StreamResourceProcessForm form = (StreamResourceProcessForm) 
context.getProcessForm();
+        InlongStreamInfo streamInfo = form.getStreamInfo();
+        final String groupId = streamInfo.getInlongGroupId();
+        final String streamId = streamInfo.getInlongStreamId();
+        LOGGER.info("begin to build sort config for groupId={}, streamId={}", 
groupId, streamId);
+
+        GroupOperateType operateType = form.getGroupOperateType();
+        if (operateType == GroupOperateType.SUSPEND || operateType == 
GroupOperateType.DELETE) {
+            LOGGER.info("not build sort config for groupId={}, streamId={}, as 
the group operate type={}",
+                    groupId, streamId, operateType);
+            return ListenerResult.success();
+        }
+
+        InlongGroupInfo groupInfo = form.getGroupInfo();
+        List<StreamSink> streamSinks = streamInfo.getSinkList();
+        if (CollectionUtils.isEmpty(streamSinks)) {
+            LOGGER.warn("not build sort config for groupId={}, streamId={}, as 
not found any sinks", groupId, streamId);
+            return ListenerResult.success();
+        }
+
+        List<InlongStreamInfo> streamInfos = 
Collections.singletonList(streamInfo);
+        try {
+            SortConfigOperator operator = 
operatorFactory.getInstance(groupInfo.getLightweight(),
+                    groupInfo.getEnableZookeeper());
+            operator.buildConfig(groupInfo, streamInfos, true);
+        } catch (Exception e) {
+            String msg = String.format("failed to build sort config for 
groupId=%s, streamId=%s, ", groupId, streamId);
+            LOGGER.error(msg + "streamInfos=" + streamInfos, e);
+            throw new WorkflowListenerException(msg + e.getMessage());
+        }
+
+        LOGGER.info("success to build sort config for groupId={}, 
streamId={}", groupId, streamId);
+        return ListenerResult.success();
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperDisabledSelector.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperDisabledSelector.java
index b30b39084..cac93aa5f 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperDisabledSelector.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperDisabledSelector.java
@@ -18,6 +18,7 @@
 package org.apache.inlong.manager.service.sort;
 
 import lombok.extern.slf4j.Slf4j;
+import org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.enums.MQType;
 import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
 import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
@@ -40,7 +41,7 @@ public class ZookeeperDisabledSelector implements 
EventSelector {
         if (processForm instanceof GroupResourceProcessForm) {
             GroupResourceProcessForm groupResourceForm = 
(GroupResourceProcessForm) processForm;
             InlongGroupInfo groupInfo = groupResourceForm.getGroupInfo();
-            boolean enable = groupInfo.getEnableZookeeper() == 0
+            boolean enable = 
InlongConstants.DISABLE_ZK.equals(groupInfo.getEnableZookeeper())
                     && MQType.forType(groupInfo.getMqType()) != MQType.NONE;
 
             log.info("zookeeper disabled was [{}] for groupId [{}]", enable, 
groupId);
@@ -49,7 +50,7 @@ public class ZookeeperDisabledSelector implements 
EventSelector {
             StreamResourceProcessForm streamResourceForm = 
(StreamResourceProcessForm) processForm;
             InlongGroupInfo groupInfo = streamResourceForm.getGroupInfo();
             InlongStreamInfo streamInfo = streamResourceForm.getStreamInfo();
-            boolean enable = groupInfo.getEnableZookeeper() == 0
+            boolean enable = 
InlongConstants.DISABLE_ZK.equals(groupInfo.getEnableZookeeper())
                     && MQType.forType(groupInfo.getMqType()) != MQType.NONE;
             log.info("zookeeper disabled was [{}] for groupId [{}] and 
streamId [{}] ", enable, groupId,
                     streamInfo.getInlongStreamId());
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperEnabledSelector.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperEnabledSelector.java
index 59904b7fd..2c3d21952 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperEnabledSelector.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/ZookeeperEnabledSelector.java
@@ -18,6 +18,7 @@
 package org.apache.inlong.manager.service.sort;
 
 import lombok.extern.slf4j.Slf4j;
+import org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.enums.MQType;
 import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
 import 
org.apache.inlong.manager.common.pojo.workflow.form.process.GroupResourceProcessForm;
@@ -43,8 +44,8 @@ public class ZookeeperEnabledSelector implements 
EventSelector {
 
         GroupResourceProcessForm groupResourceForm = 
(GroupResourceProcessForm) processForm;
         InlongGroupInfo groupInfo = groupResourceForm.getGroupInfo();
-        boolean enable =
-                groupInfo.getEnableZookeeper() == 1 && 
MQType.forType(groupInfo.getMqType()) != MQType.NONE;
+        boolean enable = 
InlongConstants.ENABLE_ZK.equals(groupInfo.getEnableZookeeper())
+                && MQType.forType(groupInfo.getMqType()) != MQType.NONE;
         log.info("zookeeper enabled was [{}] for groupId [{}]", enable, 
groupId);
         return enable;
     }
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/light/LightGroupSortListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/light/LightGroupSortListener.java
index 460ffd58e..90f116821 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/light/LightGroupSortListener.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/light/LightGroupSortListener.java
@@ -83,11 +83,11 @@ public class LightGroupSortListener implements 
SortOperateListener {
             List<InlongStreamInfo> streamInfos = processForm.getStreamInfos();
             final String groupId = groupInfo.getInlongGroupId();
             GroupInfo configInfo = this.createGroupInfo(groupInfo, 
streamInfos);
-            String dataFlows = OBJECT_MAPPER.writeValueAsString(configInfo);
+            String dataflow = OBJECT_MAPPER.writeValueAsString(configInfo);
             InlongGroupExtInfo extInfo = new InlongGroupExtInfo();
             extInfo.setInlongGroupId(groupId);
-            extInfo.setKeyName(InlongConstants.DATA_FLOW);
-            extInfo.setKeyValue(dataFlows);
+            extInfo.setKeyName(InlongConstants.DATAFLOW);
+            extInfo.setKeyValue(dataflow);
             if (groupInfo.getExtList() == null) {
                 groupInfo.setExtList(Lists.newArrayList());
             }
@@ -145,7 +145,7 @@ public class LightGroupSortListener implements 
SortOperateListener {
     }
 
     private void upsertExtInfo(InlongGroupInfo groupInfo, InlongGroupExtInfo 
extInfo) {
-        groupInfo.getExtList().removeIf(ext -> 
InlongConstants.DATA_FLOW.equals(ext.getKeyName()));
+        groupInfo.getExtList().removeIf(ext -> 
InlongConstants.DATAFLOW.equals(ext.getKeyName()));
         groupInfo.getExtList().add(extInfo);
     }
 
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/DataFlowUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/DataFlowUtils.java
deleted file mode 100644
index 6764219b9..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/DataFlowUtils.java
+++ /dev/null
@@ -1,94 +0,0 @@
-/*
- * 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.sort.util;
-
-import org.apache.inlong.manager.service.core.InlongStreamService;
-import org.apache.inlong.manager.service.group.GroupCheckService;
-import org.apache.inlong.manager.service.source.StreamSourceService;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Service;
-
-/**
- * Util for build data flow info.
- */
-@Service
-@Deprecated
-public class DataFlowUtils {
-
-    @Autowired
-    private GroupCheckService groupCheckService;
-    @Autowired
-    private StreamSourceService streamSourceService;
-    @Autowired
-    private InlongStreamService streamService;
-
-    /*
-     * Create dataflow info for sort.
-     */
-    /*public DataFlowInfo createDataFlow(InlongGroupInfo groupInfo, StreamSink 
streamSink) {
-        String groupId = streamSink.getInlongGroupId();
-        String streamId = streamSink.getInlongStreamId();
-        List<StreamSource> sourceList = 
streamSourceService.listSource(groupId, streamId);
-        if (CollectionUtils.isEmpty(sourceList)) {
-            throw new WorkflowListenerException(String.format("Source not 
found by groupId=%s and streamId=%s",
-                    groupId, streamId));
-        }
-
-        // Get all field info
-        List<FieldInfo> sourceFields = new ArrayList<>();
-        List<FieldInfo> sinkFields = new ArrayList<>();
-
-        // TODO Support more than one source and one sink
-        final StreamSource streamSource = sourceList.get(0);
-        boolean isAllMigration = 
SourceInfoUtils.isBinlogAllMigration(streamSource);
-
-        List<FieldMappingUnit> mappingUnitList;
-        InlongStreamInfo streamInfo = streamService.get(groupId, streamId);
-        if (isAllMigration) {
-            mappingUnitList = 
FieldInfoUtils.setAllMigrationFieldMapping(sourceFields, sinkFields);
-        } else {
-            mappingUnitList = 
FieldInfoUtils.createFieldInfo(streamInfo.getFieldList(),
-                    streamSink.getFieldList(), sourceFields, sinkFields);
-        }
-
-        FieldMappingRule fieldMappingRule = new 
FieldMappingRule(mappingUnitList.toArray(new FieldMappingUnit[0]));
-
-        // Get source info
-        String masterAddress = 
commonOperateService.getSpecifiedParam(InlongConstants.TUBE_MASTER_URL);
-        PulsarClusterInfo pulsarCluster = 
commonOperateService.getPulsarClusterInfo(groupInfo.getMqType());
-        org.apache.inlong.sort.protocol.source.SourceInfo sourceInfo = 
SourceInfoUtils.createSourceInfo(pulsarCluster,
-                masterAddress, clusterBean,
-                groupInfo, streamInfo, streamSource, sourceFields);
-
-        // Get sink info
-        SinkInfo sinkInfo = SinkInfoUtils.createSinkInfo(streamSource, 
streamSink, sinkFields);
-
-        // Get transformation info
-        TransformationInfo transInfo = new 
TransformationInfo(fieldMappingRule);
-
-        // Get properties
-        Map<String, Object> properties = new HashMap<>();
-        if (MapUtils.isNotEmpty(streamSink.getProperties())) {
-            properties.putAll(streamSink.getProperties());
-        }
-        properties.put(InlongConstants.DATA_FLOW_GROUP_ID_KEY, groupId);
-
-        return new DataFlowInfo(streamSink.getId(), sourceInfo, transInfo, 
sinkInfo, properties);
-    }*/
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
deleted file mode 100644
index 3631b0217..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
+++ /dev/null
@@ -1,246 +0,0 @@
-/*
- * 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.sort.util;
-
-/**
- * Utils for create sink info, such as kafka sink, clickhouse sink, etc.
- */
-public class SinkInfoUtils {
-
-    private static final String DATA_FORMAT = "yyyyMMddHH";
-    private static final String TIME_FORMAT = "HHmmss";
-    private static final String DATA_TIME_FORMAT = "yyyyMMddHHmmss";
-
-    /*
-     * Create sink info for DataFlowInfo.
-     */
-    /*public static SinkInfo createSinkInfo(StreamSource streamSource, 
StreamSink streamSink,
-            List<FieldInfo> sinkFields) {
-        String sinkType = streamSink.getSinkType();
-        SinkInfo sinkInfo;
-        switch (SinkType.forType(sinkType)) {
-            case HIVE:
-                sinkInfo = createHiveSinkInfo((HiveSink) streamSink, 
sinkFields);
-                break;
-            case KAFKA:
-                sinkInfo = createKafkaSinkInfo(streamSource, (KafkaSink) 
streamSink, sinkFields);
-                break;
-            case ICEBERG:
-                sinkInfo = createIcebergSinkInfo((IcebergSink) streamSink, 
sinkFields);
-                break;
-            case CLICKHOUSE:
-                sinkInfo = createClickhouseSinkInfo((ClickHouseSink) 
streamSink, sinkFields);
-                break;
-            case HBASE:
-                sinkInfo = createHbaseSinkInfo((HBaseSink) streamSink, 
sinkFields);
-                break;
-            case ELASTICSEARCH:
-                sinkInfo = createEsSinkInfo((ElasticsearchSink) streamSink, 
sinkFields);
-                break;
-            default:
-                throw new BusinessException(String.format("Unsupported 
SinkType {%s}", sinkType));
-        }
-        return sinkInfo;
-    }*/
-
-    /*private static ClickHouseSinkInfo 
createClickhouseSinkInfo(ClickHouseSink ckSink, List<FieldInfo> sinkFields) {
-        if (StringUtils.isEmpty(ckSink.getJdbcUrl())) {
-            throw new BusinessException(String.format("ClickHouse={%s} jdbc 
url cannot be empty", ckSink));
-        } else if (CollectionUtils.isEmpty(ckSink.getFieldList())) {
-            throw new BusinessException(String.format("ClickHouse={%s} fields 
cannot be empty", ckSink));
-        } else if (StringUtils.isEmpty(ckSink.getTableName())) {
-            throw new BusinessException(String.format("ClickHouse={%s} table 
name cannot be empty", ckSink));
-        } else if (StringUtils.isEmpty(ckSink.getDbName())) {
-            throw new BusinessException(String.format("ClickHouse={%s} 
database name cannot be empty", ckSink));
-        }
-
-        Integer isDistributed = ckSink.getIsDistributed();
-        if (isDistributed == null) {
-            throw new BusinessException(String.format("ClickHouse={%s} 
isDistributed cannot be null", ckSink));
-        }
-
-        // Default partition strategy is RANDOM
-        ClickHouseSinkInfo.PartitionStrategy partitionStrategy = 
PartitionStrategy.RANDOM;
-        boolean distributedTable = isDistributed == 1;
-        if (distributedTable) {
-            if 
(PartitionStrategy.BALANCE.name().equalsIgnoreCase(ckSink.getPartitionStrategy()))
 {
-                partitionStrategy = PartitionStrategy.BALANCE;
-            } else if 
(PartitionStrategy.HASH.name().equalsIgnoreCase(ckSink.getPartitionStrategy())) 
{
-                partitionStrategy = PartitionStrategy.HASH;
-            }
-        }
-
-        // TODO Add keyFieldNames instead of `new String[0]`
-        return new ClickHouseSinkInfo(ckSink.getJdbcUrl(), ckSink.getDbName(),
-                ckSink.getTableName(), ckSink.getUsername(), 
ckSink.getPassword(),
-                distributedTable, partitionStrategy, 
ckSink.getPartitionFields(),
-                sinkFields.toArray(new FieldInfo[0]), new String[0],
-                ckSink.getFlushInterval(), ckSink.getFlushRecord(),
-                ckSink.getRetryTimes());
-    }*/
-
-    /*// TODO Need set more configs for IcebergSinkInfo
-    private static IcebergSinkInfo createIcebergSinkInfo(IcebergSink 
icebergSink, List<FieldInfo> sinkFields) {
-        if (StringUtils.isEmpty(icebergSink.getDataPath())) {
-            throw new BusinessException(String.format("Iceberg={%s} data path 
cannot be empty", icebergSink));
-        }
-
-        return new IcebergSinkInfo(sinkFields.toArray(new FieldInfo[0]), 
icebergSink.getDataPath());
-    }*/
-
-    /*private static KafkaSinkInfo createKafkaSinkInfo(StreamSource 
streamSource, KafkaSink kafkaSink,
-            List<FieldInfo> sinkFields) {
-        String addressUrl = kafkaSink.getBootstrapServers();
-        String topicName = kafkaSink.getTopicName();
-        SerializationInfo serializationInfo = 
SerializationUtils.createSerialInfo(streamSource, kafkaSink);
-        return new KafkaSinkInfo(sinkFields.toArray(new FieldInfo[0]), 
addressUrl, topicName, serializationInfo);
-    }*/
-
-    /*
-     * Create Hive sink info.
-     */
-    /*private static HiveSinkInfo createHiveSinkInfo(HiveSink hiveInfo, 
List<FieldInfo> sinkFields) {
-        if (hiveInfo.getJdbcUrl() == null) {
-            throw new BusinessException(String.format("HiveSink={%s} server 
url cannot be empty", hiveInfo));
-        }
-        if (CollectionUtils.isEmpty(hiveInfo.getFieldList())) {
-            throw new BusinessException(String.format("HiveSink={%s} fields 
cannot be empty", hiveInfo));
-        }
-        // Use the field separator in Hive, the default is TextFile
-        Character separator = (char) 
Integer.parseInt(hiveInfo.getDataSeparator());
-        HiveFileFormat fileFormat;
-        FileFormat format = FileFormat.forName(hiveInfo.getFileFormat());
-
-        if (format == FileFormat.ORCFile) {
-            fileFormat = new HiveSinkInfo.OrcFileFormat(1000);
-        } else if (format == FileFormat.SequenceFile) {
-            fileFormat = new HiveSinkInfo.SequenceFileFormat(separator, 100);
-        } else if (format == FileFormat.Parquet) {
-            fileFormat = new HiveSinkInfo.ParquetFileFormat();
-        } else {
-            fileFormat = new HiveSinkInfo.TextFileFormat(separator);
-        }
-
-        // Handle hive partition list
-        List<HivePartitionInfo> partitionList = new ArrayList<>();
-        List<HivePartitionField> partitionFieldList = 
hiveInfo.getPartitionFieldList();
-        if (CollectionUtils.isNotEmpty(partitionFieldList)) {
-            SinkInfoUtils.checkPartitionField(hiveInfo.getFieldList(), 
partitionFieldList);
-            partitionList = partitionFieldList.stream().map(s -> {
-                HivePartitionInfo partition;
-                String fieldFormat = s.getFieldFormat();
-                switch (FieldType.forName(s.getFieldType())) {
-                    case TIMESTAMP:
-                        fieldFormat = StringUtils.isNotBlank(fieldFormat) ? 
fieldFormat : DATA_TIME_FORMAT;
-                        partition = new 
HiveTimePartitionInfo(s.getFieldName(), fieldFormat);
-                        break;
-                    case DATE:
-                        fieldFormat = StringUtils.isNotBlank(fieldFormat) ? 
fieldFormat : DATA_FORMAT;
-                        partition = new 
HiveTimePartitionInfo(s.getFieldName(), fieldFormat);
-                        break;
-                    default:
-                        partition = new 
HiveFieldPartitionInfo(s.getFieldName());
-                }
-                return partition;
-            }).collect(Collectors.toList());
-        }
-
-        // dataPath = dataPath + / + tableName
-        StringBuilder dataPathBuilder = new StringBuilder();
-        String dataPath = hiveInfo.getDataPath();
-        if (!dataPath.endsWith("/")) {
-            dataPathBuilder.append(dataPath).append("/");
-        }
-        dataPath = dataPathBuilder.append(hiveInfo.getTableName()).toString();
-
-        return new HiveSinkInfo(sinkFields.toArray(new FieldInfo[0]), 
hiveInfo.getJdbcUrl(),
-                hiveInfo.getDbName(), hiveInfo.getTableName(), 
hiveInfo.getUsername(), hiveInfo.getPassword(),
-                dataPath, partitionList.toArray(new 
HiveSinkInfo.HivePartitionInfo[0]), fileFormat);
-    }*/
-
-    /*
-     * Check the validation of Hive partition field.
-     */
-    /*public static void checkPartitionField(List<SinkField> fieldList, 
List<HivePartitionField> partitionList) {
-        if (CollectionUtils.isEmpty(partitionList)) {
-            return;
-        }
-
-        if (CollectionUtils.isEmpty(fieldList)) {
-            throw new 
BusinessException(ErrorCodeEnum.SINK_FIELD_LIST_IS_EMPTY);
-        }
-
-        Map<String, SinkField> sinkFieldMap = new HashMap<>(fieldList.size());
-        fieldList.forEach(field -> sinkFieldMap.put(field.getFieldName(), 
field));
-
-        for (HivePartitionField partitionField : partitionList) {
-            String fieldName = partitionField.getFieldName();
-            if (StringUtils.isBlank(fieldName)) {
-                throw new 
BusinessException(ErrorCodeEnum.PARTITION_FIELD_NAME_IS_EMPTY);
-            }
-
-            SinkField sinkField = sinkFieldMap.get(fieldName);
-            if (sinkField == null) {
-                throw new BusinessException(
-                        
String.format(ErrorCodeEnum.PARTITION_FIELD_NOT_FOUND.getMessage(), fieldName));
-            }
-
-            if (StringUtils.isBlank(sinkField.getSourceFieldName())) {
-                throw new BusinessException(
-                        
String.format(ErrorCodeEnum.PARTITION_FIELD_NO_SOURCE_FIELD.getMessage(), 
fieldName));
-            }
-        }
-    }*/
-
-    /*
-     * Creat HBase sink info.
-     */
-    /*private static HbaseSinkInfo createHbaseSinkInfo(HBaseSink hbaseSink, 
List<FieldInfo> sinkFields) {
-        if (StringUtils.isEmpty(hbaseSink.getZkQuorum())) {
-            throw new BusinessException(String.format("HBase={%s} zookeeper 
quorum url cannot be empty", hbaseSink));
-        } else if (StringUtils.isEmpty(hbaseSink.getZkNodeParent())) {
-            throw new BusinessException(String.format("HBase={%s} zookeeper 
node cannot be empty", hbaseSink));
-        } else if (StringUtils.isEmpty(hbaseSink.getTableName())) {
-            throw new BusinessException(String.format("HBase={%s} table name 
cannot be empty", hbaseSink));
-        }
-
-        return new HbaseSinkInfo(sinkFields.toArray(new FieldInfo[0]), 
hbaseSink.getZkQuorum(),
-                hbaseSink.getZkNodeParent(), hbaseSink.getNamespace(), 
hbaseSink.getTableName(),
-                hbaseSink.getBufferFlushMaxSize(), 
hbaseSink.getBufferFlushMaxSize(),
-                hbaseSink.getBufferFlushInterval());
-
-    }*/
-
-    /*
-     * Creat Elasticsearch sink info.
-     */
-    /*private static ElasticsearchSinkInfo createEsSinkInfo(ElasticsearchSink 
esSink, List<FieldInfo> sinkFields) {
-        if (StringUtils.isEmpty(esSink.getHost())) {
-            throw new BusinessException(String.format("es={%s} host cannot be 
empty", esSink));
-        } else if (StringUtils.isEmpty(esSink.getIndexName())) {
-            throw new BusinessException(String.format("es={%s} indexName 
cannot be empty", esSink));
-        }
-
-        return new ElasticsearchSinkInfo(esSink.getHost(), esSink.getPort(),
-                esSink.getIndexName(), esSink.getUsername(), 
esSink.getPassword(),
-                sinkFields.toArray(new FieldInfo[0]), new String[0],
-                esSink.getFlushInterval(), esSink.getFlushRecord(),
-                esSink.getRetryTimes());
-    }*/
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SourceInfoUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SourceInfoUtils.java
deleted file mode 100644
index 53652c218..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SourceInfoUtils.java
+++ /dev/null
@@ -1,112 +0,0 @@
-/*
- * 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.sort.util;
-
-import org.apache.inlong.manager.common.enums.SourceType;
-import org.apache.inlong.manager.common.pojo.source.StreamSource;
-import org.apache.inlong.manager.common.pojo.source.mysql.MySQLBinlogSource;
-
-/**
- * Utils for creat source info, such as pulsar source, tube MQ source.
- */
-@Deprecated
-public class SourceInfoUtils {
-
-    /**
-     * Whether the source is all binlog migration.
-     */
-    public static boolean isBinlogAllMigration(StreamSource sourceInfo) {
-        if (sourceInfo == null) {
-            return false;
-        }
-        if 
(SourceType.BINLOG.getType().equalsIgnoreCase(sourceInfo.getSourceType())) {
-            MySQLBinlogSource binlogSource = (MySQLBinlogSource) sourceInfo;
-            return binlogSource.isAllMigration();
-        }
-        return false;
-    }
-
-    /*
-     * Create source info for DataFlowInfo.
-     */
-    /*public static org.apache.inlong.sort.protocol.source.SourceInfo 
createSourceInfo(PulsarClusterInfo pulsarCluster,
-            String masterAddress,
-            ClusterBean clusterBean, InlongGroupInfo groupInfo, 
InlongStreamInfo streamInfo,
-            StreamSource streamSource, List<FieldInfo> sourceFields) {
-
-        MQType mqType = MQType.forType(groupInfo.getMqType());
-        DeserializationInfo deserializationInfo = 
SerializationUtils.createDeserialInfo(streamSource, streamInfo);
-        org.apache.inlong.sort.protocol.source.SourceInfo sourceInfo;
-        if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
-            sourceInfo = createPulsarSourceInfo(pulsarCluster,
-            clusterBean, groupInfo, streamInfo, deserializationInfo,
-                    sourceFields);
-        } else if (mqType == MQType.TUBE) {
-            // InlongGroupInfo groupInfo, String masterAddress,
-            sourceInfo = createTubeSourceInfo(groupInfo, masterAddress,
-            clusterBean, deserializationInfo, sourceFields);
-        } else {
-            throw new WorkflowListenerException(String.format("Unsupported 
middleware {%s}", mqType));
-        }
-
-        return sourceInfo;
-    }*/
-
-    /*
-     * Create source info for Pulsar
-     */
-    /* private static org.apache.inlong.sort.protocol.source.SourceInfo 
createPulsarSourceInfo(
-            PulsarClusterInfo pulsarCluster, ClusterBean clusterBean,
-            InlongGroupInfo groupInfo, InlongStreamInfo streamInfo,
-            DeserializationInfo deserializationInfo, List<FieldInfo> 
fieldInfos) {
-        String topicName = streamInfo.getMqResource();
-        InlongPulsarInfo pulsarInfo = (InlongPulsarInfo) groupInfo;
-        String tenant = clusterBean.getDefaultTenant();
-        if (StringUtils.isNotEmpty(pulsarInfo.getTenant())) {
-            tenant = pulsarInfo.getTenant();
-        }
-
-        final String namespace = groupInfo.getMqResource();
-        // Full name of topic in Pulsar
-        final String fullTopicName = "persistent://" + tenant + "/" + 
namespace + "/" + topicName;
-        final String consumerGroup = clusterBean.getAppName() + "_" + 
topicName + "_consumer_group";
-        FieldInfo[] fieldInfosArr = fieldInfos.toArray(new FieldInfo[0]);
-
-        String type = pulsarCluster.getType();
-        if (StringUtils.isNotEmpty(type) && MQType.forType(type) == 
MQType.TDMQ_PULSAR) {
-            return new 
TDMQPulsarSourceInfo(pulsarCluster.getBrokerServiceUrl(),
-                    fullTopicName, consumerGroup, pulsarCluster.getToken(), 
deserializationInfo, fieldInfosArr);
-        } else {
-            return new PulsarSourceInfo(pulsarCluster.getAdminUrl(), 
pulsarCluster.getBrokerServiceUrl(),
-                    fullTopicName, consumerGroup, deserializationInfo, 
fieldInfosArr, pulsarCluster.getToken());
-        }
-    }*/
-
-    /*
-     * Create source info TubeMQ
-     */
-    /*private static TubeSourceInfo createTubeSourceInfo(InlongGroupInfo 
groupInfo, String masterAddress,
-            ClusterBean clusterBean, DeserializationInfo deserializationInfo, 
List<FieldInfo> fieldInfos) {
-        Preconditions.checkNotNull(masterAddress, "tube cluster address cannot 
be empty");
-        String topic = groupInfo.getMqResource();
-        String consumerGroup = clusterBean.getAppName() + "_" + topic + 
"_consumer_group";
-        return new TubeSourceInfo(topic, masterAddress, consumerGroup, 
deserializationInfo,
-                fieldInfos.toArray(new FieldInfo[0]));
-    }*/
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/listener/AbstractSourceOperateListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/listener/AbstractSourceOperateListener.java
index 8a8137851..ed8b010cb 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/listener/AbstractSourceOperateListener.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/listener/AbstractSourceOperateListener.java
@@ -160,8 +160,7 @@ public abstract class AbstractSourceOperateListener 
implements DataSourceOperate
             return ((GroupResourceProcessForm) 
processForm).getGroupOperateType();
         } else {
             log.error("illegal process form {} to get inlong group info", 
processForm.getFormName());
-            throw new RuntimeException(String.format("Unsupported ProcessForm 
{%s} in CreateSortConfigListener",
-                    processForm.getFormName()));
+            throw new RuntimeException("Unsupported ProcessForm " + 
processForm.getFormName());
         }
     }
 
@@ -171,8 +170,7 @@ public abstract class AbstractSourceOperateListener 
implements DataSourceOperate
             return groupResourceProcessForm.getGroupInfo();
         } else {
             log.error("illegal process form {} to get inlong group info", 
processForm.getFormName());
-            throw new RuntimeException(String.format("Unsupported ProcessForm 
{%s} in CreateSortConfigListener",
-                    processForm.getFormName()));
+            throw new RuntimeException("Unsupported ProcessForm " + 
processForm.getFormName());
         }
     }
 
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/listener/approve/GroupApproveProcessListener.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/listener/approve/GroupApproveProcessListener.java
index 015cbd876..cc560e436 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/listener/approve/GroupApproveProcessListener.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/listener/approve/GroupApproveProcessListener.java
@@ -71,7 +71,7 @@ public class GroupApproveProcessListener implements 
ProcessEventListener {
             case NORMAL:
                 createGroupResource(context, groupInfo);
                 break;
-            case LIGHT:
+            case LIGHTWEIGHT:
                 createLightGroupResource(context, groupInfo);
                 break;
             default:
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/GroupTaskListenerFactory.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/GroupTaskListenerFactory.java
index 8ec31e3ef..d84d7bed0 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/GroupTaskListenerFactory.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/GroupTaskListenerFactory.java
@@ -29,7 +29,7 @@ import 
org.apache.inlong.manager.service.mq.PulsarResourceCreateSelector;
 import org.apache.inlong.manager.service.mq.PulsarResourceDeleteSelector;
 import org.apache.inlong.manager.service.mq.TubeEventSelector;
 import org.apache.inlong.manager.service.resource.SinkResourceListener;
-import org.apache.inlong.manager.service.sort.CreateSortConfigListenerV2;
+import org.apache.inlong.manager.service.sort.SortConfigListener;
 import org.apache.inlong.manager.service.sort.ZookeeperDisabledSelector;
 import org.apache.inlong.manager.service.sort.light.LightGroupSortListener;
 import org.apache.inlong.manager.service.sort.light.LightGroupSortSelector;
@@ -95,7 +95,7 @@ public class GroupTaskListenerFactory implements 
PluginBinder, ServiceTaskListen
     private SinkResourceListener sinkResourceListener;
 
     @Autowired
-    private CreateSortConfigListenerV2 createSortConfigListener;
+    private SortConfigListener sortConfigListener;
 
     @Autowired
     private LightGroupSortListener lightGroupSortListener;
@@ -113,7 +113,7 @@ public class GroupTaskListenerFactory implements 
PluginBinder, ServiceTaskListen
         queueOperateListeners.put(createPulsarGroupTaskListener, new 
PulsarResourceCreateSelector());
         queueOperateListeners.put(deletePulsarResourceTaskListener, new 
PulsarResourceDeleteSelector());
         sortOperateListeners = new LinkedHashMap<>();
-        sortOperateListeners.put(createSortConfigListener, new 
ZookeeperDisabledSelector());
+        sortOperateListeners.put(sortConfigListener, new 
ZookeeperDisabledSelector());
         sortOperateListeners.put(lightGroupSortListener, new 
LightGroupSortSelector());
     }
 
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/StreamTaskListenerFactory.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/StreamTaskListenerFactory.java
index aa89ab8bf..f556c9b34 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/StreamTaskListenerFactory.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/listener/StreamTaskListenerFactory.java
@@ -28,7 +28,7 @@ import 
org.apache.inlong.manager.service.mq.DeletePulsarTopicTaskListener;
 import org.apache.inlong.manager.service.mq.PulsarTopicCreateSelector;
 import org.apache.inlong.manager.service.mq.PulsarTopicDeleteSelector;
 import org.apache.inlong.manager.service.resource.StreamSinkResourceListener;
-import org.apache.inlong.manager.service.sort.CreateStreamSortConfigListener;
+import org.apache.inlong.manager.service.sort.StreamSortConfigListener;
 import org.apache.inlong.manager.service.sort.ZookeeperEnabledSelector;
 import org.apache.inlong.manager.workflow.WorkflowContext;
 import 
org.apache.inlong.manager.workflow.definition.ServiceTaskListenerProvider;
@@ -69,7 +69,7 @@ public class StreamTaskListenerFactory implements 
PluginBinder, ServiceTaskListe
     @Autowired
     private DeletePulsarTopicTaskListener deletePulsarTopicTaskListener;
     @Autowired
-    private CreateStreamSortConfigListener createSortConfigListener;
+    private StreamSortConfigListener streamSortConfigListener;
     @Autowired
     private StreamSinkResourceListener sinkResourceListener;
 
@@ -81,7 +81,7 @@ public class StreamTaskListenerFactory implements 
PluginBinder, ServiceTaskListe
         queueOperateListeners.put(createPulsarSubscriptionTaskListener, new 
PulsarTopicCreateSelector());
         queueOperateListeners.put(deletePulsarTopicTaskListener, new 
PulsarTopicDeleteSelector());
         sortOperateListeners = new LinkedHashMap<>();
-        sortOperateListeners.put(createSortConfigListener, new 
ZookeeperEnabledSelector());
+        sortOperateListeners.put(streamSortConfigListener, new 
ZookeeperEnabledSelector());
         sinkOperateListeners = new LinkedHashMap<>();
         sinkOperateListeners.put(sinkResourceListener, context -> {
             ProcessForm processForm = context.getProcessForm();
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/DisableZkForSortTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/DisableZkForSortTest.java
index 25737e19b..adc4dd159 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/DisableZkForSortTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/DisableZkForSortTest.java
@@ -18,6 +18,7 @@
 package org.apache.inlong.manager.service.sort;
 
 import com.google.common.collect.Lists;
+import org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.enums.GroupOperateType;
 import org.apache.inlong.manager.common.enums.GroupStatus;
 import org.apache.inlong.manager.common.enums.ProcessStatus;
@@ -122,7 +123,8 @@ public class DisableZkForSortTest extends 
WorkflowServiceImplTest {
     //    @Test
     public void testCreateSortConfigInUpdateWorkflow() {
         InlongGroupInfo groupInfo = initGroupForm("PULSAR", "test20");
-        groupInfo.setEnableZookeeper(0);
+        groupInfo.setEnableZookeeper(InlongConstants.ENABLE_ZK);
+        
groupInfo.setEnableCreateResource(InlongConstants.ENABLE_CREATE_RESOURCE);
         groupService.updateStatus(GROUP_ID, 
GroupStatus.CONFIG_SUCCESSFUL.getCode(), OPERATOR);
         groupService.update(groupInfo.genRequest(), OPERATOR);
 
@@ -143,7 +145,7 @@ public class DisableZkForSortTest extends 
WorkflowServiceImplTest {
         Assertions.assertTrue(task instanceof ServiceTask);
         Assertions.assertEquals(2, task.getNameToListenerMap().size());
         List<TaskEventListener> listeners = 
Lists.newArrayList(task.getNameToListenerMap().values());
-        Assertions.assertTrue(listeners.get(1) instanceof 
CreateSortConfigListener);
+        Assertions.assertEquals(2, listeners.size());
         ProcessForm currentProcessForm = context.getProcessForm();
         InlongGroupInfo curGroupRequest = ((GroupResourceProcessForm) 
currentProcessForm).getGroupInfo();
         Assertions.assertEquals(1, curGroupRequest.getExtList().size());

Reply via email to