【fix sonar】 JobScheduleHelper.java文件
This commit is contained in:
+21
-23
@@ -32,32 +32,27 @@ public class JobScheduleHelper {
|
|||||||
private Thread ringThread;
|
private Thread ringThread;
|
||||||
private volatile boolean scheduleThreadToStop = false;
|
private volatile boolean scheduleThreadToStop = false;
|
||||||
private volatile boolean ringThreadToStop = false;
|
private volatile boolean ringThreadToStop = false;
|
||||||
private volatile static Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>();
|
private static volatile Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
public void start(){
|
public void start(){
|
||||||
|
|
||||||
// schedule thread
|
// schedule thread
|
||||||
scheduleThread = new Thread(new Runnable() {
|
scheduleThread = new Thread(()-> {
|
||||||
@Override
|
|
||||||
public void run() {
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 );
|
TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 );
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
if (!scheduleThreadToStop) {
|
if (!scheduleThreadToStop) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
logger.info(">>>>>>>>> init xxl-job admin scheduler success.");
|
logger.info(">>>>>>>>> init xxl-job admin scheduler success.");
|
||||||
|
|
||||||
// pre-read count: treadpool-size * trigger-qps (each trigger cost 50ms, qps = 1000/50 = 20)
|
// pre-read count: treadpool-size * trigger-qps (each trigger cost 50ms, qps = 1000/50 = 20)
|
||||||
int preReadCount = (XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax() + XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax()) * 20;
|
int preReadCount = (XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax() + XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax()) * 20;
|
||||||
|
|
||||||
while (!scheduleThreadToStop) {
|
while (!scheduleThreadToStop) {
|
||||||
|
|
||||||
// Scan Job
|
// Scan Job
|
||||||
long start = System.currentTimeMillis();
|
long start = System.currentTimeMillis();
|
||||||
|
|
||||||
Connection conn = null;
|
Connection conn = null;
|
||||||
Boolean connAutoCommit = null;
|
Boolean connAutoCommit = null;
|
||||||
PreparedStatement preparedStatement = null;
|
PreparedStatement preparedStatement = null;
|
||||||
@@ -77,7 +72,7 @@ public class JobScheduleHelper {
|
|||||||
// 1、pre read
|
// 1、pre read
|
||||||
long nowTime = System.currentTimeMillis();
|
long nowTime = System.currentTimeMillis();
|
||||||
List<XxlJobInfo> scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount);
|
List<XxlJobInfo> 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
|
// 2、push time-ring
|
||||||
for (XxlJobInfo jobInfo: scheduleList) {
|
for (XxlJobInfo jobInfo: scheduleList) {
|
||||||
|
|
||||||
@@ -94,7 +89,7 @@ public class JobScheduleHelper {
|
|||||||
|
|
||||||
// 1、trigger
|
// 1、trigger
|
||||||
JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null);
|
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
|
// 2、fresh next
|
||||||
refreshNextValidTime(jobInfo, new Date());
|
refreshNextValidTime(jobInfo, new Date());
|
||||||
@@ -194,14 +189,12 @@ public class JobScheduleHelper {
|
|||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
if (!scheduleThreadToStop) {
|
if (!scheduleThreadToStop) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop");
|
logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop");
|
||||||
}
|
|
||||||
});
|
});
|
||||||
scheduleThread.setDaemon(true);
|
scheduleThread.setDaemon(true);
|
||||||
scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
|
scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
|
||||||
@@ -209,16 +202,14 @@ public class JobScheduleHelper {
|
|||||||
|
|
||||||
|
|
||||||
// ring thread
|
// ring thread
|
||||||
ringThread = new Thread(new Runnable() {
|
ringThread = new Thread(()-> {
|
||||||
@Override
|
|
||||||
public void run() {
|
|
||||||
|
|
||||||
// align second
|
// align second
|
||||||
try {
|
try {
|
||||||
TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000 );
|
TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000 );
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
if (!ringThreadToStop) {
|
if (!ringThreadToStop) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -237,7 +228,7 @@ public class JobScheduleHelper {
|
|||||||
|
|
||||||
// ring trigger
|
// ring trigger
|
||||||
logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) );
|
logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) );
|
||||||
if (ringItemData.size() > 0) {
|
if (!ringItemData.isEmpty()) {
|
||||||
// do trigger
|
// do trigger
|
||||||
for (int jobId: ringItemData) {
|
for (int jobId: ringItemData) {
|
||||||
// do trigger
|
// do trigger
|
||||||
@@ -258,11 +249,12 @@ public class JobScheduleHelper {
|
|||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
if (!ringThreadToStop) {
|
if (!ringThreadToStop) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
|
logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
|
||||||
}
|
|
||||||
});
|
});
|
||||||
ringThread.setDaemon(true);
|
ringThread.setDaemon(true);
|
||||||
ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
|
ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
|
||||||
@@ -285,7 +277,7 @@ public class JobScheduleHelper {
|
|||||||
// push async ring
|
// push async ring
|
||||||
List<Integer> ringItemData = ringData.get(ringSecond);
|
List<Integer> ringItemData = ringData.get(ringSecond);
|
||||||
if (ringItemData == null) {
|
if (ringItemData == null) {
|
||||||
ringItemData = new ArrayList<Integer>();
|
ringItemData = new ArrayList<>();
|
||||||
ringData.put(ringSecond, ringItemData);
|
ringData.put(ringSecond, ringItemData);
|
||||||
}
|
}
|
||||||
ringItemData.add(jobId);
|
ringItemData.add(jobId);
|
||||||
@@ -301,6 +293,7 @@ public class JobScheduleHelper {
|
|||||||
TimeUnit.SECONDS.sleep(1); // wait
|
TimeUnit.SECONDS.sleep(1); // wait
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
if (scheduleThread.getState() != Thread.State.TERMINATED){
|
if (scheduleThread.getState() != Thread.State.TERMINATED){
|
||||||
// interrupt and wait
|
// interrupt and wait
|
||||||
@@ -309,15 +302,17 @@ public class JobScheduleHelper {
|
|||||||
scheduleThread.join();
|
scheduleThread.join();
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// if has ring data
|
// if has ring data
|
||||||
boolean hasRingData = false;
|
boolean hasRingData = false;
|
||||||
if (!ringData.isEmpty()) {
|
if (!ringData.isEmpty()) {
|
||||||
for (int second : ringData.keySet()) {
|
Set<Map.Entry<Integer, List<Integer>>> entries = ringData.entrySet();
|
||||||
List<Integer> tmpData = ringData.get(second);
|
for (Map.Entry<Integer, List<Integer>> e : entries) {
|
||||||
if (tmpData!=null && tmpData.size()>0) {
|
List<Integer> tmpData = e.getValue();
|
||||||
|
if (tmpData!=null && !tmpData.isEmpty()) {
|
||||||
hasRingData = true;
|
hasRingData = true;
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -328,6 +323,7 @@ public class JobScheduleHelper {
|
|||||||
TimeUnit.SECONDS.sleep(8);
|
TimeUnit.SECONDS.sleep(8);
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -337,6 +333,7 @@ public class JobScheduleHelper {
|
|||||||
TimeUnit.SECONDS.sleep(1);
|
TimeUnit.SECONDS.sleep(1);
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
if (ringThread.getState() != Thread.State.TERMINATED){
|
if (ringThread.getState() != Thread.State.TERMINATED){
|
||||||
// interrupt and wait
|
// interrupt and wait
|
||||||
@@ -345,6 +342,7 @@ public class JobScheduleHelper {
|
|||||||
ringThread.join();
|
ringThread.join();
|
||||||
} catch (InterruptedException e) {
|
} catch (InterruptedException e) {
|
||||||
logger.error(e.getMessage(), e);
|
logger.error(e.getMessage(), e);
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user