This is an automated email from the ASF dual-hosted git repository.
gosonzhang 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 43e312117 [INLONG-6516][Manager] Skip starting the Sort task in the
InlongStream workflow (#6606)
43e312117 is described below
commit 43e3121178a52a3f88e164b81c4b57c3f13e8cca
Author: fuweng11 <[email protected]>
AuthorDate: Wed Nov 23 15:52:14 2022 +0800
[INLONG-6516][Manager] Skip starting the Sort task in the InlongStream
workflow (#6606)
Co-authored-by: healchow <[email protected]>
---
.../plugin/listener/StartupStreamListener.java | 122 ++-------------------
1 file changed, 7 insertions(+), 115 deletions(-)
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 4c0a19c08..f4dd053aa 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
@@ -17,36 +17,14 @@
package org.apache.inlong.manager.plugin.listener;
-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.inlong.manager.common.consts.InlongConstants;
-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.FlinkService;
-import org.apache.inlong.manager.plugin.flink.dto.FlinkInfo;
-import org.apache.inlong.manager.plugin.flink.enums.Constants;
-import org.apache.inlong.manager.pojo.group.InlongGroupExtInfo;
-import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
-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.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 java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-
-import static
org.apache.inlong.manager.plugin.util.FlinkUtils.getExceptionStackMsg;
-
/**
- * Listener for startup stream sort
+ * Listener for startup the Sort task for InlongStream
*/
@Slf4j
public class StartupStreamListener implements SortOperateListener {
@@ -56,105 +34,19 @@ public class StartupStreamListener implements
SortOperateListener {
return TaskEvent.COMPLETE;
}
+ /**
+ * Currently, the process of starting Sort tasks has been initiated in
{@link StartupSortListener}.
+ * <p/>Because the Sort task is only associated with InlongGroup, no need
to start it for InlongStream.
+ */
@Override
public boolean accept(WorkflowContext workflowContext) {
- ProcessForm processForm = workflowContext.getProcessForm();
- String groupId = processForm.getInlongGroupId();
- if (!(processForm instanceof StreamResourceProcessForm)) {
- log.info("not add startup stream listener, not
StreamResourceProcessForm for groupId [{}]", groupId);
- return false;
- }
-
- StreamResourceProcessForm streamProcessForm =
(StreamResourceProcessForm) processForm;
- String streamId =
streamProcessForm.getStreamInfo().getInlongStreamId();
- if (streamProcessForm.getGroupOperateType() != GroupOperateType.INIT) {
- log.info("not add startup stream listener, as the operate was not
INIT for groupId [{}] streamId [{}]",
- groupId, streamId);
- return false;
- }
-
- log.info("add startup stream listener for groupId [{}] streamId [{}]",
groupId, streamId);
- return true;
+ log.info("not need to start the sort task for InlongStream");
+ return false;
}
@Override
public ListenerResult listen(WorkflowContext context) throws Exception {
- ProcessForm processForm = context.getProcessForm();
- StreamResourceProcessForm streamResourceProcessForm =
(StreamResourceProcessForm) processForm;
- InlongGroupInfo groupInfo = streamResourceProcessForm.getGroupInfo();
- List<InlongGroupExtInfo> groupExtList = groupInfo.getExtList();
- log.info("inlong group :{} ext info: {}",
groupInfo.getInlongGroupId(), groupExtList);
- 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();
-
- if
(CollectionUtils.isEmpty(streamResourceProcessForm.getStreamInfo().getSinkList()))
{
- log.warn("not any sink configured for group {} and stream {}, skip
launching sort job", groupId, streamId);
- return ListenerResult.success();
- }
-
- Map<String, String> kvConf = new HashMap<>();
- groupExtList.forEach(groupExtInfo ->
kvConf.put(groupExtInfo.getKeyName(), groupExtInfo.getKeyValue()));
- streamExtList.forEach(extInfo -> {
- kvConf.put(extInfo.getKeyName(), extInfo.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]
and streamId [%s]", groupId, streamId);
- log.error(message);
- return ListenerResult.fail(message);
- }
-
- FlinkInfo flinkInfo = new FlinkInfo();
- String jobName = Constants.INLONG +
context.getProcessForm().getInlongGroupId();
- flinkInfo.setJobName(jobName);
- String sortUrl = kvConf.get(InlongConstants.SORT_URL);
- flinkInfo.setEndpoint(sortUrl);
-
- FlinkService flinkService = new FlinkService(flinkInfo.getEndpoint());
- FlinkOperation flinkOperation = new FlinkOperation(flinkService);
-
- try {
- flinkOperation.genPath(flinkInfo, dataflow);
- flinkOperation.start(flinkInfo);
- log.info("job submit success, jobId is [{}]",
flinkInfo.getJobId());
- } catch (Exception e) {
- flinkOperation.pollJobStatus(flinkInfo);
- flinkInfo.setException(true);
- flinkInfo.setExceptionMsg(getExceptionStackMsg(e));
- flinkOperation.pollJobStatus(flinkInfo);
-
- String message = String.format("startup sort failed for groupId
[%s] streamId [%s]", groupId, streamId);
- log.error(message, e);
- return ListenerResult.fail(message + e.getMessage());
- }
-
- saveInfo(groupId, streamId, InlongConstants.SORT_JOB_ID,
flinkInfo.getJobId(), streamExtList);
- flinkOperation.pollJobStatus(flinkInfo);
return ListenerResult.success();
}
- /**
- * Save ext info into list.
- */
- private void saveInfo(String inlongGroupId, String inlongStreamId, String
keyName, String keyValue,
- List<InlongStreamExtInfo> extInfoList) {
- InlongStreamExtInfo extInfo = new InlongStreamExtInfo();
- extInfo.setInlongGroupId(inlongGroupId);
- extInfo.setInlongStreamId(inlongStreamId);
- extInfo.setKeyName(keyName);
- extInfo.setKeyValue(keyValue);
- extInfoList.add(extInfo);
- }
-
}