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 cdab4cd5e8 [INLONG-10353][Manager] Refactor code for building and 
submitting flink job (#10354)
cdab4cd5e8 is described below

commit cdab4cd5e82910de85660860af6b8b76f795bd40
Author: AloysZhang <[email protected]>
AuthorDate: Thu Jun 6 11:44:30 2024 +0800

    [INLONG-10353][Manager] Refactor code for building and submitting flink job 
(#10354)
---
 .../plugin/listener/StartupSortListener.java       | 87 +++----------------
 .../plugin/listener/StartupStreamListener.java     | 97 +---------------------
 .../inlong/manager/plugin/util/FlinkUtils.java     | 94 +++++++++++++++++++++
 3 files changed, 111 insertions(+), 167 deletions(-)

diff --git 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
index 0b0e55e369..038f35543f 100644
--- 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
+++ 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupSortListener.java
@@ -18,14 +18,9 @@
 package org.apache.inlong.manager.plugin.listener;
 
 import org.apache.inlong.manager.common.consts.InlongConstants;
-import org.apache.inlong.manager.common.consts.SinkType;
 import org.apache.inlong.manager.common.enums.GroupOperateType;
 import org.apache.inlong.manager.common.enums.TaskEvent;
-import org.apache.inlong.manager.common.util.JsonUtils;
-import org.apache.inlong.manager.plugin.flink.FlinkOperation;
-import org.apache.inlong.manager.plugin.flink.dto.FlinkInfo;
-import org.apache.inlong.manager.plugin.flink.enums.Constants;
-import org.apache.inlong.manager.pojo.sink.StreamSink;
+import org.apache.inlong.manager.plugin.util.FlinkUtils;
 import org.apache.inlong.manager.pojo.stream.InlongStreamExtInfo;
 import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
 import 
org.apache.inlong.manager.pojo.workflow.form.process.GroupResourceProcessForm;
@@ -34,19 +29,12 @@ 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 com.fasterxml.jackson.core.type.TypeReference;
 import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.commons.lang3.StringUtils;
-import org.apache.flink.api.common.JobStatus;
 
-import java.util.Collections;
+import java.util.ArrayList;
 import java.util.List;
-import java.util.Map;
 import java.util.stream.Collectors;
 
-import static 
org.apache.inlong.manager.plugin.util.FlinkUtils.getExceptionStackMsg;
-
 /**
  * Listener of startup sort.
  */
@@ -96,69 +84,20 @@ public class StartupSortListener implements 
SortOperateListener {
             return ListenerResult.success();
         }
 
+        List<ListenerResult> listenerResults = new ArrayList<>();
         for (InlongStreamInfo streamInfo : streamInfos) {
-            List<StreamSink> sinkList = streamInfo.getSinkList();
-            List<String> sinkTypes = 
sinkList.stream().map(StreamSink::getSinkType).collect(Collectors.toList());
-            if (CollectionUtils.isEmpty(sinkList) || 
!SinkType.containSortFlinkSink(sinkTypes)) {
-                log.warn("not any valid sink configured for groupId {} and 
streamId {}, reason: {},"
-                        + " skip launching sort job",
-                        CollectionUtils.isEmpty(sinkList) ? "no sink 
configured" : "no sort flink sink configured",
-                        groupId, streamInfo.getInlongStreamId());
-                continue;
-            }
-
-            List<InlongStreamExtInfo> extList = streamInfo.getExtList();
-            log.info("stream ext info: {}", extList);
-            Map<String, String> kvConf = extList.stream().filter(v -> 
StringUtils.isNotEmpty(v.getKeyName())
-                    && 
StringUtils.isNotEmpty(v.getKeyValue())).collect(Collectors.toMap(
-                            InlongStreamExtInfo::getKeyName,
-                            InlongStreamExtInfo::getKeyValue));
-
-            String sortExt = kvConf.get(InlongConstants.SORT_PROPERTIES);
-            if (StringUtils.isNotEmpty(sortExt)) {
-                Map<String, String> result = 
JsonUtils.OBJECT_MAPPER.convertValue(
-                        JsonUtils.OBJECT_MAPPER.readTree(sortExt), new 
TypeReference<Map<String, String>>() {
-                        });
-                kvConf.putAll(result);
-            }
-
-            String dataflow = kvConf.get(InlongConstants.DATAFLOW);
-            if (StringUtils.isEmpty(dataflow)) {
-                String message = String.format("dataflow is empty for groupId 
[%s], streamId [%s]", groupId,
-                        streamInfo.getInlongStreamId());
-                log.error(message);
-                return ListenerResult.fail(message);
-            }
-
-            FlinkInfo flinkInfo = new FlinkInfo();
-
-            String jobName = 
Constants.SORT_JOB_NAME_GENERATOR.apply(processForm) + InlongConstants.HYPHEN
-                    + streamInfo.getInlongStreamId();
-            flinkInfo.setJobName(jobName);
-            String sortUrl = kvConf.get(InlongConstants.SORT_URL);
-            flinkInfo.setEndpoint(sortUrl);
-            
flinkInfo.setInlongStreamInfoList(Collections.singletonList(streamInfo));
-            FlinkOperation flinkOperation = FlinkOperation.getInstance();
-            try {
-                flinkOperation.genPath(flinkInfo, dataflow);
-                flinkOperation.start(flinkInfo);
-                log.info("job submit success for groupId = {}, streamId = {}, 
jobId = {}", groupId,
-                        streamInfo.getInlongStreamId(), flinkInfo.getJobId());
-            } catch (Exception e) {
-                flinkInfo.setException(true);
-                flinkInfo.setExceptionMsg(getExceptionStackMsg(e));
-                flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
-
-                String message = String.format("startup sort failed for 
groupId [%s], streamId [%s]", groupId,
-                        streamInfo.getInlongStreamId());
-                log.error(message, e);
-                return ListenerResult.fail(message + e.getMessage());
-            }
+            listenerResults.add(FlinkUtils.submitFlinkJob(streamInfo,
+                    FlinkUtils.genFlinkJobName(processForm, streamInfo)));
+        }
 
-            saveInfo(streamInfo, InlongConstants.SORT_JOB_ID, 
flinkInfo.getJobId(), extList);
-            flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
+        // only one stream in group for now
+        // we can return the list of ListenerResult if support multi-stream in 
the future
+        List<ListenerResult> failedStreams = listenerResults.stream()
+                .filter(t -> !t.isSuccess()).collect(Collectors.toList());
+        if (failedStreams.isEmpty()) {
+            ListenerResult.success();
         }
-        return ListenerResult.success();
+        return ListenerResult.fail(failedStreams.get(0).getRemark());
     }
 
     /**
diff --git 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
index c66f76d467..931ba44550 100644
--- 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
+++ 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/listener/StartupStreamListener.java
@@ -18,15 +18,9 @@
 package org.apache.inlong.manager.plugin.listener;
 
 import org.apache.inlong.manager.common.consts.InlongConstants;
-import org.apache.inlong.manager.common.consts.SinkType;
 import org.apache.inlong.manager.common.enums.GroupOperateType;
 import org.apache.inlong.manager.common.enums.TaskEvent;
-import org.apache.inlong.manager.common.util.JsonUtils;
-import org.apache.inlong.manager.plugin.flink.FlinkOperation;
-import org.apache.inlong.manager.plugin.flink.dto.FlinkInfo;
-import org.apache.inlong.manager.plugin.flink.enums.Constants;
-import org.apache.inlong.manager.pojo.sink.StreamSink;
-import org.apache.inlong.manager.pojo.stream.InlongStreamExtInfo;
+import org.apache.inlong.manager.plugin.util.FlinkUtils;
 import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
 import org.apache.inlong.manager.pojo.workflow.form.process.ProcessForm;
 import 
org.apache.inlong.manager.pojo.workflow.form.process.StreamResourceProcessForm;
@@ -34,18 +28,7 @@ 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 com.fasterxml.jackson.core.type.TypeReference;
 import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.collections.CollectionUtils;
-import org.apache.commons.lang3.StringUtils;
-import org.apache.flink.api.common.JobStatus;
-
-import java.util.Collections;
-import java.util.List;
-import java.util.Map;
-import java.util.stream.Collectors;
-
-import static 
org.apache.inlong.manager.plugin.util.FlinkUtils.getExceptionStackMsg;
 
 /**
  * Listener for startup the Sort task for InlongStream
@@ -88,82 +71,10 @@ public class StartupStreamListener implements 
SortOperateListener {
         ProcessForm processForm = context.getProcessForm();
         StreamResourceProcessForm streamResourceProcessForm = 
(StreamResourceProcessForm) processForm;
         InlongStreamInfo streamInfo = 
streamResourceProcessForm.getStreamInfo();
-        List<InlongStreamExtInfo> streamExtList = streamInfo.getExtList();
-        log.info("inlong stream :{} ext info: {}", 
streamInfo.getInlongStreamId(), streamExtList);
-        final String groupId = streamInfo.getInlongGroupId();
-        final String streamId = streamInfo.getInlongStreamId();
-
-        List<StreamSink> sinkList = streamInfo.getSinkList();
-        List<String> sinkTypes = 
sinkList.stream().map(StreamSink::getSinkType).collect(Collectors.toList());
-        if (CollectionUtils.isEmpty(sinkList) || 
!SinkType.containSortFlinkSink(sinkTypes)) {
-            log.warn("not any sink configured for group {} and stream {}, skip 
launching sort job", groupId, streamId);
-            return ListenerResult.success();
-        }
-
-        List<InlongStreamExtInfo> extList = streamInfo.getExtList();
-        log.info("stream ext info: {}", extList);
-        Map<String, String> kvConf = extList.stream().filter(v -> 
StringUtils.isNotEmpty(v.getKeyName())
-                && 
StringUtils.isNotEmpty(v.getKeyValue())).collect(Collectors.toMap(
-                        InlongStreamExtInfo::getKeyName,
-                        InlongStreamExtInfo::getKeyValue));
-
-        String sortExt = kvConf.get(InlongConstants.SORT_PROPERTIES);
-        if (StringUtils.isNotEmpty(sortExt)) {
-            Map<String, String> result = JsonUtils.OBJECT_MAPPER.convertValue(
-                    JsonUtils.OBJECT_MAPPER.readTree(sortExt), new 
TypeReference<Map<String, String>>() {
-                    });
-            kvConf.putAll(result);
-        }
-
-        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
-        if (StringUtils.isEmpty(dataflow)) {
-            String message = String.format("dataflow is empty for groupId 
[%s], streamId [%s]", groupId,
-                    streamInfo.getInlongStreamId());
-            log.error(message);
-            return ListenerResult.fail(message);
-        }
+        log.info("inlong stream :{} ext info: {}", 
streamInfo.getInlongStreamId(), streamInfo.getExtList());
 
-        FlinkInfo flinkInfo = new FlinkInfo();
-
-        String jobName = Constants.SORT_JOB_NAME_GENERATOR.apply(processForm) 
+ InlongConstants.HYPHEN
-                + streamInfo.getInlongStreamId();
-        flinkInfo.setJobName(jobName);
-        String sortUrl = kvConf.get(InlongConstants.SORT_URL);
-        flinkInfo.setEndpoint(sortUrl);
-        
flinkInfo.setInlongStreamInfoList(Collections.singletonList(streamInfo));
-        FlinkOperation flinkOperation = FlinkOperation.getInstance();
-        try {
-            flinkOperation.genPath(flinkInfo, dataflow);
-            flinkOperation.start(flinkInfo);
-            log.info("job submit success for groupId = {}, streamId = {}, 
jobId = {}", groupId,
-                    streamInfo.getInlongStreamId(), flinkInfo.getJobId());
-        } catch (Exception e) {
-            flinkInfo.setException(true);
-            flinkInfo.setExceptionMsg(getExceptionStackMsg(e));
-            flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
-
-            String message = String.format("startup sort failed for groupId 
[%s], streamId [%s]", groupId,
-                    streamInfo.getInlongStreamId());
-            log.error(message, e);
-            return ListenerResult.fail(message + e.getMessage());
-        }
-
-        saveInfo(streamInfo, InlongConstants.SORT_JOB_ID, 
flinkInfo.getJobId(), extList);
-        flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
-        return ListenerResult.success();
-    }
-
-    /**
-     * Save stream ext info into list.
-     */
-    private void saveInfo(InlongStreamInfo streamInfo, String keyName, String 
keyValue,
-            List<InlongStreamExtInfo> extInfoList) {
-        InlongStreamExtInfo extInfo = new InlongStreamExtInfo();
-        extInfo.setInlongGroupId(streamInfo.getInlongGroupId());
-        extInfo.setInlongStreamId(streamInfo.getInlongStreamId());
-        extInfo.setKeyName(keyName);
-        extInfo.setKeyValue(keyValue);
-        extInfoList.add(extInfo);
+        String jobName = FlinkUtils.genFlinkJobName(processForm, streamInfo);
+        return FlinkUtils.submitFlinkJob(streamInfo, jobName);
     }
 
 }
