如何从ThreadPoolTaskScheduler获取任务执行结果并监控任务完成?
在应用中,我从数据库获取元数据并构建实现Runnable接口的MyJob实例,随后提交给ThreadPoolTaskScheduler执行。相关代码如下:
@Scheduled(fixedRateString = "${check.rate.minutes}", timeUnit = TimeUnit.MINUTES) public void checkForJobs() { LOG.info("Executing check at {}", Instant.now().toString()); int avail = threadPool.getPoolSize() - threadPool.getActiveCount(); // will retrieve jobs metadata to fill vacancies List<MetaJob> availableJobs = this.repository.getAvailableJobs(avail); if(!(availableJobs == null || availableJobs.isEmpty())) { LOG.info("Retrieved {} jobs for scheduling", availableJobs.size()); for (MetaJob job : availableJobs) { scheduleJob(job); } } else { LOG.info("No available jobs at this moment"); } } private void scheduleJob(MetaJob job) { JobSchedule schedule = job.getSchedule(); JobRunner jobRunner = getExecutor(job); // Creating Runnable out of metadata if(schedule == null) { threadPool.schedule(jobRunner, Instant.now()); } else { Trigger trigger; if(schedule.getType() == SchedType.CRON) { trigger = new CronTrigger(schedule.getExpression()); } else { trigger = new PeriodicTrigger(Duration.parse(schedule.getExpression())); } threadPool.schedule(jobRunner, trigger); } }
任务可正常执行,但需要处理完成或失败的结果:
- 若执行次数超出元数据指定的限制,需取消后续执行;
- 需更新数据库中
MyJob实例存储的结果; - 若任务失败,需分析失败原因,若是临时故障则允许继续执行,否则取消任务。
我尝试继承ThreadPoolTaskScheduler并重写其protected方法afterExecute(Runnable task, @Nullable Throwable ex),但遇到问题:
- 该方法会在任务执行前后被调用两次;
- 方法接收的
Runnable参数并非提交的MyJobRunner,而是java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask实例,导致类型转换异常,同时无法关联提交任务时获取的ScheduledFuture实例。
应用会同时处理数十个任务,需要明确识别哪个任务已完成及执行状态。现咨询:能否通过afterExecute方法识别完成的任务?若可以,该如何实现?若不行,还有哪些监控定时任务完成的可行方式?
一、通过afterExecute识别任务的可行方案
可以通过afterExecute识别任务,但需要解决ScheduledFutureTask的包装问题,具体实现如下:
1. 自定义ThreadPoolTaskScheduler与ScheduledThreadPoolExecutor
继承ThreadPoolTaskScheduler,重写createExecutor方法返回自定义的ScheduledThreadPoolExecutor,在其中解析ScheduledFutureTask获取实际提交的JobRunner:
public class CustomThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { private static final Logger LOG = LoggerFactory.getLogger(CustomThreadPoolTaskScheduler.class); @Override protected ScheduledExecutorService createExecutor(int poolSize, ThreadFactory threadFactory, RejectedExecutionHandler rejectedExecutionHandler) { return new CustomScheduledThreadPoolExecutor(poolSize, threadFactory, rejectedExecutionHandler); } private static class CustomScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { public CustomScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory, RejectedExecutionHandler handler) { super(corePoolSize, threadFactory, handler); } @Override protected void afterExecute(Runnable r, Throwable t) { super.afterExecute(r, t); // 仅处理任务执行完成的情况,避免执行前的无效调用 if (!(r instanceof Future<?> && ((Future<?>) r).isDone())) { return; } // 解析ScheduledFutureTask获取实际任务 if (r instanceof ScheduledFutureTask<?>) { ScheduledFutureTask<?> futureTask = (ScheduledFutureTask<?>) r; try { // 通过反射获取FutureTask中的callable字段 Field callableField = FutureTask.class.getDeclaredField("callable"); callableField.setAccessible(true); Object callable = callableField.get(futureTask); // 解析RunnableAdapter获取实际的JobRunner if (callable instanceof Executors.RunnableAdapter) { Executors.RunnableAdapter<?> adapter = (Executors.RunnableAdapter<?>) callable; Runnable actualTask = adapter.getRunnable(); if (actualTask instanceof JobRunner) { JobRunner jobRunner = (JobRunner) actualTask; MetaJob metaJob = jobRunner.getMetaJob(); // 假设JobRunner持有MetaJob引用 // 处理任务结果逻辑 handleTaskResult(metaJob, t, futureTask); } } } catch (NoSuchFieldException | IllegalAccessException e) { LOG.error("Failed to extract actual task from ScheduledFutureTask", e); } } } private void handleTaskResult(MetaJob metaJob, Throwable failureCause, ScheduledFutureTask<?> futureTask) { // 1. 更新数据库任务状态与结果 updateJobInDatabase(metaJob, failureCause == null, failureCause); // 2. 检查执行次数是否超限 int currentExecutions = metaJob.getExecutions() + 1; metaJob.setExecutions(currentExecutions); if (currentExecutions >= metaJob.getMaxExecutions()) { futureTask.cancel(false); return; } // 3. 判断失败是否为永久故障,是则取消任务 if (failureCause != null && isPermanentFailure(failureCause)) { futureTask.cancel(false); } } // 业务方法:更新数据库 private void updateJobInDatabase(MetaJob metaJob, boolean success, Throwable failureCause) { /* 实现逻辑 */ } // 业务方法:判断是否为永久故障 private boolean isPermanentFailure(Throwable cause) { /* 实现逻辑 */ } } }
注意事项
- 反射依赖
FutureTask的内部字段,可能因JDK版本变更失效; - 需要确保
JobRunner持有对应的MetaJob元数据,以便识别任务。
二、更可靠的方案:在JobRunner内部处理执行结果
这种方式无需依赖线程池的内部实现,直接在任务自身中处理结果,兼容性与可控性更强:
1. 修改JobRunner类
让JobRunner持有MetaJob和提交任务后得到的ScheduledFuture,在run方法中捕获执行结果:
public class JobRunner implements Runnable { private static final Logger LOG = LoggerFactory.getLogger(JobRunner.class); private final MetaJob metaJob; private ScheduledFuture<?> scheduledFuture; public JobRunner(MetaJob metaJob) { this.metaJob = metaJob; } public void setScheduledFuture(ScheduledFuture<?> scheduledFuture) { this.scheduledFuture = scheduledFuture; } @Override public void run() { Throwable failureCause = null; boolean success = false; try { // 执行实际任务逻辑 executeBusinessLogic(); success = true; } catch (Exception e) { failureCause = e; LOG.error("Job execution failed for job: {}", metaJob.getId(), e); } finally { // 统一处理执行结果 handleExecutionResult(success, failureCause); } } private void handleExecutionResult(boolean success, Throwable failureCause) { // 更新数据库任务状态 updateJobStatusInDb(success, failureCause); // 检查执行次数是否超限 int currentExecutions = metaJob.getExecutions() + 1; metaJob.setExecutions(currentExecutions); if (currentExecutions >= metaJob.getMaxExecutions()) { cancelFutureIfNecessary(); return; } // 处理失败场景:永久故障则取消任务 if (!success && isPermanentFailure(failureCause)) { cancelFutureIfNecessary(); } } private void cancelFutureIfNecessary() { if (scheduledFuture != null && !scheduledFuture.isCancelled() && !scheduledFuture.isDone()) { scheduledFuture.cancel(false); LOG.info("Cancelled future executions for job: {}", metaJob.getId()); } } // 业务方法:执行实际任务逻辑 private void executeBusinessLogic() { /* 实现逻辑 */ } // 业务方法:更新数据库 private void updateJobStatusInDb(boolean success, Throwable failureCause) { /* 实现逻辑 */ } // 业务方法:判断是否为永久故障 private boolean isPermanentFailure(Throwable cause) { /* 实现逻辑 */ } }
2. 更新scheduleJob方法
提交任务后,将ScheduledFuture设置给JobRunner:
private void scheduleJob(MetaJob job) { JobSchedule schedule = job.getSchedule(); JobRunner jobRunner = getExecutor(job); ScheduledFuture<?> future; if(schedule == null) { future = threadPool.schedule(jobRunner, Instant.now()); } else { Trigger trigger = schedule.getType() == SchedType.CRON ? new CronTrigger(schedule.getExpression()) : new PeriodicTrigger(Duration.parse(schedule.getExpression())); future = threadPool.schedule(jobRunner, trigger); } // 将Future关联到JobRunner,用于后续取消任务 jobRunner.setScheduledFuture(future); }
三、其他可选方案
1. 使用AOP切面解耦结果处理
通过Spring AOP环绕JobRunner的run方法,将结果处理逻辑与任务逻辑分离:
@Aspect @Component public class JobExecutionAspect { private static final Logger LOG = LoggerFactory.getLogger(JobExecutionAspect.class); @Around("execution(* com.yourpackage.JobRunner.run(..))") public Object handleJobExecution(ProceedingJoinPoint joinPoint) throws Throwable { JobRunner jobRunner = (JobRunner) joinPoint.getTarget(); MetaJob metaJob = jobRunner.getMetaJob(); ScheduledFuture<?> future = jobRunner.getScheduledFuture(); Throwable failureCause = null; boolean success = false; try { Object result = joinPoint.proceed(); success = true; return result; } catch (Throwable e) { failureCause = e; throw e; } finally { // 处理结果逻辑,与之前的handleExecutionResult一致 processJobResult(metaJob, future, success, failureCause); } } private void processJobResult(MetaJob metaJob, ScheduledFuture<?> future, boolean success, Throwable failureCause) { // 更新数据库、检查执行次数、判断故障类型等逻辑 } }
2. 自定义任务装饰器
通过ThreadPoolTaskScheduler的setTaskDecorator方法,包装任务添加结果处理逻辑:
threadPool.setTaskDecorator(runnable -> { if (runnable instanceof JobRunner) { JobRunner jobRunner = (JobRunner) runnable; return () -> { Throwable failureCause = null; boolean success = false; try { jobRunner.run(); success = true; } catch (Exception e) { failureCause = e; throw e; } finally { // 处理结果逻辑 handleJobResult(jobRunner.getMetaJob(), jobRunner.getScheduledFuture(), success, failureCause); } }; } return runnable; });
内容的提问来源于stack exchange,提问作者Gary Greenberg

