最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

PowerJob分布式任務(wù)調(diào)度源碼流程解讀

 更新時(shí)間:2024年02月16日 09:46:25   作者:codecraft  
這篇文章主要為大家介紹了PowerJob分布式任務(wù)調(diào)度源碼流程解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

架構(gòu)圖

官方提供:

本文主要研究一下PowerJob的任務(wù)調(diào)度

CoreScheduleTaskManager

tech/powerjob/server/core/scheduler/CoreScheduleTaskManager.java

@Service
@Slf4j
@RequiredArgsConstructor
public class CoreScheduleTaskManager implements InitializingBean, DisposableBean {
    private final PowerScheduleService powerScheduleService;
    private final InstanceStatusCheckService instanceStatusCheckService;
    private final List<Thread> coreThreadContainer = new ArrayList<>();
    @SuppressWarnings("AlibabaAvoidManuallyCreateThread")
    @Override
    public void afterPropertiesSet() {
        // 定時(shí)調(diào)度
        coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleCronJob", PowerScheduleService.SCHEDULE_RATE, () -> powerScheduleService.scheduleNormalJob(TimeExpressionType.CRON)), "Thread-ScheduleCronJob"));
        coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleDailyTimeIntervalJob", PowerScheduleService.SCHEDULE_RATE, () -> powerScheduleService.scheduleNormalJob(TimeExpressionType.DAILY_TIME_INTERVAL)), "Thread-ScheduleDailyTimeIntervalJob"));
        coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleCronWorkflow", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::scheduleCronWorkflow), "Thread-ScheduleCronWorkflow"));
        coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleFrequentJob", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::scheduleFrequentJob), "Thread-ScheduleFrequentJob"));
        // 數(shù)據(jù)清理
        coreThreadContainer.add(new Thread(new LoopRunnable("CleanWorkerData", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::cleanData), "Thread-CleanWorkerData"));
        // 狀態(tài)檢查
        coreThreadContainer.add(new Thread(new LoopRunnable("CheckRunningInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkRunningInstance), "Thread-CheckRunningInstance"));
        coreThreadContainer.add(new Thread(new LoopRunnable("CheckWaitingDispatchInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWaitingDispatchInstance), "Thread-CheckWaitingDispatchInstance"));
        coreThreadContainer.add(new Thread(new LoopRunnable("CheckWaitingWorkerReceiveInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWaitingWorkerReceiveInstance), "Thread-CheckWaitingWorkerReceiveInstance"));
        coreThreadContainer.add(new Thread(new LoopRunnable("CheckWorkflowInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWorkflowInstance), "Thread-CheckWorkflowInstance"));
        coreThreadContainer.forEach(Thread::start);
    }
    //......
}
CoreScheduleTaskManager在afterPropertiesSet的時(shí)候會(huì)啟動(dòng)一系列的線(xiàn)程,它們都是LoopRunnable類(lèi)型的,分別調(diào)度powerScheduleService.scheduleNormalJob(TimeExpressionType.CRON)、powerScheduleService.scheduleNormalJob(TimeExpressionType.DAILY_TIME_INTERVAL)、powerScheduleService::scheduleCronWorkflow、powerScheduleService::scheduleFrequentJob、powerScheduleService::cleanData、instanceStatusCheckService::checkRunningInstance、instanceStatusCheckService::checkWaitingDispatchInstance、instanceStatusCheckService::checkWaitingWorkerReceiveInstance、instanceStatusCheckService::checkWorkflowInstance

LoopRunnable

@RequiredArgsConstructor
    private static class LoopRunnable implements Runnable {

        private final String taskName;

        private final Long runningInterval;

        private final Runnable innerRunnable;

        @SuppressWarnings("BusyWait")
        @Override
        public void run() {
            log.info("start task : {}.", taskName);
            while (true) {
                try {
                    innerRunnable.run();
                    Thread.sleep(runningInterval);
                } catch (InterruptedException e) {
                    log.warn("[{}] task has been interrupted!", taskName, e);
                    break;
                } catch (Exception e) {
                    log.error("[{}] task failed!", taskName, e);
                }
            }
        }
    }
LoopRunnable的構(gòu)造器接收taskName、runningInterval、innerRunnable三個(gè)參數(shù),其run方法通過(guò)while true循環(huán)內(nèi)部執(zhí)行innerRunnable.run(),執(zhí)行完sleep指定的runningInterval,若捕獲到InterruptedException則break跳出循環(huán),若其他異常則打印error日志

PowerScheduleService

PowerScheduleService主要提供了scheduleNormalJob、scheduleCronWorkflow、scheduleFrequentJob、cleanData方法

scheduleNormalJob

tech/powerjob/server/core/scheduler/PowerScheduleService.java

public void scheduleNormalJob(TimeExpressionType timeExpressionType) {
        long start = System.currentTimeMillis();
        // 調(diào)度 CRON 表達(dá)式 JOB
        try {
            final List<Long> allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
            if (CollectionUtils.isEmpty(allAppIds)) {
                log.info("[NormalScheduler] current server has no app's job to schedule.");
                return;
            }
            scheduleNormalJob0(timeExpressionType, allAppIds);
        } catch (Exception e) {
            log.error("[NormalScheduler] schedule cron job failed.", e);
        }
        long cost = System.currentTimeMillis() - start;
        log.info("[NormalScheduler] {} job schedule use {} ms.", timeExpressionType, cost);
        if (cost > SCHEDULE_RATE) {
            log.warn("[NormalScheduler] The database query is using too much time({}ms), please check if the database load is too high!", cost);
        }
    }
scheduleNormalJob方法主要是查詢(xún)當(dāng)前server負(fù)責(zé)的appId列表,然后內(nèi)部委托改為scheduleNormalJob0

scheduleNormalJob0

private void scheduleNormalJob0(TimeExpressionType timeExpressionType, List<Long> appIds) {

        long nowTime = System.currentTimeMillis();
        long timeThreshold = nowTime + 2 * SCHEDULE_RATE;
        Lists.partition(appIds, MAX_APP_NUM).forEach(partAppIds -> {

            try {

                // 查詢(xún)條件:任務(wù)開(kāi)啟 + 使用CRON表達(dá)調(diào)度時(shí)間 + 指定appId + 即將需要調(diào)度執(zhí)行
                List<JobInfoDO> jobInfos = jobInfoRepository.findByAppIdInAndStatusAndTimeExpressionTypeAndNextTriggerTimeLessThanEqual(partAppIds, SwitchableStatus.ENABLE.getV(), timeExpressionType.getV(), timeThreshold);

                if (CollectionUtils.isEmpty(jobInfos)) {
                    return;
                }

                // 1. 批量寫(xiě)日志表
                Map<Long, Long> jobId2InstanceId = Maps.newHashMap();
                log.info("[NormalScheduler] These {} jobs will be scheduled: {}.", timeExpressionType.name(), jobInfos);

                jobInfos.forEach(jobInfo -> {
                    Long instanceId = instanceService.create(jobInfo.getId(), jobInfo.getAppId(), jobInfo.getJobParams(), null, null, jobInfo.getNextTriggerTime()).getInstanceId();
                    jobId2InstanceId.put(jobInfo.getId(), instanceId);
                });
                instanceInfoRepository.flush();

                // 2. 推入時(shí)間輪中等待調(diào)度執(zhí)行
                jobInfos.forEach(jobInfoDO -> {

                    Long instanceId = jobId2InstanceId.get(jobInfoDO.getId());

                    long targetTriggerTime = jobInfoDO.getNextTriggerTime();
                    long delay = 0;
                    if (targetTriggerTime < nowTime) {
                        log.warn("[Job-{}] schedule delay, expect: {}, current: {}", jobInfoDO.getId(), targetTriggerTime, System.currentTimeMillis());
                    } else {
                        delay = targetTriggerTime - nowTime;
                    }

                    InstanceTimeWheelService.schedule(instanceId, delay, () -> dispatchService.dispatch(jobInfoDO, instanceId, Optional.empty(), Optional.empty()));
                });

                // 3. 計(jì)算下一次調(diào)度時(shí)間(忽略5S內(nèi)的重復(fù)執(zhí)行,即CRON模式下最小的連續(xù)執(zhí)行間隔為 SCHEDULE_RATE ms)
                jobInfos.forEach(jobInfoDO -> {
                    try {
                        refreshJob(timeExpressionType, jobInfoDO);
                    } catch (Exception e) {
                        log.error("[Job-{}] refresh job failed.", jobInfoDO.getId(), e);
                    }
                });
                jobInfoRepository.flush();


            } catch (Exception e) {
                log.error("[NormalScheduler] schedule {} job failed.", timeExpressionType.name(), e);
            }
        });
    }
scheduleNormalJob0主要是調(diào)度CRON、DAILY_TIME_INTERVAL類(lèi)型的任務(wù),它通過(guò)jobInfoRepository查找指定appId、狀態(tài)啟用、指定TimeExpressionType,以及NextTriggerTime小于等于nowTime + 2 * SCHEDULE_RATE的任務(wù),然后挨個(gè)執(zhí)行instanceService.create創(chuàng)建任務(wù)實(shí)例,然后放入到InstanceTimeWheelService.schedule進(jìn)行調(diào)度,最后計(jì)算和更新一下每個(gè)job的nextTriggerTime

scheduleCronWorkflow

public void scheduleCronWorkflow() {
        long start = System.currentTimeMillis();
        // 調(diào)度 CRON 表達(dá)式 WORKFLOW
        try {
            final List<Long> allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
            if (CollectionUtils.isEmpty(allAppIds)) {
                log.info("[CronWorkflowSchedule] current server has no app's workflow to schedule.");
                return;
            }
            scheduleWorkflowCore(allAppIds);
        } catch (Exception e) {
            log.error("[CronWorkflowSchedule] schedule cron workflow failed.", e);
        }
        long cost = System.currentTimeMillis() - start;
        log.info("[CronWorkflowSchedule] cron workflow schedule use {} ms.", cost);
        if (cost > SCHEDULE_RATE) {
            log.warn("[CronWorkflowSchedule] The database query is using too much time({}ms), please check if the database load is too high!", cost);
        }
    }
scheduleCronWorkflow主要是調(diào)度CRON 表達(dá)式 WORKFLOW,內(nèi)部委托給scheduleWorkflowCore

scheduleFrequentJob

public void scheduleFrequentJob() {
        long start = System.currentTimeMillis();
        // 調(diào)度 FIX_RATE/FIX_DELAY 表達(dá)式 JOB
        try {
            final List&lt;Long&gt; allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
            if (CollectionUtils.isEmpty(allAppIds)) {
                log.info("[FrequentJobSchedule] current server has no app's job to schedule.");
                return;
            }
            scheduleFrequentJobCore(allAppIds);
        } catch (Exception e) {
            log.error("[FrequentJobSchedule] schedule frequent job failed.", e);
        }
        long cost = System.currentTimeMillis() - start;
        log.info("[FrequentJobSchedule] frequent job schedule use {} ms.", cost);
        if (cost &gt; SCHEDULE_RATE) {
            log.warn("[FrequentJobSchedule] The database query is using too much time({}ms), please check if the database load is too high!", cost);
        }
    }
scheduleFrequentJob主要是調(diào)度FIX_RATE/FIX_DELAY 表達(dá)式 JOB,內(nèi)部委托給了scheduleFrequentJobCore

scheduleFrequentJobCore

private void scheduleFrequentJobCore(List<Long> appIds) {

        Lists.partition(appIds, MAX_APP_NUM).forEach(partAppIds -> {
            try {
                // 查詢(xún)所有的秒級(jí)任務(wù)(只包含ID)
                List<Long> jobIds = jobInfoRepository.findByAppIdInAndStatusAndTimeExpressionTypeIn(partAppIds, SwitchableStatus.ENABLE.getV(), TimeExpressionType.FREQUENT_TYPES);
                if (CollectionUtils.isEmpty(jobIds)) {
                    return;
                }
                // 查詢(xún)?nèi)罩居涗洷碇惺欠翊嬖谙嚓P(guān)的任務(wù)
                List<Long> runningJobIdList = instanceInfoRepository.findByJobIdInAndStatusIn(jobIds, InstanceStatus.GENERALIZED_RUNNING_STATUS);
                Set<Long> runningJobIdSet = Sets.newHashSet(runningJobIdList);

                List<Long> notRunningJobIds = Lists.newLinkedList();
                jobIds.forEach(jobId -> {
                    if (!runningJobIdSet.contains(jobId)) {
                        notRunningJobIds.add(jobId);
                    }
                });

                if (CollectionUtils.isEmpty(notRunningJobIds)) {
                    return;
                }

                notRunningJobIds.forEach(jobId -> {
                    Optional<JobInfoDO> jobInfoOpt = jobInfoRepository.findById(jobId);
                    jobInfoOpt.ifPresent(jobInfoDO -> {
                        LifeCycle lifeCycle = LifeCycle.parse(jobInfoDO.getLifecycle());
                        // 生命周期已經(jīng)結(jié)束
                        if (lifeCycle.getEnd() != null && lifeCycle.getEnd() < System.currentTimeMillis()) {
                            jobInfoDO.setStatus(SwitchableStatus.DISABLE.getV());
                            jobInfoDO.setGmtModified(new Date());
                            jobInfoRepository.saveAndFlush(jobInfoDO);
                            log.info("[FrequentScheduler] disable frequent job,id:{}.", jobInfoDO.getId());
                        } else if (lifeCycle.getStart() == null || lifeCycle.getStart() < System.currentTimeMillis() + SCHEDULE_RATE * 2) {
                            log.info("[FrequentScheduler] schedule frequent job,id:{}.", jobInfoDO.getId());
                            jobService.runJob(jobInfoDO.getAppId(), jobId, null, Optional.ofNullable(lifeCycle.getStart()).orElse(0L) - System.currentTimeMillis());
                        }
                    });
                });
            } catch (Exception e) {
                log.error("[FrequentScheduler] schedule frequent job failed.", e);
            }
        });
    }
scheduleFrequentJobCore主要是調(diào)度秒級(jí)任務(wù),它先找出秒級(jí)任務(wù)的id,然后過(guò)濾掉正在運(yùn)行的任務(wù),剩下的未運(yùn)行的任務(wù)挨個(gè)判斷是否需要調(diào)度,需要?jiǎng)t執(zhí)行jobService.runJob

cleanData

public void cleanData() {
        try {
            final List<Long> allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
            if (allAppIds.isEmpty()) {
                return;
            }
            WorkerClusterManagerService.clean(allAppIds);
        } catch (Exception e) {
            log.error("[CleanData] clean data failed.", e);
        }
    }
cleanData主要是通過(guò)WorkerClusterManagerService.clean來(lái)維護(hù)當(dāng)前server負(fù)責(zé)的appId緩存

InstanceStatusCheckService

InstanceStatusCheckService提供了checkRunningInstance、checkWaitingDispatchInstance、checkWaitingWorkerReceiveInstance、checkWorkflowInstance方法

小結(jié)

PowerJob的CoreScheduleTaskManager在afterPropertiesSet的時(shí)候會(huì)啟動(dòng)一系列的線(xiàn)程,它們都是LoopRunnable類(lèi)型的,其中scheduleNormalJob主要是調(diào)度CRON、DAILY_TIME_INTERVAL類(lèi)型的任務(wù),scheduleCronWorkflow主要是調(diào)度CRON 表達(dá)式 WORKFLOW任務(wù),scheduleFrequentJob主要是調(diào)度FIX_RATE/FIX_DELAY 表達(dá)式 JOB。

以上就是PowerJob分布式任務(wù)調(diào)度源碼流程解讀的詳細(xì)內(nèi)容,更多關(guān)于PowerJob分布式任務(wù)調(diào)度的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Java 實(shí)戰(zhàn)項(xiàng)目之畢業(yè)設(shè)計(jì)管理系統(tǒng)的實(shí)現(xiàn)流程

    Java 實(shí)戰(zhàn)項(xiàng)目之畢業(yè)設(shè)計(jì)管理系統(tǒng)的實(shí)現(xiàn)流程

    讀萬(wàn)卷書(shū)不如行萬(wàn)里路,只學(xué)書(shū)上的理論是遠(yuǎn)遠(yuǎn)不夠的,只有在實(shí)戰(zhàn)中才能獲得能力的提升,本篇文章手把手帶你用java+SSM+jsp+mysql+maven實(shí)現(xiàn)畢業(yè)設(shè)計(jì)管理系統(tǒng),大家可以在過(guò)程中查缺補(bǔ)漏,提升水平
    2021-11-11
  • Java中JSch與jsch.addIdentity()完全詳細(xì)解析

    Java中JSch與jsch.addIdentity()完全詳細(xì)解析

    JSch是SSH2的一個(gè)純Java實(shí)現(xiàn),它允許你連接到一個(gè)sshd服務(wù)器,使用端口轉(zhuǎn)發(fā),X11轉(zhuǎn)發(fā),文件傳輸?shù)鹊?這篇文章主要介紹了Java中JSch與jsch.addIdentity()完全詳細(xì)解析,需要的朋友可以參考下
    2026-01-01
  • SpringBoot使用freemarker導(dǎo)出word文件方法詳解

    SpringBoot使用freemarker導(dǎo)出word文件方法詳解

    這篇文章主要介紹了SpringBoot使用freemarker導(dǎo)出word文件方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)吧
    2022-11-11
  • Spring定時(shí)任務(wù)只執(zhí)行一次的原因分析與解決方案

    Spring定時(shí)任務(wù)只執(zhí)行一次的原因分析與解決方案

    在使用Spring的@Scheduled定時(shí)任務(wù)時(shí),你是否遇到過(guò)任務(wù)只執(zhí)行一次,后續(xù)不再觸發(fā)的情況?這種情況可能由多種原因?qū)е?如未啟用調(diào)度、線(xiàn)程池問(wèn)題、異常中斷等,本文將深入分析Spring定時(shí)任務(wù)只執(zhí)行一次的原因,并提供完整的解決方案,需要的朋友可以參考下
    2025-03-03
  • MyBatis啟動(dòng)時(shí)控制臺(tái)無(wú)限輸出日志的原因及解決辦法

    MyBatis啟動(dòng)時(shí)控制臺(tái)無(wú)限輸出日志的原因及解決辦法

    這篇文章主要介紹了MyBatis啟動(dòng)時(shí)控制臺(tái)無(wú)限輸出日志的原因及解決辦法的相關(guān)資料,需要的朋友可以參考下
    2016-07-07
  • Java中`void`和`Void`的區(qū)別和特性詳解

    Java中`void`和`Void`的區(qū)別和特性詳解

    這篇文章主要介紹了Java中`void`和`Void`的區(qū)別和特性的相關(guān)資料,void一般用于函數(shù),表示沒(méi)有返回結(jié)果,是java的一個(gè)關(guān)鍵字,Void是一種類(lèi)型,例如給Void引用賦值null,不可以繼承與實(shí)例化,需要的朋友可以參考下
    2026-03-03
  • java單例模式學(xué)習(xí)示例

    java單例模式學(xué)習(xí)示例

    java中單例模式是一種常見(jiàn)的設(shè)計(jì)模式,單例模式分三種:懶漢式單例、餓漢式單例、登記式單例三種,下面提供了單例模式的示例
    2014-01-01
  • Java異常處理學(xué)習(xí)心得

    Java異常處理學(xué)習(xí)心得

    本篇文章給大家詳細(xì)講述了學(xué)習(xí)Java異常處理學(xué)習(xí)的心得以及原理介紹,對(duì)此有興趣的朋友參考下吧。
    2018-01-01
  • java直接插入排序示例

    java直接插入排序示例

    這篇文章主要介紹了java直接插入排序示例,插入排序的比較次數(shù)仍然是n的平方,但在一般情況下,它要比冒泡排序快一倍,比選擇排序還要快一點(diǎn)。它常常被用在復(fù)雜排序算法的最后階段,比如快速排序。
    2014-05-05
  • 實(shí)例講解分布式緩存軟件Memcached的Java客戶(hù)端使用

    實(shí)例講解分布式緩存軟件Memcached的Java客戶(hù)端使用

    這篇文章主要介紹了分布式緩存軟件Memcached的Java客戶(hù)端使用,Memcached在GitHub上開(kāi)源,作者用其Windows平臺(tái)下的版本進(jìn)行演示,需要的朋友可以參考下
    2016-01-01

最新評(píng)論

禹城市| 炎陵县| 徐闻县| 大埔县| 当涂县| 四平市| 定兴县| 鹰潭市| 佛学| 色达县| 霍邱县| 赤峰市| 金坛市| 西乡县| 通州市| 闵行区| 青海省| 漳浦县| 文山县| 卢龙县| 灵丘县| 渭源县| 内乡县| 烟台市| 临沭县| 土默特左旗| 普兰店市| 大田县| 青神县| 九台市| 阿勒泰市| 青浦区| 夹江县| 赫章县| 古交市| 尼玛县| 鄄城县| 荃湾区| 乃东县| 沛县| 涟源市|