diff --git 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/util/FlinkUtils.java
 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/util/FlinkUtils.java
index 345c997a8e..d4c7863371 100644
--- 
a/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/util/FlinkUtils.java
+++ 
b/inlong-manager/manager-plugins/base/src/main/java/org/apache/inlong/manager/plugin/util/FlinkUtils.java
@@ -17,11 +17,24 @@
 
 package org.apache.inlong.manager.plugin.util;
 
+import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.consts.SinkType;
+import org.apache.inlong.manager.common.util.JsonUtils;
+import org.apache.inlong.manager.plugin.flink.FlinkOperation;
 import org.apache.inlong.manager.plugin.flink.dto.FlinkConfig;
+import org.apache.inlong.manager.plugin.flink.dto.FlinkInfo;
 import org.apache.inlong.manager.plugin.flink.enums.Constants;
+import org.apache.inlong.manager.pojo.sink.StreamSink;
+import org.apache.inlong.manager.pojo.stream.InlongStreamExtInfo;
+import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
+import org.apache.inlong.manager.pojo.workflow.form.process.ProcessForm;
+import org.apache.inlong.manager.workflow.event.ListenerResult;
 
+import com.fasterxml.jackson.core.type.TypeReference;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.commons.collections.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flink.api.common.JobStatus;
 import org.apache.flink.configuration.Configuration;
 
 import java.io.BufferedReader;
