This is an automated email from the ASF dual-hosted git repository.
wenjun pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
The following commit(s) were added to refs/heads/dev by this push:
new b8e3032b9e [Chore] Remove unused method in IWorkflowFailureStrategy
(#18159)
b8e3032b9e is described below
commit b8e3032b9e751eea3cd504cfa91b103ee7f9fc73
Author: Wenjun Ruan <[email protected]>
AuthorDate: Sat Apr 11 22:27:16 2026 +0800
[Chore] Remove unused method in IWorkflowFailureStrategy (#18159)
---
.../policy/ContinueWorkflowFailureStrategy.java | 6 ----
.../policy/EndWorkflowFailureStrategy.java | 6 ----
.../workflow/policy/IWorkflowFailureStrategy.java | 3 --
.../statemachine/AbstractWorkflowStateAction.java | 5 ---
.../integration/WorkflowTestCaseContext.java | 13 ++++++--
.../WorkflowTestCaseContextFactory.java | 38 +++++++++++++++-------
6 files changed, 36 insertions(+), 35 deletions(-)
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/ContinueWorkflowFailureStrategy.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/ContinueWorkflowFailureStrategy.java
index 927a19b8bf..3e904ed228 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/ContinueWorkflowFailureStrategy.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/ContinueWorkflowFailureStrategy.java
@@ -32,10 +32,4 @@ public class ContinueWorkflowFailureStrategy implements
IWorkflowFailureStrategy
// do nothing, just continue workflow execution
}
- @Override
- public boolean canTriggerSuccessor(IWorkflowExecutionRunnable
workflowExecutionRunnable,
- ITaskExecutionRunnable
taskExecutionRunnable) {
- return true;
- }
-
}
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/EndWorkflowFailureStrategy.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/EndWorkflowFailureStrategy.java
index 142f0dc982..f7a3e34daa 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/EndWorkflowFailureStrategy.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/EndWorkflowFailureStrategy.java
@@ -37,10 +37,4 @@ public class EndWorkflowFailureStrategy implements
IWorkflowFailureStrategy {
workflowExecutionRunnable.killActiveTasks();
}
- @Override
- public boolean canTriggerSuccessor(IWorkflowExecutionRunnable
workflowExecutionRunnable,
- ITaskExecutionRunnable
taskExecutionRunnable) {
- return
!workflowExecutionRunnable.getWorkflowExecutionGraph().isExistFailureTaskExecutionRunnableChain();
- }
-
}
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/IWorkflowFailureStrategy.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/IWorkflowFailureStrategy.java
index bb6f36176a..85a3b1219d 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/IWorkflowFailureStrategy.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/policy/IWorkflowFailureStrategy.java
@@ -28,7 +28,4 @@ public interface IWorkflowFailureStrategy {
void onTaskFailure(IWorkflowExecutionRunnable workflowExecutionRunnable,
ITaskExecutionRunnable taskExecutionRunnable);
- boolean canTriggerSuccessor(IWorkflowExecutionRunnable
workflowExecutionRunnable,
- ITaskExecutionRunnable taskExecutionRunnable);
-
}
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/statemachine/AbstractWorkflowStateAction.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/statemachine/AbstractWorkflowStateAction.java
index 0f943e9245..36ad66d17e 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/statemachine/AbstractWorkflowStateAction.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/workflow/statemachine/AbstractWorkflowStateAction.java
@@ -140,11 +140,6 @@ public abstract class AbstractWorkflowStateAction
implements IWorkflowStateActio
return;
}
- if
(!workflowFailureStrategy.canTriggerSuccessor(workflowExecutionRunnable,
taskExecutionRunnable)) {
- emitWorkflowFinishedEventIfApplicable(workflowExecutionRunnable);
- return;
- }
-
triggerTasks(workflowExecutionRunnable,
workflowExecutionGraph.getSuccessors(taskExecutionRunnable));
}
diff --git
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContext.java
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContext.java
index 35cee61c93..ae065542e2 100644
---
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContext.java
+++
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContext.java
@@ -29,6 +29,7 @@ import
org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation;
import org.apache.commons.collections4.CollectionUtils;
import java.util.List;
+import java.util.stream.Collectors;
import lombok.AllArgsConstructor;
import lombok.Data;
@@ -66,9 +67,15 @@ public class WorkflowTestCaseContext {
if (CollectionUtils.isEmpty(workflows)) {
throw new IllegalStateException("workflows is empty");
}
- return workflows.stream()
+ List<WorkflowDefinition> collect = workflows.stream()
.filter(workflow -> workflow.getName().equals(name))
- .findFirst()
- .orElseThrow(() -> new IllegalStateException("Workflow with
name " + name + " not found"));
+ .collect(Collectors.toList());
+ if (CollectionUtils.isEmpty(collect)) {
+ throw new IllegalStateException("Workflow with name " + name + "
not found");
+ }
+ if (collect.size() > 1) {
+ throw new IllegalStateException("Multiple workflows with name " +
name + " found");
+ }
+ return collect.get(0);
}
}
diff --git
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContextFactory.java
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContextFactory.java
index 8ff18222d0..b84f789b2b 100644
---
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContextFactory.java
+++
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/integration/WorkflowTestCaseContextFactory.java
@@ -90,34 +90,36 @@ public class WorkflowTestCaseContextFactory {
initializeWorkflowDefinitionToDB(workflowTestCaseContext.getWorkflows());
initializeTaskDefinitionsToDB(workflowTestCaseContext.getTasks());
initializeTaskRelationsToDB(workflowTestCaseContext.getTaskRelations());
- if
(CollectionUtils.isNotEmpty(workflowTestCaseContext.getWorkflowInstances())) {
-
initializeWorkflowInstancesToDB(workflowTestCaseContext.getWorkflowInstances());
- }
- if
(CollectionUtils.isNotEmpty(workflowTestCaseContext.getTaskInstances())) {
-
initializeTaskInstancesToDB(workflowTestCaseContext.getTaskInstances());
- }
- if
(CollectionUtils.isNotEmpty(workflowTestCaseContext.getTaskGroups())) {
- initializeTaskGroupsToDB(workflowTestCaseContext.getTaskGroups());
- }
- if
(CollectionUtils.isNotEmpty(workflowTestCaseContext.getEnvironments())) {
-
initializeEnvironmentToDB(workflowTestCaseContext.getEnvironments());
- }
+
+
initializeWorkflowInstancesToDB(workflowTestCaseContext.getWorkflowInstances());
+
initializeTaskInstancesToDB(workflowTestCaseContext.getTaskInstances());
+ initializeTaskGroupsToDB(workflowTestCaseContext.getTaskGroups());
+ initializeEnvironmentToDB(workflowTestCaseContext.getEnvironments());
return workflowTestCaseContext;
}
private void initializeTaskInstancesToDB(List<TaskInstance> taskInstances)
{
+ if (CollectionUtils.isEmpty(taskInstances)) {
+ return;
+ }
for (TaskInstance taskInstance : taskInstances) {
taskInstanceDao.insert(taskInstance);
}
}
private void initializeWorkflowInstancesToDB(List<WorkflowInstance>
workflowInstances) {
+ if (CollectionUtils.isEmpty(workflowInstances)) {
+ return;
+ }
for (WorkflowInstance workflowInstance : workflowInstances) {
workflowInstanceDao.insert(workflowInstance);
}
}
private void initializeWorkflowDefinitionToDB(final
List<WorkflowDefinition> workflowDefinitions) {
+ if (CollectionUtils.isEmpty(workflowDefinitions)) {
+ return;
+ }
for (final WorkflowDefinition workflowDefinition :
workflowDefinitions) {
workflowDefinitionDao.insert(workflowDefinition);
final WorkflowDefinitionLog workflowDefinitionLog = new
WorkflowDefinitionLog(workflowDefinition);
@@ -128,6 +130,9 @@ public class WorkflowTestCaseContextFactory {
}
private void initializeTaskDefinitionsToDB(final List<TaskDefinition>
taskDefinitions) {
+ if (CollectionUtils.isEmpty(taskDefinitions)) {
+ return;
+ }
for (final TaskDefinition taskDefinition : taskDefinitions) {
taskDefinitionDao.insert(taskDefinition);
@@ -139,6 +144,9 @@ public class WorkflowTestCaseContextFactory {
}
private void initializeTaskRelationsToDB(final List<WorkflowTaskRelation>
taskRelations) {
+ if (CollectionUtils.isEmpty(taskRelations)) {
+ return;
+ }
for (final WorkflowTaskRelation taskRelation : taskRelations) {
workflowTaskRelationMapper.insert(taskRelation);
@@ -153,12 +161,18 @@ public class WorkflowTestCaseContextFactory {
}
private void initializeTaskGroupsToDB(final List<TaskGroup> taskGroups) {
+ if (CollectionUtils.isEmpty(taskGroups)) {
+ return;
+ }
for (final TaskGroup taskGroup : taskGroups) {
taskGroupDao.insert(taskGroup);
}
}
private void initializeEnvironmentToDB(final List<Environment>
environments) {
+ if (CollectionUtils.isEmpty(environments)) {
+ return;
+ }
for (final Environment environment : environments) {
environmentDao.insert(environment);
}