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());