You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 00:28:11