This is an automated email from the ASF dual-hosted git repository.
SbloodyS 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 db8ec27394 [Fix-18540][Master] Reset the runtime state when recreating
a failed task instance (#18541)
db8ec27394 is described below
commit db8ec273940812069f086f0eadcadb23d2b88b81
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Mon Aug 10 14:47:23 2026 +0530
[Fix-18540][Master] Reset the runtime state when recreating a failed task
instance (#18541)
---
.../FailedRecoverTaskInstanceFactory.java | 7 ++
.../FailedRecoverTaskInstanceFactoryTest.java | 125 +++++++++++++++++++++
2 files changed, 132 insertions(+)
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
index 1e0b4107a4..041c338180 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
@@ -51,8 +51,15 @@ public class FailedRecoverTaskInstanceFactory
taskInstance.setHost(null);
taskInstance.setVarPool(null);
taskInstance.setSubmitTime(new Date());
+ taskInstance.setStartTime(null);
+ taskInstance.setEndTime(null);
taskInstance.setLogPath(null);
taskInstance.setExecutePath(null);
+ taskInstance.setPid(0);
+ // The recreated task instance is a new attempt, it should get the
whole retry budget back, otherwise a task
+ // which already exhausted its retry times will never be retried again
in the recovered workflow.
+ taskInstance.setRetryTimes(0);
+ taskInstance.setAlertFlag(Flag.NO);
taskInstanceDao.insert(taskInstance);
needRecoverTaskInstance.setFlag(Flag.NO);
diff --git
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
new file mode 100644
index 0000000000..27684fa258
--- /dev/null
+++
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
@@ -0,0 +1,125 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.server.master.engine.task.execution;
+
+import static com.google.common.truth.Truth.assertThat;
+import static org.mockito.Mockito.verify;
+
+import org.apache.dolphinscheduler.common.enums.Flag;
+import org.apache.dolphinscheduler.dao.entity.TaskInstance;
+import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao;
+import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus;
+
+import java.util.Date;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+@ExtendWith(MockitoExtension.class)
+class FailedRecoverTaskInstanceFactoryTest {
+
+ private static final int MAX_RETRY_TIMES = 3;
+
+ @InjectMocks
+ private FailedRecoverTaskInstanceFactory failedRecoverTaskInstanceFactory;
+
+ @Mock
+ private TaskInstanceDao taskInstanceDao;
+
+ @Test
+ @DisplayName("Test the recreated task instance gets the whole retry budget
back")
+ void testCreateTaskInstance_resetRetryTimes() {
+ final TaskInstance needRecoverTaskInstance =
createExhaustedFailedTaskInstance();
+
+ final TaskInstance taskInstance =
failedRecoverTaskInstanceFactory.builder()
+ .withTaskInstance(needRecoverTaskInstance)
+ .build();
+
+ assertThat(taskInstance.getRetryTimes()).isEqualTo(0);
+ assertThat(taskInstance.getMaxRetryTimes()).isEqualTo(MAX_RETRY_TIMES);
+ }
+
+ @Test
+ @DisplayName("Test the recreated task instance doesn't inherit the runtime
state of the failed attempt")
+ void testCreateTaskInstance_clearPreviousRuntimeState() {
+ final TaskInstance needRecoverTaskInstance =
createExhaustedFailedTaskInstance();
+
+ final TaskInstance taskInstance =
failedRecoverTaskInstanceFactory.builder()
+ .withTaskInstance(needRecoverTaskInstance)
+ .build();
+
+ assertThat(taskInstance.getId()).isNull();
+
assertThat(taskInstance.getState()).isEqualTo(TaskExecutionStatus.SUBMITTED_SUCCESS);
+ assertThat(taskInstance.getStartTime()).isNull();
+ assertThat(taskInstance.getEndTime()).isNull();
+ assertThat(taskInstance.getHost()).isNull();
+ assertThat(taskInstance.getLogPath()).isNull();
+ assertThat(taskInstance.getExecutePath()).isNull();
+ assertThat(taskInstance.getVarPool()).isNull();
+ assertThat(taskInstance.getPid()).isEqualTo(0);
+ // The appLink is deliberately kept: SubWorkflowLogicTask stores its
runtime context there and needs it
+ // to recover the existing sub workflow instance instead of triggering
a new one.
+
assertThat(taskInstance.getAppLink()).isEqualTo("application_1717430400000_0001");
+ assertThat(taskInstance.getAlertFlag()).isEqualTo(Flag.NO);
+ assertThat(taskInstance.getSubmitTime()).isNotNull();
+ }
+
+ @Test
+ @DisplayName("Test the recovered task instance is inserted and the origin
one is marked as invalid")
+ void testCreateTaskInstance_markOriginTaskInstanceInvalid() {
+ final TaskInstance needRecoverTaskInstance =
createExhaustedFailedTaskInstance();
+
+ final TaskInstance taskInstance =
failedRecoverTaskInstanceFactory.builder()
+ .withTaskInstance(needRecoverTaskInstance)
+ .build();
+
+ assertThat(taskInstance.getFlag()).isEqualTo(Flag.YES);
+ assertThat(needRecoverTaskInstance.getFlag()).isEqualTo(Flag.NO);
+ verify(taskInstanceDao).insert(taskInstance);
+ verify(taskInstanceDao).updateById(needRecoverTaskInstance);
+ }
+
+ /**
+ * Create a task instance which is failed and has already used up all its
retry times.
+ */
+ private TaskInstance createExhaustedFailedTaskInstance() {
+ final TaskInstance taskInstance = new TaskInstance();
+ taskInstance.setId(1);
+ taskInstance.setName("A");
+ taskInstance.setState(TaskExecutionStatus.FAILURE);
+ taskInstance.setFlag(Flag.YES);
+ taskInstance.setRetryTimes(MAX_RETRY_TIMES);
+ taskInstance.setMaxRetryTimes(MAX_RETRY_TIMES);
+ taskInstance.setRetryInterval(1);
+ taskInstance.setStartTime(new Date());
+ taskInstance.setEndTime(new Date());
+ taskInstance.setSubmitTime(new Date());
+ taskInstance.setHost("127.0.0.1:1234");
+ taskInstance.setLogPath("/tmp/log/A.log");
+ taskInstance.setExecutePath("/tmp/exec/A");
+ taskInstance.setVarPool("[]");
+ taskInstance.setPid(1234);
+ taskInstance.setAppLink("application_1717430400000_0001");
+ taskInstance.setAlertFlag(Flag.YES);
+ return taskInstance;
+ }
+}