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