From 2d3be6b36c5502711be113a0a568f533904613db Mon Sep 17 00:00:00 2001 From: Wenjun Ruan Date: Sat, 4 Jun 2022 16:39:33 +0800 Subject: [PATCH] Add dolphinscheduler-scheduler module (#10360) * Add dolphinscheduler-scheduler module --- dolphinscheduler-api/pom.xml | 23 +-- .../api/service/impl/ExecutorServiceImpl.java | 2 +- .../service/impl/SchedulerServiceImpl.java | 31 +--- .../api/service/SchedulerServiceTest.java | 8 +- .../dolphinscheduler/common/Constants.java | 20 --- dolphinscheduler-master/pom.xml | 4 + .../server/master/MasterServer.java | 6 +- .../master/runner/WorkflowExecuteThread.java | 3 +- .../dolphinscheduler-scheduler-api/pom.xml | 37 +++++ .../scheduler/api/SchedulerApi.java | 56 ++++++++ .../scheduler/api/SchedulerException.java | 33 ++--- .../dolphinscheduler-scheduler-quartz/pom.xml | 78 +++++++++++ .../scheduler/quartz/ProcessScheduleTask.java | 23 ++- .../scheduler/quartz/QuartzScheduler.java | 97 ++++++------- .../quartz/QuartzSchedulerConfiguration.java | 32 +++++ .../quartz/utils/QuartzTaskUtils.java | 59 ++++++++ dolphinscheduler-scheduler-plugin/pom.xml | 35 +++++ dolphinscheduler-service/pom.xml | 14 +- .../{quartz/cron => corn}/AbstractCycle.java | 2 +- .../{quartz/cron => corn}/CronUtils.java | 15 +- .../{quartz/cron => corn}/CycleFactory.java | 2 +- .../{quartz/cron => corn}/CycleLinks.java | 2 +- .../service/process/ProcessServiceImpl.java | 2 +- .../{quartz => }/cron/CronUtilsTest.java | 132 ++++++++++-------- .../service/process/ProcessServiceTest.java | 2 +- .../server/worker/WorkerServer.java | 2 - pom.xml | 12 ++ 27 files changed, 484 insertions(+), 248 deletions(-) create mode 100644 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/pom.xml create mode 100644 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerApi.java rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/QuartzExecutor.java => dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerException.java (57%) create mode 100644 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/pom.xml rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/ProcessScheduleJob.java => dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/ProcessScheduleTask.java (86%) rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/impl/QuartzExecutorImpl.java => dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzScheduler.java (67%) create mode 100644 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzSchedulerConfiguration.java create mode 100644 dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/utils/QuartzTaskUtils.java create mode 100644 dolphinscheduler-scheduler-plugin/pom.xml rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/{quartz/cron => corn}/AbstractCycle.java (99%) rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/{quartz/cron => corn}/CronUtils.java (94%) rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/{quartz/cron => corn}/CycleFactory.java (99%) rename dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/{quartz/cron => corn}/CycleLinks.java (97%) rename dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/{quartz => }/cron/CronUtilsTest.java (67%) diff --git a/dolphinscheduler-api/pom.xml b/dolphinscheduler-api/pom.xml index 7f7b1b2aae..6007ac1d9e 100644 --- a/dolphinscheduler-api/pom.xml +++ b/dolphinscheduler-api/pom.xml @@ -64,6 +64,10 @@ org.apache.dolphinscheduler dolphinscheduler-registry-zookeeper + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-quartz + org.codehaus.janino @@ -121,25 +125,6 @@ spring-context - - org.quartz-scheduler - quartz - - - com.mchange - c3p0 - - - com.mchange - mchange-commons-java - - - com.zaxxer - HikariCP-java7 - - - - io.springfox springfox-swagger2 diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java index 8c5dde91dc..9fafacb4f8 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java @@ -67,8 +67,8 @@ import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.enums.ExecutionStatus; import org.apache.dolphinscheduler.remote.command.StateEventChangeCommand; import org.apache.dolphinscheduler.remote.processor.StateEventCallbackService; +import org.apache.dolphinscheduler.service.corn.CronUtils; import org.apache.dolphinscheduler.service.process.ProcessService; -import org.apache.dolphinscheduler.service.quartz.cron.CronUtils; import org.apache.commons.beanutils.BeanUtils; import org.apache.commons.collections.CollectionUtils; diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java index 5416afcf52..7d4d1d6013 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java @@ -45,10 +45,9 @@ import org.apache.dolphinscheduler.dao.mapper.ProcessDefinitionMapper; import org.apache.dolphinscheduler.dao.mapper.ProcessTaskRelationMapper; import org.apache.dolphinscheduler.dao.mapper.ProjectMapper; import org.apache.dolphinscheduler.dao.mapper.ScheduleMapper; +import org.apache.dolphinscheduler.scheduler.api.SchedulerApi; +import org.apache.dolphinscheduler.service.corn.CronUtils; import org.apache.dolphinscheduler.service.process.ProcessService; -import org.apache.dolphinscheduler.service.quartz.ProcessScheduleJob; -import org.apache.dolphinscheduler.service.quartz.QuartzExecutor; -import org.apache.dolphinscheduler.service.quartz.cron.CronUtils; import org.apache.commons.lang3.StringUtils; @@ -60,12 +59,10 @@ import java.util.List; import java.util.Map; import org.quartz.CronExpression; -import org.quartz.JobKey; -import org.quartz.Scheduler; -import org.quartz.SchedulerException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.support.CronTrigger; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @@ -102,13 +99,11 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe private ProcessDefinitionMapper processDefinitionMapper; @Autowired - private Scheduler scheduler; + private SchedulerApi schedulerApi; @Autowired private ProcessTaskRelationMapper processTaskRelationMapper; - @Autowired - private QuartzExecutor quartzExecutor; /** * save schedule @@ -460,8 +455,7 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe public void setSchedule(int projectId, Schedule schedule) { logger.info("set schedule, project id: {}, scheduleId: {}", projectId, schedule.getId()); - - quartzExecutor.addJob(ProcessScheduleJob.class, projectId, schedule); + schedulerApi.insertOrUpdateScheduleTask(projectId, schedule); } /** @@ -474,20 +468,7 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe @Override public void deleteSchedule(int projectId, int scheduleId) { logger.info("delete schedules of project id:{}, schedule id:{}", projectId, scheduleId); - - String jobName = quartzExecutor.buildJobName(scheduleId); - String jobGroupName = quartzExecutor.buildJobGroupName(projectId); - - JobKey jobKey = new JobKey(jobName, jobGroupName); - try { - if (scheduler.checkExists(jobKey)) { - logger.info("Try to delete job: {}, group name: {},", jobName, jobGroupName); - scheduler.deleteJob(jobKey); - } - } catch (SchedulerException e) { - logger.error("Failed to delete job: {}", jobKey); - throw new ServiceException("Failed to delete job: " + jobKey); - } + schedulerApi.deleteScheduleTask(projectId, scheduleId); } /** diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java index 888b5b23eb..9d292f23fe 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java @@ -31,8 +31,9 @@ import org.apache.dolphinscheduler.dao.mapper.ProcessDefinitionMapper; import org.apache.dolphinscheduler.dao.mapper.ProcessTaskRelationMapper; import org.apache.dolphinscheduler.dao.mapper.ProjectMapper; import org.apache.dolphinscheduler.dao.mapper.ScheduleMapper; +import org.apache.dolphinscheduler.scheduler.api.SchedulerApi; +import org.apache.dolphinscheduler.scheduler.quartz.QuartzScheduler; import org.apache.dolphinscheduler.service.process.ProcessService; -import org.apache.dolphinscheduler.service.quartz.impl.QuartzExecutorImpl; import java.util.ArrayList; import java.util.HashMap; @@ -53,7 +54,6 @@ import org.powermock.modules.junit4.PowerMockRunner; * scheduler service test */ @RunWith(PowerMockRunner.class) -@PrepareForTest(QuartzExecutorImpl.class) public class SchedulerServiceTest { @InjectMocks @@ -80,8 +80,8 @@ public class SchedulerServiceTest { @Mock private ProjectServiceImpl projectService; - @InjectMocks - private QuartzExecutorImpl quartzExecutors; + @Mock + private SchedulerApi schedulerApi; @Before public void setUp() { diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/Constants.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/Constants.java index 49c56f897a..997ead7ef2 100644 --- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/Constants.java +++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/Constants.java @@ -494,26 +494,6 @@ public final class Constants { * underline "_" */ public static final String UNDERLINE = "_"; - /** - * quartz job prifix - */ - public static final String QUARTZ_JOB_PREFIX = "job"; - /** - * quartz job group prifix - */ - public static final String QUARTZ_JOB_GROUP_PREFIX = "jobgroup"; - /** - * projectId - */ - public static final String PROJECT_ID = "projectId"; - /** - * processId - */ - public static final String SCHEDULE_ID = "scheduleId"; - /** - * schedule - */ - public static final String SCHEDULE = "schedule"; /** * application regex */ diff --git a/dolphinscheduler-master/pom.xml b/dolphinscheduler-master/pom.xml index a1e179f63f..f21d281ccd 100644 --- a/dolphinscheduler-master/pom.xml +++ b/dolphinscheduler-master/pom.xml @@ -51,6 +51,10 @@ org.apache.dolphinscheduler dolphinscheduler-registry-zookeeper + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-quartz + org.apache.dolphinscheduler diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/MasterServer.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/MasterServer.java index ad4d02e2e9..a7145fc049 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/MasterServer.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/MasterServer.java @@ -23,6 +23,7 @@ import org.apache.dolphinscheduler.common.thread.Stopper; import org.apache.dolphinscheduler.remote.NettyRemotingServer; import org.apache.dolphinscheduler.remote.command.CommandType; import org.apache.dolphinscheduler.remote.config.NettyServerConfig; +import org.apache.dolphinscheduler.scheduler.api.SchedulerApi; import org.apache.dolphinscheduler.server.log.LoggerRequestProcessor; import org.apache.dolphinscheduler.server.master.config.MasterConfig; import org.apache.dolphinscheduler.server.master.processor.CacheProcessor; @@ -77,7 +78,7 @@ public class MasterServer implements IStoppable { private MasterSchedulerService masterSchedulerService; @Autowired - private Scheduler scheduler; + private SchedulerApi schedulerApi; @Autowired private TaskExecuteRunningProcessor taskExecuteRunningProcessor; @@ -154,7 +155,7 @@ public class MasterServer implements IStoppable { this.eventExecuteService.start(); this.failoverExecuteThread.start(); - this.scheduler.start(); + this.schedulerApi.start(); Runtime.getRuntime().addShutdownHook(new Thread(() -> { if (Stopper.isRunning()) { @@ -188,6 +189,7 @@ public class MasterServer implements IStoppable { logger.warn("thread sleep exception ", e); } // close + this.schedulerApi.close(); this.masterSchedulerService.close(); this.nettyRemotingServer.close(); this.masterRegistryClient.closeRegistry(); diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java index a54e2083ed..cb3f4fc362 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java @@ -17,7 +17,6 @@ package org.apache.dolphinscheduler.server.master.runner; -import net.bytebuddy.implementation.bytecode.Throw; import static org.apache.dolphinscheduler.common.Constants.CMDPARAM_COMPLEMENT_DATA_END_DATE; import static org.apache.dolphinscheduler.common.Constants.CMDPARAM_COMPLEMENT_DATA_START_DATE; import static org.apache.dolphinscheduler.common.Constants.CMD_PARAM_RECOVERY_START_NODE_STRING; @@ -71,8 +70,8 @@ import org.apache.dolphinscheduler.server.master.runner.task.ITaskProcessor; import org.apache.dolphinscheduler.server.master.runner.task.TaskAction; import org.apache.dolphinscheduler.server.master.runner.task.TaskProcessorFactory; import org.apache.dolphinscheduler.service.alert.ProcessAlertManager; +import org.apache.dolphinscheduler.service.corn.CronUtils; import org.apache.dolphinscheduler.service.process.ProcessService; -import org.apache.dolphinscheduler.service.quartz.cron.CronUtils; import org.apache.dolphinscheduler.service.queue.PeerTaskInstancePriorityQueue; import org.apache.commons.collections.CollectionUtils; diff --git a/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/pom.xml b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/pom.xml new file mode 100644 index 0000000000..230c592827 --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/pom.xml @@ -0,0 +1,37 @@ + + + + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-plugin + dev-SNAPSHOT + + 4.0.0 + + dolphinscheduler-scheduler-api + + + + org.apache.dolphinscheduler + dolphinscheduler-dao + + + + \ No newline at end of file diff --git a/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerApi.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerApi.java new file mode 100644 index 0000000000..36f04c22bb --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerApi.java @@ -0,0 +1,56 @@ +/* + * 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.scheduler.api; + +import org.apache.dolphinscheduler.dao.entity.Schedule; + +/** + * This is the interface for scheduler, contains methods to operate schedule task. + */ +public interface SchedulerApi extends AutoCloseable{ + + /** + * Start the scheduler, if not start, the scheduler will not execute task. + * + * @throws SchedulerException if start failed. + */ + void start() throws SchedulerException; + + /** + * @param projectId project id, the schedule task belongs to. + * @param schedule schedule metadata. + * @throws SchedulerException if insert/update failed. + */ + void insertOrUpdateScheduleTask(int projectId, Schedule schedule) throws SchedulerException; + + /** + * Delete a schedule task. + * + * @param projectId project id, the schedule task belongs to. + * @param scheduleId schedule id. + * @throws SchedulerException if delete failed. + */ + void deleteScheduleTask(int projectId, int scheduleId) throws SchedulerException; + + /** + * Close the scheduler and release the resource. + * + * @throws SchedulerException if close failed. + */ + void close() throws Exception; +} diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/QuartzExecutor.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerException.java similarity index 57% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/QuartzExecutor.java rename to dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerException.java index e0e76e5652..c81f114204 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/QuartzExecutor.java +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-api/src/main/java/org/apache/dolphinscheduler/scheduler/api/SchedulerException.java @@ -15,33 +15,16 @@ * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz; +package org.apache.dolphinscheduler.scheduler.api; -import org.apache.dolphinscheduler.dao.entity.Schedule; +public class SchedulerException extends RuntimeException { -import java.util.Map; + public SchedulerException(String message) { + super(message); + } -import org.quartz.Job; + public SchedulerException(String message, Throwable cause) { + super(message, cause); + } -public interface QuartzExecutor { - - /** - * build job name - */ - String buildJobName(int scheduleId); - - /** - * build job group name - */ - String buildJobGroupName(int projectId); - - /** - * build data map of job detail - */ - Map buildDataMap(int projectId, Schedule schedule); - - /** - * add job to quartz - */ - void addJob(Class clazz, int projectId, final Schedule schedule); } diff --git a/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/pom.xml b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/pom.xml new file mode 100644 index 0000000000..d3afa0e70f --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/pom.xml @@ -0,0 +1,78 @@ + + + + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-plugin + dev-SNAPSHOT + + 4.0.0 + + dolphinscheduler-scheduler-quartz + + + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-api + + + org.apache.dolphinscheduler + dolphinscheduler-service + + + org.apache.dolphinscheduler + dolphinscheduler-meter + + + + org.springframework.boot + spring-boot-starter-quartz + + + org.springframework.boot + spring-boot-starter + + + + + org.quartz-scheduler + quartz + + + com.mchange + c3p0 + + + com.mchange + mchange-commons-java + + + com.zaxxer + HikariCP-java7 + + + + + com.cronutils + cron-utils + + + + \ No newline at end of file diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/ProcessScheduleJob.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/ProcessScheduleTask.java similarity index 86% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/ProcessScheduleJob.java rename to dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/ProcessScheduleTask.java index ec5ed36f45..d65fae32b8 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/ProcessScheduleJob.java +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/ProcessScheduleTask.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz; +package org.apache.dolphinscheduler.scheduler.quartz; import org.apache.dolphinscheduler.common.Constants; import org.apache.dolphinscheduler.common.enums.CommandType; @@ -23,8 +23,8 @@ import org.apache.dolphinscheduler.common.enums.ReleaseState; import org.apache.dolphinscheduler.dao.entity.Command; import org.apache.dolphinscheduler.dao.entity.ProcessDefinition; import org.apache.dolphinscheduler.dao.entity.Schedule; +import org.apache.dolphinscheduler.scheduler.quartz.utils.QuartzTaskUtils; import org.apache.dolphinscheduler.service.process.ProcessService; -import org.apache.dolphinscheduler.service.quartz.impl.QuartzExecutorImpl; import java.util.Date; @@ -41,23 +41,21 @@ import org.springframework.util.StringUtils; import io.micrometer.core.annotation.Counted; import io.micrometer.core.annotation.Timed; -public class ProcessScheduleJob extends QuartzJobBean { - private static final Logger logger = LoggerFactory.getLogger(ProcessScheduleJob.class); +public class ProcessScheduleTask extends QuartzJobBean { + + private static final Logger logger = LoggerFactory.getLogger(ProcessScheduleTask.class); @Autowired private ProcessService processService; - @Autowired - private QuartzExecutor quartzExecutor; - @Counted(value = "quartz_job_executed") @Timed(value = "quartz_job_execution", percentiles = {0.5, 0.75, 0.95, 0.99}, histogram = true) @Override protected void executeInternal(JobExecutionContext context) { JobDataMap dataMap = context.getJobDetail().getJobDataMap(); - int projectId = dataMap.getInt(Constants.PROJECT_ID); - int scheduleId = dataMap.getInt(Constants.SCHEDULE_ID); + int projectId = dataMap.getInt(QuartzTaskUtils.PROJECT_ID); + int scheduleId = dataMap.getInt(QuartzTaskUtils.SCHEDULE_ID); Date scheduledFireTime = context.getScheduledFireTime(); @@ -100,13 +98,10 @@ public class ProcessScheduleJob extends QuartzJobBean { private void deleteJob(JobExecutionContext context, int projectId, int scheduleId) { final Scheduler scheduler = context.getScheduler(); - String jobName = quartzExecutor.buildJobName(scheduleId); - String jobGroupName = quartzExecutor.buildJobGroupName(projectId); - - JobKey jobKey = new JobKey(jobName, jobGroupName); + JobKey jobKey = QuartzTaskUtils.getJobKey(scheduleId, projectId); try { if (scheduler.checkExists(jobKey)) { - logger.info("Try to delete job: {}, group name: {},", jobName, jobGroupName); + logger.info("Try to delete job: {}, projectId: {}, schedulerId", projectId, scheduleId); scheduler.deleteJob(jobKey); } } catch (Exception e) { diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/impl/QuartzExecutorImpl.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzScheduler.java similarity index 67% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/impl/QuartzExecutorImpl.java rename to dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzScheduler.java index c31de54643..6b76968401 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/impl/QuartzExecutorImpl.java +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzScheduler.java @@ -15,33 +15,25 @@ * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.impl; +package org.apache.dolphinscheduler.scheduler.quartz; -import static org.apache.dolphinscheduler.common.Constants.PROJECT_ID; -import static org.apache.dolphinscheduler.common.Constants.QUARTZ_JOB_GROUP_PREFIX; -import static org.apache.dolphinscheduler.common.Constants.QUARTZ_JOB_PREFIX; -import static org.apache.dolphinscheduler.common.Constants.SCHEDULE; -import static org.apache.dolphinscheduler.common.Constants.SCHEDULE_ID; -import static org.apache.dolphinscheduler.common.Constants.UNDERLINE; import static org.quartz.CronScheduleBuilder.cronSchedule; import static org.quartz.JobBuilder.newJob; import static org.quartz.TriggerBuilder.newTrigger; import org.apache.dolphinscheduler.common.utils.DateUtils; -import org.apache.dolphinscheduler.common.utils.JSONUtils; import org.apache.dolphinscheduler.dao.entity.Schedule; -import org.apache.dolphinscheduler.service.exceptions.ServiceException; -import org.apache.dolphinscheduler.service.quartz.QuartzExecutor; +import org.apache.dolphinscheduler.scheduler.api.SchedulerApi; +import org.apache.dolphinscheduler.scheduler.api.SchedulerException; +import org.apache.dolphinscheduler.scheduler.quartz.utils.QuartzTaskUtils; import java.util.Date; -import java.util.HashMap; import java.util.Map; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; -import com.google.common.base.Strings; + import org.quartz.CronTrigger; -import org.quartz.Job; import org.quartz.JobDetail; import org.quartz.JobKey; import org.quartz.Scheduler; @@ -49,30 +41,31 @@ import org.quartz.TriggerKey; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.stereotype.Service; -@Service -public class QuartzExecutorImpl implements QuartzExecutor { - private static final Logger logger = LoggerFactory.getLogger(QuartzExecutorImpl.class); +import com.google.common.base.Strings; + +public class QuartzScheduler implements SchedulerApi { + + private static final Logger logger = LoggerFactory.getLogger(QuartzScheduler.class); @Autowired private Scheduler scheduler; private final ReadWriteLock lock = new ReentrantReadWriteLock(); - /** - * add task trigger , if this task already exists, return this task with updated trigger - * - * @param clazz job class name - * @param projectId projectId - * @param schedule schedule - */ @Override - public void addJob(Class clazz, int projectId, final Schedule schedule) { - String jobName = this.buildJobName(schedule.getId()); - String jobGroupName = this.buildJobGroupName(projectId); + public void start() throws SchedulerException { + try { + scheduler.start(); + } catch (Exception e) { + throw new SchedulerException("Failed to start quartz scheduler ", e); + } + } - Map jobDataMap = this.buildDataMap(projectId, schedule); + @Override + public void insertOrUpdateScheduleTask(int projectId, Schedule schedule) throws SchedulerException { + JobKey jobKey = QuartzTaskUtils.getJobKey(schedule.getId(), projectId); + Map jobDataMap = QuartzTaskUtils.buildDataMap(projectId, schedule); String cronExpression = schedule.getCrontab(); String timezoneId = schedule.getTimezoneId(); @@ -89,7 +82,6 @@ public class QuartzExecutorImpl implements QuartzExecutor { lock.writeLock().lock(); try { - JobKey jobKey = new JobKey(jobName, jobGroupName); JobDetail jobDetail; //add a task (if this task already exists, return this task directly) if (scheduler.checkExists(jobKey)) { @@ -97,17 +89,16 @@ public class QuartzExecutorImpl implements QuartzExecutor { jobDetail = scheduler.getJobDetail(jobKey); jobDetail.getJobDataMap().putAll(jobDataMap); } else { - jobDetail = newJob(clazz).withIdentity(jobKey).build(); + jobDetail = newJob(ProcessScheduleTask.class).withIdentity(jobKey).build(); jobDetail.getJobDataMap().putAll(jobDataMap); scheduler.addJob(jobDetail, false, true); - logger.info("Add job, job name: {}, group name: {}", - jobName, jobGroupName); + logger.info("Add job, job name: {}, group name: {}", jobKey.getName(), jobKey.getGroup()); } - TriggerKey triggerKey = new TriggerKey(jobName, jobGroupName); + TriggerKey triggerKey = new TriggerKey(jobKey.getName(), jobKey.getGroup()); /* * Instructs the Scheduler that upon a mis-fire * situation, the CronTrigger wants to have it's @@ -135,39 +126,43 @@ public class QuartzExecutorImpl implements QuartzExecutor { // reschedule job trigger scheduler.rescheduleJob(triggerKey, cronTrigger); logger.info("reschedule job trigger, triggerName: {}, triggerGroupName: {}, cronExpression: {}, startDate: {}, endDate: {}", - jobName, jobGroupName, cronExpression, startDate, endDate); + triggerKey.getName(), triggerKey.getGroup(), cronExpression, startDate, endDate); } } else { scheduler.scheduleJob(cronTrigger); logger.info("schedule job trigger, triggerName: {}, triggerGroupName: {}, cronExpression: {}, startDate: {}, endDate: {}", - jobName, jobGroupName, cronExpression, startDate, endDate); + triggerKey.getName(), triggerKey.getGroup(), cronExpression, startDate, endDate); } } catch (Exception e) { - throw new ServiceException("add job failed", e); + logger.error("Failed to add scheduler task, projectId: {}, scheduler: {}", projectId, schedule, e); + throw new SchedulerException("Add schedule job failed", e); } finally { lock.writeLock().unlock(); } } @Override - public String buildJobName(int processId) { - return QUARTZ_JOB_PREFIX + UNDERLINE + processId; + public void deleteScheduleTask(int projectId, int scheduleId) throws SchedulerException { + JobKey jobKey = QuartzTaskUtils.getJobKey(scheduleId, projectId); + try { + if (scheduler.checkExists(jobKey)) { + logger.info("Try to delete scheduler task, projectId: {}, schedulerId: {}", projectId, scheduleId); + scheduler.deleteJob(jobKey); + } + } catch (Exception e) { + logger.error("Failed to delete scheduler task, projectId: {}, schedulerId: {}", projectId, scheduleId, e); + throw new SchedulerException("Failed to delete scheduler task"); + } } @Override - public String buildJobGroupName(int projectId) { - return QUARTZ_JOB_GROUP_PREFIX + UNDERLINE + projectId; + public void close() { + // nothing to do + try { + scheduler.shutdown(); + } catch (org.quartz.SchedulerException e) { + throw new SchedulerException("Failed to shutdown scheduler", e); + } } - - @Override - public Map buildDataMap(int projectId, Schedule schedule) { - Map dataMap = new HashMap<>(8); - dataMap.put(PROJECT_ID, projectId); - dataMap.put(SCHEDULE_ID, schedule.getId()); - dataMap.put(SCHEDULE, JSONUtils.toJsonString(schedule)); - - return dataMap; - } - } diff --git a/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzSchedulerConfiguration.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzSchedulerConfiguration.java new file mode 100644 index 0000000000..123e4cde50 --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/QuartzSchedulerConfiguration.java @@ -0,0 +1,32 @@ +/* + * 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.scheduler.quartz; + +import org.apache.dolphinscheduler.scheduler.api.SchedulerApi; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +@Configuration +public class QuartzSchedulerConfiguration { + + @Bean + public SchedulerApi schedulerApi() { + return new QuartzScheduler(); + } +} diff --git a/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/utils/QuartzTaskUtils.java b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/utils/QuartzTaskUtils.java new file mode 100644 index 0000000000..e4f3471399 --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/dolphinscheduler-scheduler-quartz/src/main/java/org/apache/dolphinscheduler/scheduler/quartz/utils/QuartzTaskUtils.java @@ -0,0 +1,59 @@ +/* + * 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.scheduler.quartz.utils; + +import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.dao.entity.Schedule; + +import java.util.HashMap; +import java.util.Map; + +import org.quartz.JobKey; + +public final class QuartzTaskUtils { + + public static final String QUARTZ_JOB_PREFIX = "job"; + public static final String QUARTZ_JOB_GROUP_PREFIX = "jobgroup"; + public static final String UNDERLINE = "_"; + public static final String PROJECT_ID = "projectId"; + public static final String SCHEDULE_ID = "scheduleId"; + public static final String SCHEDULE = "schedule"; + + /** + * @param schedulerId scheduler id + * @return quartz job name + */ + public static JobKey getJobKey(int schedulerId, int projectId) { + String jobName = QUARTZ_JOB_PREFIX + UNDERLINE + schedulerId; + String jobGroup = QUARTZ_JOB_GROUP_PREFIX + UNDERLINE + projectId; + return new JobKey(jobName, jobGroup); + } + + /** + * create quartz job data, include projectId and scheduleId, schedule. + */ + public static Map buildDataMap(int projectId, Schedule schedule) { + Map dataMap = new HashMap<>(8); + dataMap.put(PROJECT_ID, projectId); + dataMap.put(SCHEDULE_ID, schedule.getId()); + dataMap.put(SCHEDULE, JSONUtils.toJsonString(schedule)); + + return dataMap; + } + +} diff --git a/dolphinscheduler-scheduler-plugin/pom.xml b/dolphinscheduler-scheduler-plugin/pom.xml new file mode 100644 index 0000000000..15e2f4f898 --- /dev/null +++ b/dolphinscheduler-scheduler-plugin/pom.xml @@ -0,0 +1,35 @@ + + + + + dolphinscheduler + org.apache.dolphinscheduler + dev-SNAPSHOT + + 4.0.0 + pom + + dolphinscheduler-scheduler-plugin + + + dolphinscheduler-scheduler-api + dolphinscheduler-scheduler-quartz + + \ No newline at end of file diff --git a/dolphinscheduler-service/pom.xml b/dolphinscheduler-service/pom.xml index 67c932e96a..bd05a15f45 100644 --- a/dolphinscheduler-service/pom.xml +++ b/dolphinscheduler-service/pom.xml @@ -52,14 +52,8 @@ - org.springframework.boot - spring-boot-starter-quartz - - - org.springframework.boot - spring-boot-starter - - + com.cronutils + cron-utils org.quartz-scheduler @@ -79,10 +73,6 @@ - - com.cronutils - cron-utils - io.micrometer diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/AbstractCycle.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/AbstractCycle.java similarity index 99% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/AbstractCycle.java rename to dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/AbstractCycle.java index b00f1476ab..8e69c458c4 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/AbstractCycle.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/AbstractCycle.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.cron; +package org.apache.dolphinscheduler.service.corn; import org.apache.dolphinscheduler.common.enums.CycleEnum; diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtils.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CronUtils.java similarity index 94% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtils.java rename to dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CronUtils.java index 49810cd1f4..f1ed739da8 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtils.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CronUtils.java @@ -15,14 +15,14 @@ * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.cron; +package org.apache.dolphinscheduler.service.corn; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.day; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.hour; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.min; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.month; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.week; -import static org.apache.dolphinscheduler.service.quartz.cron.CycleFactory.year; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.day; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.hour; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.min; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.month; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.week; +import static org.apache.dolphinscheduler.service.corn.CycleFactory.year; import static com.cronutils.model.CronType.QUARTZ; @@ -51,6 +51,7 @@ import com.cronutils.model.definition.CronDefinitionBuilder; import com.cronutils.parser.CronParser; /** + * // todo: this utils is heavy, it rely on quartz and corn-utils. * cron utils */ public class CronUtils { diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleFactory.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleFactory.java similarity index 99% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleFactory.java rename to dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleFactory.java index 9f931d20bb..1a133ee7ea 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleFactory.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleFactory.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.cron; +package org.apache.dolphinscheduler.service.corn; import com.cronutils.model.Cron; import com.cronutils.model.field.expression.Always; diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleLinks.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleLinks.java similarity index 97% rename from dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleLinks.java rename to dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleLinks.java index 9f01b18868..7cc4a87f07 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/quartz/cron/CycleLinks.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/corn/CycleLinks.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.cron; +package org.apache.dolphinscheduler.service.corn; import com.cronutils.model.Cron; import org.apache.dolphinscheduler.common.enums.CycleEnum; diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java index f14a20eea2..7690672f77 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java @@ -128,9 +128,9 @@ import org.apache.dolphinscheduler.remote.command.TaskEventChangeCommand; import org.apache.dolphinscheduler.remote.processor.StateEventCallbackService; import org.apache.dolphinscheduler.remote.utils.Host; import org.apache.dolphinscheduler.service.bean.SpringApplicationContext; +import org.apache.dolphinscheduler.service.corn.CronUtils; import org.apache.dolphinscheduler.service.exceptions.ServiceException; import org.apache.dolphinscheduler.service.log.LogClientService; -import org.apache.dolphinscheduler.service.quartz.cron.CronUtils; import org.apache.dolphinscheduler.service.task.TaskPluginManager; import org.apache.dolphinscheduler.spi.enums.ResourceType; diff --git a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtilsTest.java b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/cron/CronUtilsTest.java similarity index 67% rename from dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtilsTest.java rename to dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/cron/CronUtilsTest.java index 4fbcd8f9c0..b53062e975 100644 --- a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/quartz/cron/CronUtilsTest.java +++ b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/cron/CronUtilsTest.java @@ -14,7 +14,25 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.dolphinscheduler.service.quartz.cron; + +package org.apache.dolphinscheduler.service.cron; + +import static com.cronutils.model.field.expression.FieldExpressionFactory.always; +import static com.cronutils.model.field.expression.FieldExpressionFactory.every; +import static com.cronutils.model.field.expression.FieldExpressionFactory.on; +import static com.cronutils.model.field.expression.FieldExpressionFactory.questionMark; + +import org.apache.dolphinscheduler.common.enums.CycleEnum; +import org.apache.dolphinscheduler.common.utils.DateUtils; +import org.apache.dolphinscheduler.service.corn.CronUtils; + +import java.text.ParseException; +import java.util.Date; + +import org.junit.Assert; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import com.cronutils.builder.CronBuilder; import com.cronutils.model.Cron; @@ -22,18 +40,12 @@ import com.cronutils.model.CronType; import com.cronutils.model.definition.CronDefinitionBuilder; import com.cronutils.model.field.CronField; import com.cronutils.model.field.CronFieldName; -import com.cronutils.model.field.expression.*; -import org.apache.dolphinscheduler.common.enums.CycleEnum; -import org.apache.dolphinscheduler.common.utils.DateUtils; -import org.junit.Assert; -import org.junit.Test; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.text.ParseException; -import java.util.Date; - -import static com.cronutils.model.field.expression.FieldExpressionFactory.*; +import com.cronutils.model.field.expression.Always; +import com.cronutils.model.field.expression.And; +import com.cronutils.model.field.expression.Between; +import com.cronutils.model.field.expression.Every; +import com.cronutils.model.field.expression.On; +import com.cronutils.model.field.expression.QuestionMark; /** * CronUtilsTest @@ -66,6 +78,7 @@ public class CronUtilsTest { /** * cron parse test + * * @throws ParseException if error throws ParseException */ @Test @@ -83,6 +96,7 @@ public class CronUtilsTest { /** * schedule type test + * * @throws ParseException if error throws ParseException */ @Test @@ -95,18 +109,18 @@ public class CronUtilsTest { CycleEnum cycleEnum3 = CronUtils.getMiniCycle(CronUtils.parse2Cron("0 * * * * ? *")); Assert.assertEquals("MINUTE", cycleEnum3.name()); - + CycleEnum cycleEnum4 = CronUtils.getMaxCycle(CronUtils.parse2Cron("0 0 7 * 1 ? *")); Assert.assertEquals("YEAR", cycleEnum4.name()); cycleEnum4 = CronUtils.getMiniCycle(CronUtils.parse2Cron("0 0 7 * 1 ? *")); Assert.assertEquals("DAY", cycleEnum4.name()); - + CycleEnum cycleEnum5 = CronUtils.getMaxCycle(CronUtils.parse2Cron("0 0 7 * 1/1 ? *")); Assert.assertEquals("MONTH", cycleEnum5.name()); - + CycleEnum cycleEnum6 = CronUtils.getMaxCycle(CronUtils.parse2Cron("0 0 7 * 1-2 ? *")); Assert.assertEquals("YEAR", cycleEnum6.name()); - + CycleEnum cycleEnum7 = CronUtils.getMaxCycle(CronUtils.parse2Cron("0 0 7 * 1,2 ? *")); Assert.assertEquals("YEAR", cycleEnum7.name()); } @@ -115,7 +129,7 @@ public class CronUtilsTest { * test */ @Test - public void test2(){ + public void test2() { Cron cron1 = CronBuilder.cron(CronDefinitionBuilder.instanceDefinitionFor(CronType.QUARTZ)) .withYear(always()) .withDoW(questionMark()) @@ -126,62 +140,62 @@ public class CronUtilsTest { .withSecond(on(0)) .instance(); // minute cycle - String[] cronArayy = new String[]{"* * * * * ? *","* 0 * * * ? *", - "* 5 * * 3/5 ? *","0 0 * * * ? *", "0 0 7 * 1 ? *", "0 0 7 * 1/1 ? *", "0 0 7 * 1-2 ? *" , "0 0 7 * 1,2 ? *"}; - for(String minCrontab:cronArayy){ + String[] cronArayy = new String[] {"* * * * * ? *", "* 0 * * * ? *", + "* 5 * * 3/5 ? *", "0 0 * * * ? *", "0 0 7 * 1 ? *", "0 0 7 * 1/1 ? *", "0 0 7 * 1-2 ? *", "0 0 7 * 1,2 ? *"}; + for (String minCrontab : cronArayy) { if (!org.quartz.CronExpression.isValidExpression(minCrontab)) { - throw new RuntimeException(minCrontab+" verify failure, cron expression not valid"); + throw new RuntimeException(minCrontab + " verify failure, cron expression not valid"); } Cron cron = CronUtils.parse2Cron(minCrontab); CronField minField = cron.retrieve(CronFieldName.MINUTE); - logger.info("minField instanceof Between:"+(minField.getExpression() instanceof Between)); - logger.info("minField instanceof Every:"+(minField.getExpression() instanceof Every)); + logger.info("minField instanceof Between:" + (minField.getExpression() instanceof Between)); + logger.info("minField instanceof Every:" + (minField.getExpression() instanceof Every)); logger.info("minField instanceof Always:" + (minField.getExpression() instanceof Always)); - logger.info("minField instanceof On:"+(minField.getExpression() instanceof On)); - logger.info("minField instanceof And:"+(minField.getExpression() instanceof And)); + logger.info("minField instanceof On:" + (minField.getExpression() instanceof On)); + logger.info("minField instanceof And:" + (minField.getExpression() instanceof And)); CronField hourField = cron.retrieve(CronFieldName.HOUR); - logger.info("hourField instanceof Between:"+(hourField.getExpression() instanceof Between)); - logger.info("hourField instanceof Always:"+(hourField.getExpression() instanceof Always)); - logger.info("hourField instanceof Every:"+(hourField.getExpression() instanceof Every)); - logger.info("hourField instanceof On:"+(hourField.getExpression() instanceof On)); - logger.info("hourField instanceof And:"+(hourField.getExpression() instanceof And)); + logger.info("hourField instanceof Between:" + (hourField.getExpression() instanceof Between)); + logger.info("hourField instanceof Always:" + (hourField.getExpression() instanceof Always)); + logger.info("hourField instanceof Every:" + (hourField.getExpression() instanceof Every)); + logger.info("hourField instanceof On:" + (hourField.getExpression() instanceof On)); + logger.info("hourField instanceof And:" + (hourField.getExpression() instanceof And)); CronField dayOfMonthField = cron.retrieve(CronFieldName.DAY_OF_MONTH); - logger.info("dayOfMonthField instanceof Between:"+(dayOfMonthField.getExpression() instanceof Between)); - logger.info("dayOfMonthField instanceof Always:"+(dayOfMonthField.getExpression() instanceof Always)); - logger.info("dayOfMonthField instanceof Every:"+(dayOfMonthField.getExpression() instanceof Every)); - logger.info("dayOfMonthField instanceof On:"+(dayOfMonthField.getExpression() instanceof On)); - logger.info("dayOfMonthField instanceof And:"+(dayOfMonthField.getExpression() instanceof And)); - logger.info("dayOfMonthField instanceof QuestionMark:"+(dayOfMonthField.getExpression() instanceof QuestionMark)); + logger.info("dayOfMonthField instanceof Between:" + (dayOfMonthField.getExpression() instanceof Between)); + logger.info("dayOfMonthField instanceof Always:" + (dayOfMonthField.getExpression() instanceof Always)); + logger.info("dayOfMonthField instanceof Every:" + (dayOfMonthField.getExpression() instanceof Every)); + logger.info("dayOfMonthField instanceof On:" + (dayOfMonthField.getExpression() instanceof On)); + logger.info("dayOfMonthField instanceof And:" + (dayOfMonthField.getExpression() instanceof And)); + logger.info("dayOfMonthField instanceof QuestionMark:" + (dayOfMonthField.getExpression() instanceof QuestionMark)); CronField monthField = cron.retrieve(CronFieldName.MONTH); - logger.info("monthField instanceof Between:"+(monthField.getExpression() instanceof Between)); - logger.info("monthField instanceof Always:"+(monthField.getExpression() instanceof Always)); - logger.info("monthField instanceof Every:"+(monthField.getExpression() instanceof Every)); - logger.info("monthField instanceof On:"+(monthField.getExpression() instanceof On)); - logger.info("monthField instanceof And:"+(monthField.getExpression() instanceof And)); - logger.info("monthField instanceof QuestionMark:"+(monthField.getExpression() instanceof QuestionMark)); + logger.info("monthField instanceof Between:" + (monthField.getExpression() instanceof Between)); + logger.info("monthField instanceof Always:" + (monthField.getExpression() instanceof Always)); + logger.info("monthField instanceof Every:" + (monthField.getExpression() instanceof Every)); + logger.info("monthField instanceof On:" + (monthField.getExpression() instanceof On)); + logger.info("monthField instanceof And:" + (monthField.getExpression() instanceof And)); + logger.info("monthField instanceof QuestionMark:" + (monthField.getExpression() instanceof QuestionMark)); CronField dayOfWeekField = cron.retrieve(CronFieldName.DAY_OF_WEEK); - logger.info("dayOfWeekField instanceof Between:"+(dayOfWeekField.getExpression() instanceof Between)); - logger.info("dayOfWeekField instanceof Always:"+(dayOfWeekField.getExpression() instanceof Always)); - logger.info("dayOfWeekField instanceof Every:"+(dayOfWeekField.getExpression() instanceof Every)); - logger.info("dayOfWeekField instanceof On:"+(dayOfWeekField.getExpression() instanceof On)); - logger.info("dayOfWeekField instanceof And:"+(dayOfWeekField.getExpression() instanceof And)); - logger.info("dayOfWeekField instanceof QuestionMark:"+(dayOfWeekField.getExpression() instanceof QuestionMark)); - + logger.info("dayOfWeekField instanceof Between:" + (dayOfWeekField.getExpression() instanceof Between)); + logger.info("dayOfWeekField instanceof Always:" + (dayOfWeekField.getExpression() instanceof Always)); + logger.info("dayOfWeekField instanceof Every:" + (dayOfWeekField.getExpression() instanceof Every)); + logger.info("dayOfWeekField instanceof On:" + (dayOfWeekField.getExpression() instanceof On)); + logger.info("dayOfWeekField instanceof And:" + (dayOfWeekField.getExpression() instanceof And)); + logger.info("dayOfWeekField instanceof QuestionMark:" + (dayOfWeekField.getExpression() instanceof QuestionMark)); + CronField yearField = cron.retrieve(CronFieldName.YEAR); - logger.info("yearField instanceof Between:"+(yearField.getExpression() instanceof Between)); - logger.info("yearField instanceof Always:"+(yearField.getExpression() instanceof Always)); - logger.info("yearField instanceof Every:"+(yearField.getExpression() instanceof Every)); - logger.info("yearField instanceof On:"+(yearField.getExpression() instanceof On)); - logger.info("yearField instanceof And:"+(yearField.getExpression() instanceof And)); - logger.info("yearField instanceof QuestionMark:"+(yearField.getExpression() instanceof QuestionMark)); + logger.info("yearField instanceof Between:" + (yearField.getExpression() instanceof Between)); + logger.info("yearField instanceof Always:" + (yearField.getExpression() instanceof Always)); + logger.info("yearField instanceof Every:" + (yearField.getExpression() instanceof Every)); + logger.info("yearField instanceof On:" + (yearField.getExpression() instanceof On)); + logger.info("yearField instanceof And:" + (yearField.getExpression() instanceof And)); + logger.info("yearField instanceof QuestionMark:" + (yearField.getExpression() instanceof QuestionMark)); CycleEnum cycleEnum = CronUtils.getMaxCycle(minCrontab); - if(cycleEnum !=null){ + if (cycleEnum != null) { logger.info(cycleEnum.name()); - }else{ + } else { logger.info("can't get scheduleType"); } } @@ -213,7 +227,7 @@ public class CronUtilsTest { } @Test - public void getExpirationTime(){ + public void getExpirationTime() { Date startTime = DateUtils.stringToDate("2020-02-07 18:30:00"); Date expirationTime = CronUtils.getExpirationTime(startTime, CycleEnum.HOUR); Assert.assertEquals("2020-02-07 19:30:00", DateUtils.dateToString(expirationTime)); diff --git a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java index 0c9ddfbc7f..09ba19e911 100644 --- a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java +++ b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java @@ -72,7 +72,7 @@ import org.apache.dolphinscheduler.plugin.task.api.enums.dp.OptionSourceType; import org.apache.dolphinscheduler.plugin.task.api.enums.dp.ValueType; import org.apache.dolphinscheduler.plugin.task.api.model.ResourceInfo; import org.apache.dolphinscheduler.service.exceptions.ServiceException; -import org.apache.dolphinscheduler.service.quartz.cron.CronUtilsTest; +import org.apache.dolphinscheduler.service.cron.CronUtilsTest; import org.apache.dolphinscheduler.spi.params.base.FormType; import org.junit.Assert; import org.junit.Rule; diff --git a/dolphinscheduler-worker/src/main/java/org/apache/dolphinscheduler/server/worker/WorkerServer.java b/dolphinscheduler-worker/src/main/java/org/apache/dolphinscheduler/server/worker/WorkerServer.java index 05d7718bd2..12f80573d9 100644 --- a/dolphinscheduler-worker/src/main/java/org/apache/dolphinscheduler/server/worker/WorkerServer.java +++ b/dolphinscheduler-worker/src/main/java/org/apache/dolphinscheduler/server/worker/WorkerServer.java @@ -63,8 +63,6 @@ import org.springframework.transaction.annotation.EnableTransactionManagement; excludeFilters = { @ComponentScan.Filter(type = FilterType.REGEX, pattern = { "org.apache.dolphinscheduler.service.process.*", - // todo: split the quartz into a single module - "org.apache.dolphinscheduler.service.quartz.*", "org.apache.dolphinscheduler.service.queue.*", }) } diff --git a/pom.xml b/pom.xml index 3628920b27..f8a325ee77 100644 --- a/pom.xml +++ b/pom.xml @@ -386,6 +386,17 @@ ${project.version} + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-api + ${project.version} + + + org.apache.dolphinscheduler + dolphinscheduler-scheduler-quartz + ${project.version} + + org.apache.dolphinscheduler dolphinscheduler-datasource-all @@ -1254,5 +1265,6 @@ dolphinscheduler-log-server dolphinscheduler-tools dolphinscheduler-ui + dolphinscheduler-scheduler-plugin