惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

The GitHub Blog
The GitHub Blog
S
SegmentFault 最新的问题
MyScale Blog
MyScale Blog
有赞技术团队
有赞技术团队
V
Visual Studio Blog
T
The Blog of Author Tim Ferriss
爱范儿
爱范儿
Vercel News
Vercel News
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
Y
Y Combinator Blog
Blog — PlanetScale
Blog — PlanetScale
D
DataBreaches.Net
美团技术团队
Microsoft Security Blog
Microsoft Security Blog
大猫的无限游戏
大猫的无限游戏
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
酷 壳 – CoolShell
酷 壳 – CoolShell
GbyAI
GbyAI
A
About on SuperTechFans
云风的 BLOG
云风的 BLOG
The Cloudflare Blog
宝玉的分享
宝玉的分享
V
V2EX
Microsoft Azure Blog
Microsoft Azure Blog

博客园 - 刚泡

CompletableFuture用法详解 使用责任链模式简化if-else代码示例 使用Function Interface简化if-else代码示例 自定义AngularJS Modal弹框大小的方法 SpringCloud Alibaba 改造Sentinel Dashboard将熔断规则持久化到Nacos SpringCloud Alibaba 改造Sentinel Dashboard将流控规则持久化到Nacos [Access][Microsoft][ODBC 驱动程序管理器] 无效的字符串或缓冲区长度 Invalid string or buffer length Java操作Kafka执行不成功的解决方法,Kafka Broker Advertised.Listeners属性的设置 pom.xml错误:org.codehaus.plexus.archiver.jar.Manifest.write(java.io.PrintWriter)的解决方法 Spring Cloud @HystrixCommand和@CacheResult注解使用,参数配置 Spring Boot的属性加载顺序 Spring Boot 使用properties如何多环境配置 Spring Cloud微服务实战阅读笔记(一) 基础知识 Java使用POI为Excel打水印,调整列宽并设置Excel只读(用户不可编辑) Sitemesh 3使用及配置 Java上传文件夹(Jersey) 路径名导致的异常:javax.imageio.IIOException: Can't read input file! 解决安卓微信浏览器中上传不能调用相机的问题 Mac下安装SVN插件javaHL not available的解决方法
数据库短暂波动导致xxl-job-admin不调度任务的问题
刚泡 · 2026-03-08 · via 博客园 - 刚泡

1. 问题描述

xxl-job-admin作为定时任务调度中心,当它与数据库之间存在短暂的网络波动,导致短时间内无法连接数据库。那么当数据库连接恢复后,任务调度中心就不再调度任务了,也就是说自某一个日期之后调度就停止了。

2. 问题成因分析

根因:调度线程因未捕获的异常退出,且不会自动重启。

  1. 调度线程在 JobScheduleHelper 里是单线程循环(scheduleThread),负责:
  • 开启事务 → scheduleLock()(FOR UPDATE)→ scheduleJobQuery() → 处理任务并 scheduleUpdate() → 在 finally 里 commit(transactionStatus)。
  1. 问题出在 finally 块(约 146–147 行):
} finally {
	// transaction commit
	XxlJobAdminBootstrap.getInstance().getTransactionManager().commit(transactionStatus);   // avlid schedule repeat
}
  • 未判空:若 getTransaction() 在赋值前抛异常,transactionStatus 为 null,commit(null) 可能触发 NPE。
  • 未捕获:数据库短暂失联或网络异常时,commit() 会抛异常(如连接关闭、通信失败),异常会从 finally 向外抛出。
  • 线程退出:异常导致 Runnable.run() 结束,调度线程终止;该线程只 start 一次,没有重启逻辑,因此调度会永久停止。
  1. 与执行器网络失败的关系:执行器调用在 JobTriggerPoolHelper 的线程池中异步执行,不会直接杀死 schedule 线程;但若同一时段发生 DB 异常(或连接池/事务受影响),就会在 commit 时抛错并导致上述退出。
    修复计划:在 finally 中仅在 transactionStatus != null 时执行提交,并用 try-catch 包裹 commit,失败时尝试 rollback 并打日志,避免异常向外传播导致调度线程退出。
} finally {
                        // 仅在有事务时提交,避免 commit(null) 导致 NPE;提交失败时回滚并记录,避免异常抛出导致调度线程退出(如 DB 短暂失联、网络问题后调度永久停止)
                        if (transactionStatus != null) {
                            try {
                                XxlJobAdminBootstrap.getInstance().getTransactionManager().commit(transactionStatus);
                            } catch (Throwable e) {
                                if (!scheduleThreadToStop) {
                                    logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread commit error, rollback and continue. error:{}", e.getMessage(), e);
                                }
                                try {
                                    XxlJobAdminBootstrap.getInstance().getTransactionManager().rollback(transactionStatus);
                                } catch (Throwable rollbackEx) {
                                    if (!scheduleThreadToStop) {
                                        logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread rollback error:{}", rollbackEx.getMessage(), rollbackEx);
                                    }
                                }
                            }
                        }
                    }

