尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

从 0 手写一个巡检调度系统(二):从“能跑”到“敢上线”

从 0 手写一个巡检调度系统(二):从“能跑”到“敢上线” 一、优化背景第一版我们已经实现了Cron 调度数据库分布式锁避免多节点重复执行多规则扩展JobEngine JobExecutor但在实际运行中很快暴露问题❌ 串行执行 → 调度被慢任务拖死❌ 无执行记录 → 无法排查问题❌ 失败只能等 Cron → 恢复慢❌ 全量扫描 → 数据库压力大二、本次优化目标本次升级核心是 4 点并发执行提升吞吐可观测执行日志可恢复失败重试可控批量扫描三、整体执行链路调度线程→ 扫描任务LIMIT→ 提交线程池→ tryLockDB锁→ 执行任务→ 成功 / 失败分支→ 更新 next_run_time→ 写执行日志→ unlock核心认知调度线程 ≠ 执行线程四、并发执行调度器AutowiredprivateInspectionJobServiceinspectionJobService;AutowiredQualifier(inspectionJobExecutor)privateExecutorinspectionJobExecutor;Scheduled(fixedDelayString${inspection.scheduling.scan-interval-ms:10000})publicvoidschedule(){cleanLocks();Listlt;InspectionJobgt;jobsinspectionJobService.findDueJobs();for(InspectionJobjob:jobs){InspectionJobsnapshotjob;inspectionJobExecutor.execute(()-gt;inspectionJobService.runInspectionJob(snapshot));}}线程池配置Bean(nameinspectionJobExecutor)publicExecutorinspectionJobExecutor(InspectionPropertiesproperties){InspectionProperties.Executorcfgproperties.getExecutor();ThreadPoolTaskExecutorexecutornewThreadPoolTaskExecutor();executor.setCorePoolSize(cfg.getCorePoolSize());executor.setMaxPoolSize(cfg.getMaxPoolSize());executor.setQueueCapacity(cfg.getQueueCapacity());executor.setThreadNamePrefix(cfg.getThreadNamePrefix());executor.setRejectedExecutionHandler(newThreadPoolExecutor.CallerRunsPolicy());executor.initialize();returnexecutor;}五、核心执行逻辑重点publicvoidrunInspectionJob(InspectionJobinspectionJob){longstartMsSystem.currentTimeMillis();StringnodeIdresolveNodeId();booleanlockedfalse;intmaxRetryinspectionJob.getMaxRetry()!null?inspectionJob.getMaxRetry():3;intretryIntervalinspectionJob.getRetryIntervalSeconds()!null?inspectionJob.getRetryIntervalSeconds():60;intpriorFailuresinspectionJob.getRetryCount()!null?inspectionJob.getRetryCount():0;try{intlockRowsinspectionJobMapper.tryLock(inspectionJob.getJobId(),nodeId);if(lockRows0){return;}lockedtrue;// 执行任务jobEngine.executeRule(inspectionJob.getInspectionId(),inspectionJob);// 成功清空重试次数inspectionJobMapper.resetRetryCount(inspectionJob.getJobId());// 计算下一次 cronLocalDateTimenextRunTimeCronUtils.calcNextRunTime(inspectionJob.getCronExpr());inspectionJobMapper.updateNextRunTime(inspectionJob.getJobId(),nextRunTime,LocalDateTime.now());longdurationSystem.currentTimeMillis()-startMs;jobExecutionLogService.logSuccess(inspectionJob,nodeId,priorFailures1,duration);}catch(Exceptione){longdurationSystem.currentTimeMillis()-startMs;intfailOrdinalpriorFailures1;if(failOrdinallt;maxRetry){// 计算重试时间指数退避LocalDateTimeretryAtCronUtils.calcRetryAfter(LocalDateTime.now(),retryInterval,failOrdinal);inspectionJobMapper.scheduleRetry(inspectionJob.getJobId(),failOrdinal,retryAt);jobExecutionLogService.logFailure(inspectionJob,nodeId,failOrdinal,duration,e,true);}else{// 超过最大重试 → 回到 croninspectionJobMapper.resetRetryCount(inspectionJob.getJobId());LocalDateTimenextRunTimeCronUtils.calcNextRunTime(inspectionJob.getCronExpr());inspectionJobMapper.updateNextRunTime(inspectionJob.getJobId(),nextRunTime,LocalDateTime.now());jobExecutionLogService.logFailure(inspectionJob,nodeId,failOrdinal,duration,e,false);}}finally{if(locked){inspectionJobMapper.unlock(inspectionJob.getJobId(),nodeId);}}}六、执行日志publicvoidlogSuccess(InspectionJobjob,StringnodeId,intattemptNumber,longdurationMs){JobExecutionLoglogbase(job,nodeId,attemptNumber,durationMs);log.setStatus(SUCCESS);log.setWillRetry(false);jobExecutionLogMapper.insert(log);}publicvoidlogFailure(InspectionJobjob,StringnodeId,intattemptNumber,longdurationMs,Throwableerror,booleanwillRetry){JobExecutionLoglogbase(job,nodeId,attemptNumber,durationMs);log.setStatus(FAILED);log.setWillRetry(willRetry);log.setErrorMessage(error.getMessage());jobExecutionLogMapper.insert(log);}七、指数退避算法publicstaticLocalDateTimecalcRetryAfter(LocalDateTimefrom,intbaseIntervalSecs,intfailureOrdinal){intexpMath.max(0,failureOrdinal-1);longmult1LMath.min(exp,12);longsecondsMath.min((long)baseIntervalSecs*mult,3600);returnfrom.plusSeconds(seconds);}八、调度优化SELECT*FROMinspection_jobWHEREstatusRUNNINGANDnext_run_timelt;NOW()ORDERBYnext_run_timeASCLIMIT?九、关键设计点tryLock保证只有一个节点执行unlock必须带 nodeIdretry_count必须数据库更新线程池负责并发能力CallerRunsPolicy天然限流十、总结维度第一版第二版执行方式串行并发失败处理等 cron重试可观测无日志扫描方式全量LIMIT一句话总结第一版解决“能不能跑” 第二版解决“能不能稳定跑”
返回列表