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/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 134f27502 [INLONG-4315][Manager] Fix incorrect task service node order
in create-group workflow (#4323)
134f27502 is described below
commit 134f275027224ed39b1ba12ca16248267cd16bb4
Author: woofyzhao <[email protected]>
AuthorDate: Thu May 26 18:58:11 2022 +0800
[INLONG-4315][Manager] Fix incorrect task service node order in
create-group workflow (#4323)
Co-authored-by: healzhou <[email protected]>
---
.../group/CreateGroupWorkflowDefinition.java | 64 ++++++++++-----------
.../stream/CreateStreamWorkflowDefinition.java | 65 +++++++++++-----------
2 files changed, 67 insertions(+), 62 deletions(-)
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/CreateGroupWorkflowDefinition.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/CreateGroupWorkflowDefinition.java
index 7bbd9fd77..103106945 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/CreateGroupWorkflowDefinition.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/group/CreateGroupWorkflowDefinition.java
@@ -20,11 +20,11 @@ package org.apache.inlong.manager.service.workflow.group;
import lombok.extern.slf4j.Slf4j;
import
org.apache.inlong.manager.common.pojo.workflow.form.GroupResourceProcessForm;
import org.apache.inlong.manager.service.workflow.ProcessName;
-import
org.apache.inlong.manager.service.workflow.listener.GroupTaskListenerFactory;
import org.apache.inlong.manager.service.workflow.WorkflowDefinition;
import
org.apache.inlong.manager.service.workflow.group.listener.GroupCompleteProcessListener;
import
org.apache.inlong.manager.service.workflow.group.listener.GroupFailedProcessListener;
import
org.apache.inlong.manager.service.workflow.group.listener.GroupInitProcessListener;
+import
org.apache.inlong.manager.service.workflow.listener.GroupTaskListenerFactory;
import org.apache.inlong.manager.workflow.definition.EndEvent;
import org.apache.inlong.manager.workflow.definition.ServiceTask;
import org.apache.inlong.manager.workflow.definition.ServiceTaskType;
@@ -69,31 +69,15 @@ public class CreateGroupWorkflowDefinition implements
WorkflowDefinition {
StartEvent startEvent = new StartEvent();
process.setStartEvent(startEvent);
- // init DataSource
- ServiceTask initDataSourceTask = new ServiceTask();
- initDataSourceTask.setName("initSource");
- initDataSourceTask.setDisplayName("Group-InitSource");
- initDataSourceTask.addServiceTaskType(ServiceTaskType.INIT_SOURCE);
- initDataSourceTask.addListenerProvider(groupTaskListenerFactory);
- process.addTask(initDataSourceTask);
-
- // init MQ resource
- ServiceTask initMQResourceTask = new ServiceTask();
- initMQResourceTask.setName("initMQ");
- initMQResourceTask.setDisplayName("Group-InitMQ");
- initMQResourceTask.addServiceTaskType(ServiceTaskType.INIT_MQ);
- initMQResourceTask.addListenerProvider(groupTaskListenerFactory);
- process.addTask(initMQResourceTask);
-
- // init Sort resource
- ServiceTask initSortResourceTask = new ServiceTask();
- initSortResourceTask.setName("initSort");
- initSortResourceTask.setDisplayName("Group-InitSort");
- initSortResourceTask.addServiceTaskType(ServiceTaskType.INIT_SORT);
- initSortResourceTask.addListenerProvider(groupTaskListenerFactory);
- process.addTask(initSortResourceTask);
-
- // init sink
+ // init MQ
+ ServiceTask initMQTask = new ServiceTask();
+ initMQTask.setName("initMQ");
+ initMQTask.setDisplayName("Group-InitMQ");
+ initMQTask.addServiceTaskType(ServiceTaskType.INIT_MQ);
+ initMQTask.addListenerProvider(groupTaskListenerFactory);
+ process.addTask(initMQTask);
+
+ // init Sink
ServiceTask initSinkTask = new ServiceTask();
initSinkTask.setName("initSink");
initSinkTask.setDisplayName("Group-InitSink");
@@ -101,15 +85,33 @@ public class CreateGroupWorkflowDefinition implements
WorkflowDefinition {
initSinkTask.addListenerProvider(groupTaskListenerFactory);
process.addTask(initSinkTask);
+ // init Sort
+ ServiceTask initSortTask = new ServiceTask();
+ initSortTask.setName("initSort");
+ initSortTask.setDisplayName("Group-InitSort");
+ initSortTask.addServiceTaskType(ServiceTaskType.INIT_SORT);
+ initSortTask.addListenerProvider(groupTaskListenerFactory);
+ process.addTask(initSortTask);
+
+ // init Source
+ ServiceTask initSourceTask = new ServiceTask();
+ initSourceTask.setName("initSource");
+ initSourceTask.setDisplayName("Group-InitSource");
+ initSourceTask.addServiceTaskType(ServiceTaskType.INIT_SOURCE);
+ initSourceTask.addListenerProvider(groupTaskListenerFactory);
+ process.addTask(initSourceTask);
+
// End node
EndEvent endEvent = new EndEvent();
process.setEndEvent(endEvent);
- startEvent.addNext(initDataSourceTask);
- initDataSourceTask.addNext(initMQResourceTask);
- initMQResourceTask.addNext(initSortResourceTask);
- initSortResourceTask.addNext(initSinkTask);
- initSinkTask.addNext(endEvent);
+ // Task dependency order: 1.MQ -> 2.Sink -> 3.Sort -> 4.Source
+ // To ensure that after some tasks fail, data will not start to be
collected by source or consumed by sort
+ startEvent.addNext(initMQTask);
+ initMQTask.addNext(initSinkTask);
+ initSinkTask.addNext(initSortTask);
+ initSortTask.addNext(initSourceTask);
+ initSourceTask.addNext(endEvent);
return process;
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/stream/CreateStreamWorkflowDefinition.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/stream/CreateStreamWorkflowDefinition.java
index dc99207c0..d85f67944 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/stream/CreateStreamWorkflowDefinition.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/stream/CreateStreamWorkflowDefinition.java
@@ -51,13 +51,14 @@ public class CreateStreamWorkflowDefinition implements
WorkflowDefinition {
@Override
public WorkflowProcess defineProcess() {
+
// Configuration process
WorkflowProcess process = new WorkflowProcess();
process.addListener(streamInitProcessListener);
process.addListener(streamFailedProcessListener);
process.addListener(streamCompleteProcessListener);
- process.setType("Stream resource creation");
+ process.setType("Stream Resource Creation");
process.setName(getProcessName().name());
process.setDisplayName(getProcessName().getDisplayName());
process.setFormClass(StreamResourceProcessForm.class);
@@ -68,31 +69,15 @@ public class CreateStreamWorkflowDefinition implements
WorkflowDefinition {
StartEvent startEvent = new StartEvent();
process.setStartEvent(startEvent);
- // init DataSource
- ServiceTask initDataSourceTask = new ServiceTask();
- initDataSourceTask.setName("initSource");
- initDataSourceTask.setDisplayName("Stream-InitSource");
- initDataSourceTask.addServiceTaskType(ServiceTaskType.INIT_SOURCE);
- initDataSourceTask.addListenerProvider(streamTaskListenerFactory);
- process.addTask(initDataSourceTask);
-
- // init MQ topic
- ServiceTask initMQResourceTask = new ServiceTask();
- initMQResourceTask.setName("initMQ");
- initMQResourceTask.setDisplayName("Stream-InitMQ");
- initMQResourceTask.addServiceTaskType(ServiceTaskType.INIT_MQ);
- initMQResourceTask.addListenerProvider(streamTaskListenerFactory);
- process.addTask(initMQResourceTask);
-
- // init sort config
- ServiceTask initSortResourceTask = new ServiceTask();
- initSortResourceTask.setName("initSort");
- initSortResourceTask.setDisplayName("Stream-InitSort");
- initSortResourceTask.addServiceTaskType(ServiceTaskType.INIT_SORT);
- initSortResourceTask.addListenerProvider(streamTaskListenerFactory);
- process.addTask(initSortResourceTask);
-
- // init DataSink
+ // init MQ
+ ServiceTask initMQTask = new ServiceTask();
+ initMQTask.setName("initMQ");
+ initMQTask.setDisplayName("Stream-InitMQ");
+ initMQTask.addServiceTaskType(ServiceTaskType.INIT_MQ);
+ initMQTask.addListenerProvider(streamTaskListenerFactory);
+ process.addTask(initMQTask);
+
+ // init Sink
ServiceTask initSinkTask = new ServiceTask();
initSinkTask.setName("initSink");
initSinkTask.setDisplayName("Stream-InitSink");
@@ -100,15 +85,33 @@ public class CreateStreamWorkflowDefinition implements
WorkflowDefinition {
initSinkTask.addListenerProvider(streamTaskListenerFactory);
process.addTask(initSinkTask);
+ // init Sort
+ ServiceTask initSortTask = new ServiceTask();
+ initSortTask.setName("initSort");
+ initSortTask.setDisplayName("Stream-InitSort");
+ initSortTask.addServiceTaskType(ServiceTaskType.INIT_SORT);
+ initSortTask.addListenerProvider(streamTaskListenerFactory);
+ process.addTask(initSortTask);
+
+ // init Source
+ ServiceTask initSourceTask = new ServiceTask();
+ initSourceTask.setName("initSource");
+ initSourceTask.setDisplayName("Stream-InitSource");
+ initSourceTask.addServiceTaskType(ServiceTaskType.INIT_SOURCE);
+ initSourceTask.addListenerProvider(streamTaskListenerFactory);
+ process.addTask(initSourceTask);
+
// End node
EndEvent endEvent = new EndEvent();
process.setEndEvent(endEvent);
- startEvent.addNext(initDataSourceTask);
- initDataSourceTask.addNext(initMQResourceTask);
- initMQResourceTask.addNext(initSortResourceTask);
- initSortResourceTask.addNext(initSinkTask);
- initSinkTask.addNext(endEvent);
+ // Task dependency order: 1.MQ -> 2.Sink -> 3.Sort -> 4.Source
+ // To ensure that after some tasks fail, data will not start to be
collected by source or consumed by sort
+ startEvent.addNext(initMQTask);
+ initMQTask.addNext(initSinkTask);
+ initSinkTask.addNext(initSortTask);
+ initSortTask.addNext(initSourceTask);
+ initSourceTask.addNext(endEvent);
return process;
}