3. 总结

问题成因

  1. 现象:网络异常或数据库短暂失联后,到某个时间点后定时任务不再被调度。
  2. 根因:调度由 单一线程 JobScheduleHelper#scheduleThread 负责;在 finally 里无条件执行 commit(transactionStatus),且未捕获异常:
    • transactionStatus == null 时(例如 getTransaction() 先抛异常),commit(null) 可能触发 NPE。
    • 数据库/网络异常时,commit() 会抛出异常并直接冒出 run()。
  3. 结果:该线程一旦因异常退出就不会再被拉起,调度循环永久停止,表现为“到某日期后就不调度了”。

执行器与 admin 之间的网络失败是在触发线程池里处理的,不会直接杀死调度线程;但若同一时期发生 DB 异常或连接问题,就会在本次事务的 commit 里触发上述退出。

4. 完整代码

package com.xxl.job.admin.scheduler.thread;

import com.xxl.job.admin.constant.TriggerStatus;
import com.xxl.job.admin.model.XxlJobInfo;
import com.xxl.job.admin.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.admin.scheduler.misfire.MisfireStrategyEnum;
import com.xxl.job.admin.scheduler.trigger.TriggerTypeEnum;
import com.xxl.job.admin.scheduler.type.ScheduleTypeEnum;
import com.xxl.tool.core.CollectionTool;
import com.xxl.tool.core.MapTool;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition;

import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;

/**
 * @author xuxueli 2019-05-21
 */
public class JobScheduleHelper {
    private static final Logger logger = LoggerFactory.getLogger(JobScheduleHelper.class);


    public static final long PRE_READ_MS = 5000;    // pre read

    private Thread scheduleThread;
    private Thread ringThread;
    private volatile boolean scheduleThreadToStop = false;
    private volatile boolean ringThreadToStop = false;
    private final Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>();

    /**
     * start
     */
    public void start(){

        // schedule thread
        scheduleThread = new Thread(new Runnable() {
            @Override
            public void run() {

                // align time
                try {
                    TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 );
                } catch (Throwable e) {
                    if (!scheduleThreadToStop) {
                        logger.error(e.getMessage(), e);
                    }
                }
                logger.info(">>>>>>>>> init xxl-job admin scheduler success.");

                // pre-read count: treadpool-size * trigger-qps (each trigger cost 100ms, qps = 1000/100 = 100)
                int preReadCount = (XxlJobAdminBootstrap.getInstance().getTriggerPoolFastMax() + XxlJobAdminBootstrap.getInstance().getTriggerPoolSlowMax()) * 10;

                // do schedule
                while (!scheduleThreadToStop) {

                    // param
                    long start = System.currentTimeMillis();
                    boolean preReadSuc = true;

                    // transaction start
                    TransactionStatus transactionStatus = null;
                    try {
                        transactionStatus = XxlJobAdminBootstrap.getInstance().getTransactionManager().getTransaction(new DefaultTransactionDefinition());
                        // 1、job lock
                        String lockedRecord = XxlJobAdminBootstrap.getInstance().getXxlJobLockMapper().scheduleLock();
                        long nowTime = System.currentTimeMillis();

                        // scan and process job
                        List<XxlJobInfo> scheduleList = XxlJobAdminBootstrap.getInstance().getXxlJobInfoMapper().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount);
                        if (CollectionTool.isNotEmpty(scheduleList)) {

                            // 2、push time-ring
                            for (XxlJobInfo jobInfo: scheduleList) {

                                // time-ring jump
                                if (nowTime > jobInfo.getTriggerNextTime() + PRE_READ_MS) {
                                    // 2.1、trigger-expire > 5s:pass && make next-trigger-time

                                    // 1、misfire handle
                                    MisfireStrategyEnum misfireStrategyEnum = MisfireStrategyEnum.match(jobInfo.getMisfireStrategy(), MisfireStrategyEnum.DO_NOTHING);
                                    misfireStrategyEnum.getMisfireHandler().handle(jobInfo.getId());

                                    // 2、fresh next
                                    refreshNextTriggerTime(jobInfo, new Date());

                                } else if (nowTime > jobInfo.getTriggerNextTime()) {
                                    // 2.2、trigger-expire < 5s:direct-trigger && make next-trigger-time

                                    // 1、trigger direct
                                    XxlJobAdminBootstrap.getInstance().getJobTriggerPoolHelper().trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null);
                                    logger.debug(">>>>>>>>>>> xxl-job, schedule expire, direct trigger : jobId = " + jobInfo.getId() );

                                    // 2、fresh next
                                    refreshNextTriggerTime(jobInfo, new Date());

                                    // next-trigger-time in 5s, pre-read again
                                    if (jobInfo.getTriggerStatus()== TriggerStatus.RUNNING.getValue() && nowTime + PRE_READ_MS > jobInfo.getTriggerNextTime()) {

                                        // 1、make ring second
                                        int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);

                                        // 2、push time ring (pre read)
                                        pushTimeRing(ringSecond, jobInfo.getId());
                                        logger.debug(">>>>>>>>>>> xxl-job, schedule pre-read, push trigger : jobId = " + jobInfo.getId() );

                                        // 3、fresh next
                                        refreshNextTriggerTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));

                                    }

                                } else {
                                    // 2.3、trigger-pre-read:time-ring trigger && make next-trigger-time

                                    // 1、make ring second
                                    int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);

                                    // 2、push time ring
                                    pushTimeRing(ringSecond, jobInfo.getId());
                                    logger.debug(">>>>>>>>>>> xxl-job, schedule normal, push trigger : jobId = " + jobInfo.getId() );

                                    // 3、fresh next
                                    refreshNextTriggerTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));

                                }

                            }

                            // 3、update trigger info
                            for (XxlJobInfo jobInfo: scheduleList) {
                                XxlJobAdminBootstrap.getInstance().getXxlJobInfoMapper().scheduleUpdate(jobInfo);
                            }

                        } else {
                            preReadSuc = false;
                        }

                    } catch (Throwable e) {
                        if (!scheduleThreadToStop) {
                            logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread error:{}", e.getMessage(), e);
                        }
                    } finally {
                        // 仅在有事务时提交,避免 commit(null) 导致 NPE;提交失败时回滚并记录,避免异常抛出导致调度线程退出(如 DB 短暂失联、网络问题后调度永久停止)
                        if (transactionStatus != null) {
                            try {
                                XxlJobAdminBootstrap.getInstance().getTransactionManager().commit(transactionStatus);
                            } catch (Throwable e) {
                                if (!scheduleThreadToStop) {
                                    logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread commit error, rollback and continue. error:{}", e.getMessage(), e);
                                }
                                try {
                                    XxlJobAdminBootstrap.getInstance().getTransactionManager().rollback(transactionStatus);
                                } catch (Throwable rollbackEx) {
                                    if (!scheduleThreadToStop) {
                                        logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread rollback error:{}", rollbackEx.getMessage(), rollbackEx);
                                    }
                                }
                            }
                        } else {
                            logger.warn(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread no transaction, continue.");
                        }
                    }
                    // transaction end
                    long cost = System.currentTimeMillis()-start;


                    // Wait seconds, align second
                    if (cost < 1000) {  // scan-overtime, not wait
                        try {
                            // pre-read period: success > scan each second; fail > skip this period;
                            TimeUnit.MILLISECONDS.sleep((preReadSuc?1000:PRE_READ_MS) - System.currentTimeMillis()%1000);
                        } catch (Throwable e) {
                            if (!scheduleThreadToStop) {
                                logger.error(e.getMessage(), e);
                            }
                        }
                    }

                }

                logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop");
            }
        });
        scheduleThread.setDaemon(true);
        scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
        scheduleThread.start();


        // ring thread
        ringThread = new Thread(new Runnable() {
            @Override
            public void run() {

                while (!ringThreadToStop) {

                    // align second
                    try {
                        TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis() % 1000);
                    } catch (Throwable e) {
                        if (!ringThreadToStop) {
                            logger.error(e.getMessage(), e);
                        }
                    }

                    try {
                        // second data
                        List<Integer> ringItemData = new ArrayList<>();

                        // collect rind data, by second
                        int nowSecond = Calendar.getInstance().get(Calendar.SECOND);
                        for (int i = 0; i <= 2; i++) {                                                              // 避免调度遗漏:处理耗时太长、跨过刻度,除当前刻度外 + 向前校验2个刻度;
                            List<Integer> ringItemList = ringData.remove( (nowSecond+60-i)%60 );
                            if (CollectionTool.isNotEmpty(ringItemList)) {
                                // distinct for each second
                                List<Integer> ringItemListDistinct = ringItemList.stream().distinct().toList();     // 避免调度重复:重复推送时间轮刻度,去重只保留一个;;
                                if (ringItemListDistinct.size() < ringItemList.size()) {
                                    logger.warn(">>>>>>>>>>> xxl-job, time-ring found job repeat beat : " + nowSecond + " = " + ringItemData);
                                }

                                // collect ring item
                                ringItemData.addAll(ringItemListDistinct);
                            }
                        }

                        // ring trigger
                        logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + ringItemData);
                        if (CollectionTool.isNotEmpty(ringItemData)) {
                            // do trigger
                            for (int jobId: ringItemData) {
                                // do trigger
                                XxlJobAdminBootstrap.getInstance().getJobTriggerPoolHelper().trigger(jobId, TriggerTypeEnum.CRON, -1, null, null, null);
                            }
                            // clear
                            ringItemData.clear();
                        }
                    } catch (Throwable e) {
                        if (!ringThreadToStop) {
                            logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e.getMessage(), e);
                        }
                    }
                }
                logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
            }
        });
        ringThread.setDaemon(true);
        ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
        ringThread.start();
    }

    /**
     * refresh next trigger time of job
     *
     * @param jobInfo   job info
     * @param fromTime  from time
     */
    private void refreshNextTriggerTime(XxlJobInfo jobInfo, Date fromTime) {
        try {
            // generate next trigger time
            ScheduleTypeEnum scheduleTypeEnum = ScheduleTypeEnum.match(jobInfo.getScheduleType(), ScheduleTypeEnum.NONE);
            Date nextTriggerTime = scheduleTypeEnum.getScheduleType().generateNextTriggerTime(jobInfo, fromTime);

            // refresh next trigger-time + status
            if (nextTriggerTime != null) {
                // generate success
                jobInfo.setTriggerStatus(-1);                               // pass, may be Inaccurate
                jobInfo.setTriggerLastTime(jobInfo.getTriggerNextTime());
                jobInfo.setTriggerNextTime(nextTriggerTime.getTime());
            } else {
                // generate fail, stop job
                jobInfo.setTriggerStatus(TriggerStatus.STOPPED.getValue());
                jobInfo.setTriggerLastTime(0);
                jobInfo.setTriggerNextTime(0);
                logger.error(">>>>>>>>>>> xxl-job, refreshNextValidTime fail for job: jobId={}, scheduleType={}, scheduleConf={}",
                        jobInfo.getId(), jobInfo.getScheduleType(), jobInfo.getScheduleConf());
            }
        } catch (Throwable e) {
            // generate error, stop job
            jobInfo.setTriggerStatus(TriggerStatus.STOPPED.getValue());
            jobInfo.setTriggerLastTime(0);
            jobInfo.setTriggerNextTime(0);

            logger.error(">>>>>>>>>>> xxl-job, refreshNextValidTime error for job: jobId={}, scheduleType={}, scheduleConf={}",
                    jobInfo.getId(), jobInfo.getScheduleType(), jobInfo.getScheduleConf(), e);
        }
    }

    /**
     * push time ring
     *
     * @param ringSecond    ring second
     * @param jobId         job id
     */
    private void pushTimeRing(int ringSecond, int jobId){
        // get ringItemData, init when not exists
        List<Integer> ringItemList = ringData.computeIfAbsent(
                ringSecond,
                k -> new ArrayList<>());

        // push async rind
        ringItemList.add(jobId);
        logger.debug(">>>>>>>>>>> xxl-job, schedule push time-ring : " + ringSecond + " = " + List.of(ringItemList));
    }

    /**
     * stop
     */
    public void stop(){

        // 1、stop schedule
        scheduleThreadToStop = true;
        try {
            TimeUnit.SECONDS.sleep(1);  // wait
        } catch (Throwable e) {
            logger.error(e.getMessage(), e);
        }
        if (scheduleThread.getState() != Thread.State.TERMINATED){
            // interrupt and wait
            scheduleThread.interrupt();
            try {
                scheduleThread.join();
            } catch (Throwable e) {
                logger.error(e.getMessage(), e);
            }
        }

        // if has ring data
        boolean hasRingData = false;
        if (MapTool.isNotEmpty(ringData)) {
            for (int second : ringData.keySet()) {
                List<Integer> ringItemList = ringData.get(second);
                if (CollectionTool.isNotEmpty(ringItemList)) {
                    hasRingData = true;
                    break;
                }
            }
        }
        if (hasRingData) {
            try {
                TimeUnit.SECONDS.sleep(8);
            } catch (Throwable e) {
                logger.error(e.getMessage(), e);
            }
        }

        // stop ring (wait job-in-memory stop)
        ringThreadToStop = true;
        try {
            TimeUnit.SECONDS.sleep(1);
        } catch (Throwable e) {
            logger.error(e.getMessage(), e);
        }
        if (ringThread.getState() != Thread.State.TERMINATED){
            // interrupt and wait
            ringThread.interrupt();
            try {
                ringThread.join();
            } catch (Throwable e) {
                logger.error(e.getMessage(), e);
            }
        }

        logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper stop");
    }

}