@@ -37,10 +50,13 @@ import java.net.URLClassLoader;
 import java.nio.file.Path;
 import java.nio.file.Paths;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.List;
+import java.util.Map;
 import java.util.Properties;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
+import java.util.stream.Collectors;
 
 import static org.apache.inlong.manager.plugin.flink.enums.Constants.ADDRESS;
 import static org.apache.inlong.manager.plugin.flink.enums.Constants.DRAIN;
@@ -201,4 +217,82 @@ public class FlinkUtils {
         flinkConfig.setVersion(properties.getProperty(FLINK_VERSION));
         return flinkConfig;
     }
+
+    public static ListenerResult submitFlinkJob(InlongStreamInfo streamInfo, 
String jobName) throws Exception {
+        List<StreamSink> sinkList = streamInfo.getSinkList();
+        List<String> sinkTypes = 
sinkList.stream().map(StreamSink::getSinkType).collect(Collectors.toList());
+        if (CollectionUtils.isEmpty(sinkList) || 
!SinkType.containSortFlinkSink(sinkTypes)) {
+            log.warn("not any valid sink configured for groupId {} and 
streamId {}, reason: {},"
+                    + " skip launching sort job",
+                    CollectionUtils.isEmpty(sinkList) ? "no sink configured" : 
"no sort flink sink configured",
+                    streamInfo.getInlongGroupId(), 
streamInfo.getInlongStreamId());
+            return ListenerResult.success();
+        }
+        List<InlongStreamExtInfo> extList = streamInfo.getExtList();
+        log.info("stream ext info: {}", extList);
+        Map<String, String> kvConf = extList.stream().filter(v -> 
StringUtils.isNotEmpty(v.getKeyName())
+                && 
StringUtils.isNotEmpty(v.getKeyValue())).collect(Collectors.toMap(
+                        InlongStreamExtInfo::getKeyName,
+                        InlongStreamExtInfo::getKeyValue));
+
+        String sortExtProperties = kvConf.get(InlongConstants.SORT_PROPERTIES);
+        if (StringUtils.isNotEmpty(sortExtProperties)) {
+            Map<String, String> result = JsonUtils.OBJECT_MAPPER.convertValue(
+                    JsonUtils.OBJECT_MAPPER.readTree(sortExtProperties), new 
TypeReference<Map<String, String>>() {
+                    });
+            kvConf.putAll(result);
+        }
+
+        String dataflow = kvConf.get(InlongConstants.DATAFLOW);
+        if (StringUtils.isEmpty(dataflow)) {
+            String message = String.format("dataflow is empty for groupId 
[%s], streamId [%s]",
+                    streamInfo.getInlongGroupId(), 
streamInfo.getInlongStreamId());
+            log.error(message);
+            return ListenerResult.fail(message);
+        }
+
+        FlinkInfo flinkInfo = new FlinkInfo();
+        flinkInfo.setJobName(jobName);
+        String sortUrl = kvConf.get(InlongConstants.SORT_URL);
+        flinkInfo.setEndpoint(sortUrl);
+        
flinkInfo.setInlongStreamInfoList(Collections.singletonList(streamInfo));
+        FlinkOperation flinkOperation = FlinkOperation.getInstance();
+        try {
+            flinkOperation.genPath(flinkInfo, dataflow);
+            flinkOperation.start(flinkInfo);
+            log.info("job submit success for groupId = {}, streamId = {}, 
jobId = {}",
+                    streamInfo.getInlongGroupId(), 
streamInfo.getInlongStreamId(), flinkInfo.getJobId());
+        } catch (Exception e) {
+            flinkInfo.setException(true);
+            flinkInfo.setExceptionMsg(getExceptionStackMsg(e));
+            flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
+
+            String message = String.format("startup sort failed for groupId 
[%s], streamId [%s]",
+                    streamInfo.getInlongGroupId(), 
streamInfo.getInlongStreamId());
+            log.error(message, e);
+            return ListenerResult.fail(message + e.getMessage());
+        }
+
+        saveInfo(streamInfo, InlongConstants.SORT_JOB_ID, 
flinkInfo.getJobId(), extList);
+        flinkOperation.pollJobStatus(flinkInfo, JobStatus.RUNNING);
+        return ListenerResult.success();
+    }
+
+    /**
+     * Save stream ext info into list.
+     */
+    public static void saveInfo(InlongStreamInfo streamInfo, String keyName, 
String keyValue,
+            List<InlongStreamExtInfo> extInfoList) {
+        InlongStreamExtInfo extInfo = new InlongStreamExtInfo();
+        extInfo.setInlongGroupId(streamInfo.getInlongGroupId());
+        extInfo.setInlongStreamId(streamInfo.getInlongStreamId());
+        extInfo.setKeyName(keyName);
+        extInfo.setKeyValue(keyValue);
+        extInfoList.add(extInfo);
+    }
+
+    public static String genFlinkJobName(ProcessForm processForm, 
InlongStreamInfo streamInfo) {
+        return Constants.SORT_JOB_NAME_GENERATOR.apply(processForm) + 
InlongConstants.HYPHEN
+                + streamInfo.getInlongStreamId();
+    }
 }

Reply via email to