Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions docs/docs/en/guide/upgrade/incompatible.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,3 +44,7 @@ This document records the incompatible updates between each version. You need to

* Remove import and export of workflow definition. ([#17940])(https://github.com/apache/dolphinscheduler/issues/17940)

## 3.5.0

* Add the `missed_fire_policy` column to `t_ds_schedules`. Existing schedules default to `FIRE_ALL_MISSED` to preserve the previous Quartz `IgnoreMisfires` behavior. ([#18464](https://github.com/apache/dolphinscheduler/pull/18464))

4 changes: 4 additions & 0 deletions docs/docs/zh/guide/upgrade/incompatible.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,3 +44,7 @@

* 移除导入导出工作流([#17940])(https://github.com/apache/dolphinscheduler/issues/17940)

## 3.5.0

* 为 `t_ds_schedules` 表新增 `missed_fire_policy` 字段。现有定时默认使用 `FIRE_ALL_MISSED`,以保持原有 Quartz `IgnoreMisfires` 行为。([#18464](https://github.com/apache/dolphinscheduler/pull/18464))

Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,14 @@

package org.apache.dolphinscheduler.api.dto;

import org.apache.dolphinscheduler.common.enums.ScheduleMissedFirePolicy;

import java.util.Date;

import lombok.Data;

import com.fasterxml.jackson.annotation.JsonIgnore;

/**
* schedule parameters
*/
Expand All @@ -31,6 +35,10 @@ public class ScheduleParam {
private Date endTime;
private String crontab;
private String timezoneId;
private ScheduleMissedFirePolicy missedFirePolicy = ScheduleMissedFirePolicy.FIRE_ALL_MISSED;

@JsonIgnore
private boolean missedFirePolicySet;

public ScheduleParam() {
}
Expand All @@ -42,6 +50,15 @@ public ScheduleParam(Date startTime, Date endTime, String timezoneId, String cro
this.crontab = crontab;
}

public void setMissedFirePolicy(ScheduleMissedFirePolicy missedFirePolicy) {
this.missedFirePolicy = missedFirePolicy;
this.missedFirePolicySet = true;
}

public boolean isMissedFirePolicySet() {
return missedFirePolicySet;
}

@Override
public String toString() {
return "ScheduleParam{"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,8 @@ public Schedule insertSchedule(User loginUser,
throw new ServiceException(Status.REQUEST_PARAMS_NOT_VALID_ERROR, scheduleParam.getCrontab());
}
scheduleObj.setCrontab(scheduleParam.getCrontab());
validateMissedFirePolicy(scheduleParam);
scheduleObj.setMissedFirePolicy(scheduleParam.getMissedFirePolicy());
scheduleObj.setTimezoneId(scheduleParam.getTimezoneId());
scheduleObj.setWarningType(warningType);
scheduleObj.setWarningGroupId(warningGroupId);
Expand Down Expand Up @@ -557,6 +559,10 @@ private Schedule updateSchedule(Schedule schedule, WorkflowDefinition workflowDe
throw new ServiceException(Status.SCHEDULE_CRON_CHECK_FAILED, scheduleParam.getCrontab());
}
schedule.setCrontab(scheduleParam.getCrontab());
validateMissedFirePolicy(scheduleParam);
if (scheduleParam.isMissedFirePolicySet() && scheduleParam.getMissedFirePolicy() != null) {
schedule.setMissedFirePolicy(scheduleParam.getMissedFirePolicy());
}
schedule.setTimezoneId(scheduleParam.getTimezoneId());
}

Expand Down Expand Up @@ -585,4 +591,11 @@ private Schedule updateSchedule(Schedule schedule, WorkflowDefinition workflowDe
return schedule;
}

private void validateMissedFirePolicy(ScheduleParam scheduleParam) {
if (scheduleParam.isMissedFirePolicySet() && scheduleParam.getMissedFirePolicy() == null) {
log.warn("Schedule missed fire policy is invalid.");
throw new ServiceException(Status.REQUEST_PARAMS_NOT_VALID_ERROR, "missedFirePolicy");
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import org.apache.dolphinscheduler.common.enums.FailureStrategy;
import org.apache.dolphinscheduler.common.enums.Priority;
import org.apache.dolphinscheduler.common.enums.ReleaseState;
import org.apache.dolphinscheduler.common.enums.ScheduleMissedFirePolicy;
import org.apache.dolphinscheduler.common.enums.WarningType;
import org.apache.dolphinscheduler.common.utils.DateUtils;
import org.apache.dolphinscheduler.dao.entity.Schedule;
Expand Down Expand Up @@ -54,6 +55,8 @@ public class ScheduleVO {

private String crontab;

private ScheduleMissedFirePolicy missedFirePolicy;

private FailureStrategy failureStrategy;

private WarningType warningType;
Expand Down Expand Up @@ -83,6 +86,7 @@ public class ScheduleVO {
public ScheduleVO(Schedule schedule) {
this.setId(schedule.getId());
this.setCrontab(schedule.getCrontab());
this.setMissedFirePolicy(schedule.getMissedFirePolicy());
this.setProjectName(schedule.getProjectName());
this.setUserName(schedule.getUserName());
this.setWorkerGroup(schedule.getWorkerGroup());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,24 +17,35 @@

package org.apache.dolphinscheduler.api.service;

import org.apache.dolphinscheduler.api.dto.ScheduleParam;
import org.apache.dolphinscheduler.api.enums.Status;
import org.apache.dolphinscheduler.api.exceptions.ServiceException;
import org.apache.dolphinscheduler.api.service.impl.SchedulerServiceImpl;
import org.apache.dolphinscheduler.api.validator.TenantExistValidator;
import org.apache.dolphinscheduler.common.enums.FailureStrategy;
import org.apache.dolphinscheduler.common.enums.Priority;
import org.apache.dolphinscheduler.common.enums.ReleaseState;
import org.apache.dolphinscheduler.common.enums.ScheduleMissedFirePolicy;
import org.apache.dolphinscheduler.common.enums.WarningType;
import org.apache.dolphinscheduler.common.utils.JSONUtils;
import org.apache.dolphinscheduler.dao.entity.Project;
import org.apache.dolphinscheduler.dao.entity.Schedule;
import org.apache.dolphinscheduler.dao.entity.User;
import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition;
import org.apache.dolphinscheduler.dao.repository.ProjectDao;
import org.apache.dolphinscheduler.dao.repository.ScheduleDao;
import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionDao;
import org.apache.dolphinscheduler.scheduler.api.SchedulerApi;

import java.util.Optional;

import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.EnumSource;
import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.Mockito;
Expand All @@ -61,6 +72,15 @@ public class SchedulerServiceTest extends BaseServiceTestTool {
@Mock
private ProjectService projectService;

@Mock
private ExecutorService executorService;

@Mock
private TenantExistValidator tenantExistValidator;

@Mock
private SchedulerApi schedulerApi;

protected static User user;
protected Exception exception;
private static final String userName = "userName";
Expand All @@ -80,6 +100,173 @@ public void setUp() {
user.setId(userId);
}

@Test
public void testScheduleParamMissedFirePolicyPresence() {
String scheduleWithoutPolicy = "{\"startTime\":\"2019-12-16 00:00:00\","
+ "\"endTime\":\"2019-12-17 00:00:00\",\"crontab\":\"0 0 6 * * ? *\"}";
String scheduleWithPolicy = "{\"startTime\":\"2019-12-16 00:00:00\","
+ "\"endTime\":\"2019-12-17 00:00:00\",\"crontab\":\"0 0 6 * * ? *\","
+ "\"missedFirePolicy\":\"SKIP_MISSED\"}";

ScheduleParam withoutPolicy = JSONUtils.parseObject(scheduleWithoutPolicy, ScheduleParam.class);
ScheduleParam withPolicy = JSONUtils.parseObject(scheduleWithPolicy, ScheduleParam.class);

Assertions.assertEquals(ScheduleMissedFirePolicy.FIRE_ALL_MISSED, withoutPolicy.getMissedFirePolicy());
Assertions.assertFalse(withoutPolicy.isMissedFirePolicySet());
Assertions.assertEquals(ScheduleMissedFirePolicy.SKIP_MISSED, withPolicy.getMissedFirePolicy());
Assertions.assertTrue(withPolicy.isMissedFirePolicySet());
}

@ParameterizedTest
@EnumSource(ScheduleMissedFirePolicy.class)
public void testInsertScheduleWithMissedFirePolicy(ScheduleMissedFirePolicy missedFirePolicy) {
Project project = this.getProject();
WorkflowDefinition workflowDefinition = this.getProcessDefinition();
Schedule insertedSchedule = new Schedule();
insertedSchedule.setId(scheduleId);
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project);
Mockito.when(scheduleDao.queryByWorkflowDefinitionCode(processDefinitionCode)).thenReturn(null);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(workflowDefinition));
Mockito.when(scheduleDao.queryById(Mockito.any())).thenReturn(insertedSchedule);

Schedule result = schedulerService.insertSchedule(
user, projectCode, processDefinitionCode, scheduleExpression(missedFirePolicy), WarningType.NONE, 0,
FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode);

ArgumentCaptor<Schedule> scheduleCaptor = ArgumentCaptor.forClass(Schedule.class);
Mockito.verify(scheduleDao).insert(scheduleCaptor.capture());
Assertions.assertEquals(missedFirePolicy, scheduleCaptor.getValue().getMissedFirePolicy());
Assertions.assertSame(insertedSchedule, result);
}

@Test
public void testInsertScheduleDefaultsMissedFirePolicy() {
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject());
Mockito.when(scheduleDao.queryByWorkflowDefinitionCode(processDefinitionCode)).thenReturn(null);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(this.getProcessDefinition()));
Mockito.when(scheduleDao.queryById(Mockito.anyInt())).thenReturn(new Schedule());

schedulerService.insertSchedule(
user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0,
FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode);

ArgumentCaptor<Schedule> scheduleCaptor = ArgumentCaptor.forClass(Schedule.class);
Mockito.verify(scheduleDao).insert(scheduleCaptor.capture());
Assertions.assertEquals(ScheduleMissedFirePolicy.FIRE_ALL_MISSED,
scheduleCaptor.getValue().getMissedFirePolicy());
}

@Test
public void testInsertScheduleRejectsExplicitNullMissedFirePolicy() {
assertInsertScheduleRejectsInvalidMissedFirePolicy("null");
}

@Test
public void testInsertScheduleRejectsUnknownMissedFirePolicy() {
assertInsertScheduleRejectsInvalidMissedFirePolicy("\"FIRE_ONCE_NWO\"");
}

@ParameterizedTest
@EnumSource(ScheduleMissedFirePolicy.class)
public void testUpdateScheduleWithMissedFirePolicy(ScheduleMissedFirePolicy missedFirePolicy) {
Schedule schedule = this.getSchedule();
schedule.setReleaseState(ReleaseState.OFFLINE);
WorkflowDefinition workflowDefinition = this.getProcessDefinition();
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject());
Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(workflowDefinition));

schedulerService.updateSchedule(
user, projectCode, scheduleId, scheduleExpression(missedFirePolicy), WarningType.NONE, 0,
FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode);

Assertions.assertEquals(missedFirePolicy, schedule.getMissedFirePolicy());
}

@Test
public void testUpdateSchedulePreservesMissedFirePolicyWhenOmitted() {
Schedule schedule = this.getSchedule();
schedule.setReleaseState(ReleaseState.OFFLINE);
schedule.setMissedFirePolicy(ScheduleMissedFirePolicy.SKIP_MISSED);
WorkflowDefinition workflowDefinition = this.getProcessDefinition();
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject());
Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(workflowDefinition));

schedulerService.updateSchedule(
user, projectCode, scheduleId, scheduleExpression(null), WarningType.NONE, 0,
FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode);

Assertions.assertEquals(ScheduleMissedFirePolicy.SKIP_MISSED, schedule.getMissedFirePolicy());
}

@Test
public void testUpdateScheduleRejectsExplicitNullMissedFirePolicy() {
assertUpdateScheduleRejectsInvalidMissedFirePolicy("null");
}

@Test
public void testUpdateScheduleRejectsUnknownMissedFirePolicy() {
assertUpdateScheduleRejectsInvalidMissedFirePolicy("\"FIRE_ONCE_NWO\"");
}

private void assertInsertScheduleRejectsInvalidMissedFirePolicy(String missedFirePolicyValue) {
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject());
Mockito.when(scheduleDao.queryByWorkflowDefinitionCode(processDefinitionCode)).thenReturn(null);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(this.getProcessDefinition()));

exception = Assertions.assertThrows(ServiceException.class,
() -> schedulerService.insertSchedule(
user, projectCode, processDefinitionCode,
scheduleExpressionWithPolicyValue(missedFirePolicyValue),
WarningType.NONE, 0, FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode",
environmentCode));

Assertions.assertEquals(Status.REQUEST_PARAMS_NOT_VALID_ERROR.getCode(),
((ServiceException) exception).getCode());
Mockito.verify(scheduleDao, Mockito.never()).insert(Mockito.any());
}

private void assertUpdateScheduleRejectsInvalidMissedFirePolicy(String missedFirePolicyValue) {
Schedule schedule = this.getSchedule();
schedule.setReleaseState(ReleaseState.OFFLINE);
schedule.setMissedFirePolicy(ScheduleMissedFirePolicy.SKIP_MISSED);
Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject());
Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule);
Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode))
.thenReturn(Optional.of(this.getProcessDefinition()));

exception = Assertions.assertThrows(ServiceException.class,
() -> schedulerService.updateSchedule(
user, projectCode, scheduleId, scheduleExpressionWithPolicyValue(missedFirePolicyValue),
WarningType.NONE, 0, FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode",
environmentCode));

Assertions.assertEquals(Status.REQUEST_PARAMS_NOT_VALID_ERROR.getCode(),
((ServiceException) exception).getCode());
Assertions.assertEquals(ScheduleMissedFirePolicy.SKIP_MISSED, schedule.getMissedFirePolicy());
Mockito.verify(scheduleDao, Mockito.never()).updateById(Mockito.any());
}

private String scheduleExpression(ScheduleMissedFirePolicy missedFirePolicy) {
String policy = missedFirePolicy == null ? "" : ",\"missedFirePolicy\":\"" + missedFirePolicy.name() + "\"";
return scheduleExpressionWithPolicy(policy);
}

private String scheduleExpressionWithPolicyValue(String missedFirePolicyValue) {
return scheduleExpressionWithPolicy(",\"missedFirePolicy\":" + missedFirePolicyValue);
}

private String scheduleExpressionWithPolicy(String policy) {
return "{\"startTime\":\"2019-12-16 00:00:00\",\"endTime\":\"2019-12-17 00:00:00\","
+ "\"crontab\":\"0 0 6 * * ? *\",\"timezoneId\":\"Asia/Shanghai\"" + policy + "}";
}

@Test
public void testDeleteSchedules() {
Schedule schedule = this.getSchedule();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
/*
* 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.common.enums;

import lombok.Getter;

import com.baomidou.mybatisplus.annotation.EnumValue;

@Getter
public enum ScheduleMissedFirePolicy {

SKIP_MISSED(0),
FIRE_ONCE_NOW(1),
FIRE_ALL_MISSED(2);

@EnumValue
private final int code;

ScheduleMissedFirePolicy(int code) {
this.code = code;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import org.apache.dolphinscheduler.common.enums.FailureStrategy;
import org.apache.dolphinscheduler.common.enums.Priority;
import org.apache.dolphinscheduler.common.enums.ReleaseState;
import org.apache.dolphinscheduler.common.enums.ScheduleMissedFirePolicy;
import org.apache.dolphinscheduler.common.enums.WarningType;

import java.util.Date;
Expand Down Expand Up @@ -67,6 +68,8 @@ public class Schedule {

private String crontab;

private ScheduleMissedFirePolicy missedFirePolicy;

private FailureStrategy failureStrategy;

private WarningType warningType;
Expand Down
Loading
Loading