From 11b47a604f610fa751ae26adf8c5beefc10219a0 Mon Sep 17 00:00:00 2001 From: tianwenbo Date: Fri, 10 Mar 2023 17:11:05 +0800 Subject: [PATCH] =?UTF-8?q?=E3=80=90fix=20sonar=E3=80=91=20JobScheduleHelp?= =?UTF-8?q?er.java=E6=96=87=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../admin/core/thread/JobScheduleHelper.java | 44 +++++++++---------- 1 file changed, 21 insertions(+), 23 deletions(-) diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobScheduleHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobScheduleHelper.java index e425b323..ed664337 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobScheduleHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobScheduleHelper.java @@ -32,32 +32,27 @@ public class JobScheduleHelper { private Thread ringThread; private volatile boolean scheduleThreadToStop = false; private volatile boolean ringThreadToStop = false; - private volatile static Map> ringData = new ConcurrentHashMap<>(); + private static volatile Map> ringData = new ConcurrentHashMap<>(); public void start(){ // schedule thread - scheduleThread = new Thread(new Runnable() { - @Override - public void run() { - + scheduleThread = new Thread(()-> { try { TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 ); } catch (InterruptedException e) { if (!scheduleThreadToStop) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } logger.info(">>>>>>>>> init xxl-job admin scheduler success."); - // pre-read count: treadpool-size * trigger-qps (each trigger cost 50ms, qps = 1000/50 = 20) int preReadCount = (XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax() + XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax()) * 20; while (!scheduleThreadToStop) { - // Scan Job long start = System.currentTimeMillis(); - Connection conn = null; Boolean connAutoCommit = null; PreparedStatement preparedStatement = null; @@ -77,7 +72,7 @@ public class JobScheduleHelper { // 1、pre read long nowTime = System.currentTimeMillis(); List scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount); - if (scheduleList!=null && scheduleList.size()>0) { + if (scheduleList!=null && !scheduleList.isEmpty()) { // 2、push time-ring for (XxlJobInfo jobInfo: scheduleList) { @@ -94,7 +89,7 @@ public class JobScheduleHelper { // 1、trigger JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null); - logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = " + jobInfo.getId() ); + logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = {}" , jobInfo.getId() ); // 2、fresh next refreshNextValidTime(jobInfo, new Date()); @@ -194,14 +189,12 @@ public class JobScheduleHelper { } catch (InterruptedException e) { if (!scheduleThreadToStop) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } } - } - logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop"); - } }); scheduleThread.setDaemon(true); scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread"); @@ -209,16 +202,14 @@ public class JobScheduleHelper { // ring thread - ringThread = new Thread(new Runnable() { - @Override - public void run() { - + ringThread = new Thread(()-> { // align second try { TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000 ); } catch (InterruptedException e) { if (!ringThreadToStop) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } @@ -237,7 +228,7 @@ public class JobScheduleHelper { // ring trigger logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) ); - if (ringItemData.size() > 0) { + if (!ringItemData.isEmpty()) { // do trigger for (int jobId: ringItemData) { // do trigger @@ -258,11 +249,12 @@ public class JobScheduleHelper { } catch (InterruptedException e) { if (!ringThreadToStop) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } + } } logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop"); - } }); ringThread.setDaemon(true); ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread"); @@ -285,7 +277,7 @@ public class JobScheduleHelper { // push async ring List ringItemData = ringData.get(ringSecond); if (ringItemData == null) { - ringItemData = new ArrayList(); + ringItemData = new ArrayList<>(); ringData.put(ringSecond, ringItemData); } ringItemData.add(jobId); @@ -301,6 +293,7 @@ public class JobScheduleHelper { TimeUnit.SECONDS.sleep(1); // wait } catch (InterruptedException e) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } if (scheduleThread.getState() != Thread.State.TERMINATED){ // interrupt and wait @@ -309,15 +302,17 @@ public class JobScheduleHelper { scheduleThread.join(); } catch (InterruptedException e) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } // if has ring data boolean hasRingData = false; if (!ringData.isEmpty()) { - for (int second : ringData.keySet()) { - List tmpData = ringData.get(second); - if (tmpData!=null && tmpData.size()>0) { + Set>> entries = ringData.entrySet(); + for (Map.Entry> e : entries) { + List tmpData = e.getValue(); + if (tmpData!=null && !tmpData.isEmpty()) { hasRingData = true; break; } @@ -328,6 +323,7 @@ public class JobScheduleHelper { TimeUnit.SECONDS.sleep(8); } catch (InterruptedException e) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } @@ -337,6 +333,7 @@ public class JobScheduleHelper { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } if (ringThread.getState() != Thread.State.TERMINATED){ // interrupt and wait @@ -345,6 +342,7 @@ public class JobScheduleHelper { ringThread.join(); } catch (InterruptedException e) { logger.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } }