diff --git a/jero-boot/jero-cloud-module/jero-cloud-gateway/src/main/java/com/jero/loader/DynamicRouteLoader.java b/jero-boot/jero-cloud-module/jero-cloud-gateway/src/main/java/com/jero/loader/DynamicRouteLoader.java index d0e67744..733ec581 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-gateway/src/main/java/com/jero/loader/DynamicRouteLoader.java +++ b/jero-boot/jero-cloud-module/jero-cloud-gateway/src/main/java/com/jero/loader/DynamicRouteLoader.java @@ -60,8 +60,9 @@ public class DynamicRouteLoader implements ApplicationEventPublisherAware { public DynamicRouteLoader(InMemoryRouteDefinitionRepository repository, DynamicRouteService dynamicRouteService, RedisUtil redisUtil) { - - this.repository = repository; + if(this.repository == null){ + this.repository = repository; + } this.dynamicRouteService = dynamicRouteService; this.redisUtil = redisUtil; } @@ -234,49 +235,6 @@ public class DynamicRouteLoader implements ApplicationEventPublisherAware { } -// private void loadRoutesByDataBase() { -// List routeList = jdbcTemplate.query(SELECT_ROUTES, new RowMapper() { -// @Override -// public GatewayRouteVo mapRow(ResultSet rs, int i) throws SQLException { -// GatewayRouteVo result = new GatewayRouteVo(); -// result.setId(rs.getString("id")); -// result.setName(rs.getString("name")); -// result.setUri(rs.getString("uri")); -// result.setStatus(rs.getInt("status")); -// result.setRetryable(rs.getInt("retryable")); -// result.setPredicates(rs.getString("predicates")); -// result.setStripPrefix(rs.getInt("strip_prefix")); -// result.setPersist(rs.getInt("persist")); -// return result; -// } -// }); -// if (ObjectUtil.isNotEmpty(routeList)) { -// // 加载路由 -// routeList.forEach(route -> { -// RouteDefinition definition = new RouteDefinition(); -// List predicatesList = Lists.newArrayList(); -// List filtersList = Lists.newArrayList(); -// definition.setId(route.getId()); -// String predicates = route.getPredicates(); -// String filters = route.getFilters(); -// if (StringUtils.isNotEmpty(predicates)) { -// predicatesList = JSON.parseArray(predicates, PredicateDefinition.class); -// definition.setPredicates(predicatesList); -// } -// if (StringUtils.isNotEmpty(filters)) { -// filtersList = JSON.parseArray(filters, FilterDefinition.class); -// definition.setFilters(filtersList); -// } -// URI uri = UriComponentsBuilder.fromUriString(route.getUri()).build().toUri(); -// definition.setUri(uri); -// this.repository.save(Mono.just(definition)).subscribe(); -// }); -// log.info("加载路由:{}==============", routeList.size()); -// Mono.empty(); -// } -// } - - /** * 监听Nacos下发的动态路由配置 * diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/IndexController.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/IndexController.java index 61891e37..3895f2dd 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/IndexController.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/IndexController.java @@ -73,11 +73,6 @@ public class IndexController { @RequestMapping("/help") public String help() { - - /*if (!PermissionInterceptor.ifLogin(request)) { - return "redirect:/toLogin"; - }*/ - return "help"; } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobCodeController.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobCodeController.java index b1eb7365..927f56da 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobCodeController.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobCodeController.java @@ -36,10 +36,10 @@ public class JobCodeController { List jobLogGlues = xxlJobLogGlueDao.findByJobId(jobId); if (jobInfo == null) { - throw new RuntimeException(I18nUtil.getString("jobinfo_glue_jobid_unvalid")); + throw new IllegalArgumentException(I18nUtil.getString("jobinfo_glue_jobid_unvalid")); } if (GlueTypeEnum.BEAN == GlueTypeEnum.match(jobInfo.getGlueType())) { - throw new RuntimeException(I18nUtil.getString("jobinfo_glue_gluetype_unvalid")); + throw new IllegalArgumentException(I18nUtil.getString("jobinfo_glue_gluetype_unvalid")); } // valid permission diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobGroupController.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobGroupController.java index 0d6c0269..79e03d03 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobGroupController.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobGroupController.java @@ -74,14 +74,9 @@ public class JobGroupController { return new ReturnT<>(500, (I18nUtil.getString(systemPleaseInput) + I18nUtil.getString("jobgroup_field_title")) ); } if (xxlJobGroup.getAddressType()!=0) { - if (xxlJobGroup.getAddressList()==null || xxlJobGroup.getAddressList().trim().length()==0) { - return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_addressType_limit") ); - } - String[] addresss = xxlJobGroup.getAddressList().split(","); - for (String item: addresss) { - if (item==null || item.trim().length()==0) { - return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_registryList_unvalid") ); - } + ReturnT x = getStringReturnT(xxlJobGroup); + if (x != null) { + return x; } } @@ -89,6 +84,19 @@ public class JobGroupController { return (ret>0)?ReturnT.SUCCESS:ReturnT.FAIL; } + private ReturnT getStringReturnT(XxlJobGroup xxlJobGroup) { + if (xxlJobGroup.getAddressList()==null || xxlJobGroup.getAddressList().trim().length()==0) { + return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_addressType_limit") ); + } + String[] addresss = xxlJobGroup.getAddressList().split(","); + for (String item: addresss) { + if (item==null || item.trim().length()==0) { + return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_registryList_unvalid") ); + } + } + return null; + } + @RequestMapping("/update") @ResponseBody public ReturnT update(XxlJobGroup xxlJobGroup){ @@ -105,33 +113,44 @@ public class JobGroupController { if (xxlJobGroup.getAddressType() == 0) { // 0=自动注册 List registryList = findRegistryByAppName(xxlJobGroup.getAppname()); - StringBuilder addressListStr = new StringBuilder(); - if (registryList!=null && !registryList.isEmpty()) { - Collections.sort(registryList); - addressListStr = new StringBuilder(); - for (String item:registryList) { - addressListStr.append(item).append(","); - } - addressListStr = new StringBuilder(addressListStr.substring(0, addressListStr.length() - 1)); - } + StringBuilder addressListStr = getStringBuilder(registryList); xxlJobGroup.setAddressList(addressListStr.toString()); } else { // 1=手动录入 if (xxlJobGroup.getAddressList()==null || xxlJobGroup.getAddressList().trim().length()==0) { return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_addressType_limit") ); } - String[] addresss = xxlJobGroup.getAddressList().split(","); - for (String item: addresss) { - if (item==null || item.trim().length()==0) { - return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_registryList_unvalid") ); - } - } + if (getAddresss(xxlJobGroup)) + return new ReturnT<>(500, I18nUtil.getString("jobgroup_field_registryList_unvalid")); } int ret = xxlJobGroupDao.update(xxlJobGroup); return (ret>0)?ReturnT.SUCCESS:ReturnT.FAIL; } + private boolean getAddresss(XxlJobGroup xxlJobGroup) { + String[] addresss = xxlJobGroup.getAddressList().split(","); + for (String item : addresss) { + if (item == null || item.trim().length() == 0) { + return true; + } + } + return false; + } + + private StringBuilder getStringBuilder(List registryList) { + StringBuilder addressListStr = new StringBuilder(); + if (registryList!=null && !registryList.isEmpty()) { + Collections.sort(registryList); + addressListStr = new StringBuilder(); + for (String item:registryList) { + addressListStr.append(item).append(","); + } + addressListStr = new StringBuilder(addressListStr.substring(0, addressListStr.length() - 1)); + } + return addressListStr; + } + private List findRegistryByAppName(String appnameParam){ HashMap> appAddressMap = new HashMap<>(); List list = xxlJobRegistryDao.findAll(RegistryConfig.DEAD_TIMEOUT, new Date()); diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobInfoController.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobInfoController.java index 96c631f4..d65417eb 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobInfoController.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobInfoController.java @@ -75,18 +75,23 @@ public class JobInfoController { groupIdStrs = Arrays.asList(loginUser.getPermission().trim().split(",")); } for (XxlJobGroup groupItem:jobGroupListAll) { - if (groupIdStrs.contains(String.valueOf(groupItem.getId()))) { - jobGroupList.add(groupItem); - } + getAdd(jobGroupList, groupIdStrs, groupItem); } } } return jobGroupList; } + + private static void getAdd(List jobGroupList, List groupIdStrs, XxlJobGroup groupItem) { + if (groupIdStrs.contains(String.valueOf(groupItem.getId()))) { + jobGroupList.add(groupItem); + } + } + public static void validPermission(HttpServletRequest request, int jobGroup) { XxlJobUser loginUser = (XxlJobUser) request.getAttribute(LoginService.LOGIN_IDENTITY_KEY); if (!loginUser.validPermission(jobGroup)) { - throw new RuntimeException(I18nUtil.getString("system_permission_limit") + "[username="+ loginUser.getUsername() +"]"); + throw new IllegalArgumentException(I18nUtil.getString("system_permission_limit") + "[username="+ loginUser.getUsername() +"]"); } } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobLogController.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobLogController.java index 3cc01cd2..036dbe4d 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobLogController.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/JobLogController.java @@ -64,7 +64,7 @@ public class JobLogController { if (jobId > 0) { XxlJobInfo jobInfo = xxlJobInfoDao.loadById(jobId); if (jobInfo == null) { - throw new RuntimeException(I18nUtil.getString("jobinfo_field_id") + I18nUtil.getString("system_unvalid")); + throw new IllegalArgumentException(I18nUtil.getString("jobinfo_field_id") + I18nUtil.getString("system_unvalid")); } model.addAttribute("jobInfo", jobInfo); @@ -120,10 +120,9 @@ public class JobLogController { public String logDetailPage(int id, Model model){ // base check -// ReturnT logStatue = ReturnT.SUCCESS; XxlJobLog jobLog = xxlJobLogDao.load(id); if (jobLog == null) { - throw new RuntimeException(I18nUtil.getString("joblog_logid_unvalid")); + throw new IllegalArgumentException(I18nUtil.getString("joblog_logid_unvalid")); } model.addAttribute("triggerCode", jobLog.getTriggerCode()); diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/interceptor/PermissionInterceptor.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/interceptor/PermissionInterceptor.java index 3584c129..73d9aece 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/interceptor/PermissionInterceptor.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/interceptor/PermissionInterceptor.java @@ -25,7 +25,7 @@ public class PermissionInterceptor extends HandlerInterceptorAdapter { @Override public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception { - + if (!(handler instanceof HandlerMethod)) { return super.preHandle(request, response, handler); } @@ -44,16 +44,15 @@ public class PermissionInterceptor extends HandlerInterceptorAdapter { XxlJobUser loginUser = loginService.ifLogin(request, response); if (loginUser == null) { response.sendRedirect(request.getContextPath() + "/toLogin"); - //request.getRequestDispatcher("/toLogin").forward(request, response); return false; } if (needAdminuser && loginUser.getRole()!=1) { - throw new RuntimeException(I18nUtil.getString("system_permission_limit")); + throw new IllegalArgumentException(I18nUtil.getString("system_permission_limit")); } request.setAttribute(LoginService.LOGIN_IDENTITY_KEY, loginUser); } return super.preHandle(request, response, handler); } - + } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/resolver/WebExceptionResolver.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/resolver/WebExceptionResolver.java index 07de25fd..d5448d76 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/resolver/WebExceptionResolver.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/controller/resolver/WebExceptionResolver.java @@ -1,10 +1,9 @@ package com.xxl.job.admin.controller.resolver; import com.xxl.job.admin.core.exception.XxlJobException; -import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.admin.core.util.JacksonUtil; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import com.xxl.job.core.biz.model.ReturnT; +import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.method.HandlerMethod; @@ -20,16 +19,16 @@ import java.io.IOException; * * @author xuxueli 2016-1-6 19:22:18 */ +@Slf4j @Component public class WebExceptionResolver implements HandlerExceptionResolver { - private static Logger logger = LoggerFactory.getLogger(WebExceptionResolver.class); @Override public ModelAndView resolveException(HttpServletRequest request, HttpServletResponse response, Object handler, Exception ex) { if (!(ex instanceof XxlJobException)) { - logger.error("WebExceptionResolver:{}", ex); + log.error("WebExceptionResolver:{}", ex); } // if json @@ -50,7 +49,7 @@ public class WebExceptionResolver implements HandlerExceptionResolver { response.setContentType("application/json;charset=utf-8"); response.getWriter().print(JacksonUtil.writeValueAsString(errorResult)); } catch (IOException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); } return mv; } else { diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/alarm/impl/EmailJobAlarm.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/alarm/impl/EmailJobAlarm.java index 2f2bebe9..5fa2bf24 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/alarm/impl/EmailJobAlarm.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/alarm/impl/EmailJobAlarm.java @@ -27,6 +27,8 @@ import java.util.Set; public class EmailJobAlarm implements JobAlarm { private static Logger logger = LoggerFactory.getLogger(EmailJobAlarm.class); + private static final String TD = "\n"; + /** * fail alarm * @@ -94,11 +96,11 @@ public class EmailJobAlarm implements JobAlarm { "\n" + " " + " \n" + - " \n" + - " \n" + - " \n" + - " \n" + - " \n" + + " \n" + " \n" + " \n" + @@ -106,7 +108,7 @@ public class EmailJobAlarm implements JobAlarm { " \n" + " \n" + " \n" + - " \n" + + " \n" + " \n" + " \n" + diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/conf/XxlJobAdminConfig.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/conf/XxlJobAdminConfig.java index 1b0405ac..dd385f12 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/conf/XxlJobAdminConfig.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/conf/XxlJobAdminConfig.java @@ -34,8 +34,6 @@ public class XxlJobAdminConfig implements InitializingBean, DisposableBean { @Override public void afterPropertiesSet() throws Exception { - adminConfig = this; - xxlJobScheduler = new XxlJobScheduler(); xxlJobScheduler.init(); } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/cron/CronExpression.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/cron/CronExpression.java index 6816b883..384a42df 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/cron/CronExpression.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/cron/CronExpression.java @@ -17,7 +17,7 @@ package com.xxl.job.admin.core.cron; -import lombok.val; +import lombok.extern.slf4j.Slf4j; import java.io.Serializable; import java.text.ParseException; @@ -190,6 +190,7 @@ import java.util.*; *

* Borrowed from quartz v2.3.1 */ +@Slf4j public final class CronExpression implements Serializable, Cloneable { private static final long serialVersionUID = 12423409423L; @@ -232,7 +233,7 @@ public final class CronExpression implements Serializable, Cloneable { dayMap.put("SAT", 7); } - private final String cronExpression; + private final String cronExpressionData; private TimeZone timeZone = null; protected transient TreeSet seconds; protected transient TreeSet minutes; @@ -265,9 +266,9 @@ public final class CronExpression implements Serializable, Cloneable { throw new IllegalArgumentException("cronExpression cannot be null"); } - this.cronExpression = cronExpression.toUpperCase(Locale.US); + this.cronExpressionData = cronExpression.toUpperCase(Locale.US); - buildExpression(this.cronExpression); + buildExpression(this.cronExpressionData); } /** @@ -282,9 +283,9 @@ public final class CronExpression implements Serializable, Cloneable { * ParseException. We also elide some of the sanity checking as it is * not logically trippable. */ - this.cronExpression = expression.getCronExpression(); + this.cronExpressionData = expression.getCronExpression(); try { - buildExpression(cronExpression); + buildExpression(cronExpressionData); } catch (ParseException ex) { throw new AssertionError(); } @@ -394,7 +395,7 @@ public final class CronExpression implements Serializable, Cloneable { */ @Override public String toString() { - return cronExpression; + return cronExpressionData; } /** @@ -433,27 +434,7 @@ public final class CronExpression implements Serializable, Cloneable { try { - if (seconds == null) { - seconds = new TreeSet<>(); - } - if (minutes == null) { - minutes = new TreeSet<>(); - } - if (hours == null) { - hours = new TreeSet<>(); - } - if (daysOfMonth == null) { - daysOfMonth = new TreeSet<>(); - } - if (months == null) { - months = new TreeSet<>(); - } - if (daysOfWeek == null) { - daysOfWeek = new TreeSet<>(); - } - if (years == null) { - years = new TreeSet<>(); - } + getTime(); int exprOn = SECOND; @@ -462,18 +443,7 @@ public final class CronExpression implements Serializable, Cloneable { while (exprsTok.hasMoreTokens() && exprOn <= YEAR) { String expr = exprsTok.nextToken().trim(); - - // throw an exception if L is used with other days of the month - if (exprOn == DAY_OF_MONTH && expr.indexOf('L') != -1 && expr.length() > 1 && expr.contains(",")) { - throw new ParseException("Support for specifying 'L' and 'LW' with other days of the month is not implemented", -1); - } - // throw an exception if L is used with other days of the week - if (exprOn == DAY_OF_WEEK && expr.indexOf('L') != -1 && expr.length() > 1 && expr.contains(",")) { - throw new ParseException("Support for specifying 'L' with other days of the week is not implemented", -1); - } - if (exprOn == DAY_OF_WEEK && expr.indexOf('#') != -1 && expr.indexOf('#', expr.indexOf('#') + 1) != -1) { - throw new ParseException("Support for specifying multiple \"nth\" days is not implemented.", -1); - } + getExprOn(exprOn, expr); StringTokenizer vTok = new StringTokenizer(expr, ","); while (vTok.hasMoreTokens()) { @@ -512,6 +482,44 @@ public final class CronExpression implements Serializable, Cloneable { } } + private void getTime() { + if (seconds == null) { + seconds = new TreeSet<>(); + } + if (minutes == null) { + minutes = new TreeSet<>(); + } + if (hours == null) { + hours = new TreeSet<>(); + } + if (daysOfMonth == null) { + daysOfMonth = new TreeSet<>(); + } + if (months == null) { + months = new TreeSet<>(); + } + if (daysOfWeek == null) { + daysOfWeek = new TreeSet<>(); + } + if (years == null) { + years = new TreeSet<>(); + } + } + + private void getExprOn(int exprOn, String expr) throws ParseException { + // throw an exception if L is used with other days of the month + if (exprOn == DAY_OF_MONTH && expr.indexOf('L') != -1 && expr.length() > 1 && expr.contains(",")) { + throw new ParseException("Support for specifying 'L' and 'LW' with other days of the month is not implemented", -1); + } + // throw an exception if L is used with other days of the week + if (exprOn == DAY_OF_WEEK && expr.indexOf('L') != -1 && expr.length() > 1 && expr.contains(",")) { + throw new ParseException("Support for specifying 'L' with other days of the week is not implemented", -1); + } + if (exprOn == DAY_OF_WEEK && expr.indexOf('#') != -1 && expr.indexOf('#', expr.indexOf('#') + 1) != -1) { + throw new ParseException("Support for specifying multiple \"nth\" days is not implemented.", -1); + } + } + protected int storeExpressionVals(int pos, String s, int type) throws ParseException { @@ -523,154 +531,72 @@ public final class CronExpression implements Serializable, Cloneable { char c = s.charAt(i); if ((c >= 'A') && (c <= 'Z') && (!s.equals("L")) && (!s.equals("LW")) && (!s.matches("^L-\\d*W?"))) { String sub = s.substring(i, i + 3); - int sval = -1; - int eval = -1; - if (type == MONTH) { - sval = getMonthNumber(sub) + 1; - if (sval <= 0) { - throw new ParseException("Invalid Month value: '" + sub + "'", i); - } - if (s.length() > i + 3) { - c = s.charAt(i + 3); - if (c == '-') { - i += 4; - sub = s.substring(i, i + 3); - eval = getMonthNumber(sub) + 1; - if (eval <= 0) { - throw new ParseException("Invalid Month value: '" + sub + "'", i); - } - } - } - } else if (type == DAY_OF_WEEK) { - sval = getDayOfWeekNumber(sub); - if (sval < 0) { - throw new ParseException("Invalid Day-of-Week value: '" - + sub + "'", i); - } - if (s.length() > i + 3) { - c = s.charAt(i + 3); - if (c == '-') { - i += 4; - sub = s.substring(i, i + 3); - eval = getDayOfWeekNumber(sub); - if (eval < 0) { - throw new ParseException( - "Invalid Day-of-Week value: '" + sub - + "'", i); - } - } else if (c == '#') { - try { - i += 4; - nthdayOfWeek = Integer.parseInt(s.substring(i)); - if (nthdayOfWeek < 1 || nthdayOfWeek > 5) { - throw new Exception(); - } - } catch (Exception e) { - throw new ParseException( - "A numeric value between 1 and 5 must follow the '#' option", - i); - } - } else if (c == 'L') { - lastdayOfWeek = true; - i++; - } - } - - } else { - throw new ParseException( - "Illegal characters for this position: '" + sub + "'", - i); - } - if (eval != -1) { - incr = 1; - } - addToSet(sval, eval, incr, type); - return (i + 3); + return getReturn(s, type, incr, i, sub); } if (c == '?') { - i++; - if ((i + 1) < s.length() - && (s.charAt(i) != ' ' && s.charAt(i + 1) != '\t')) { - throw new ParseException("Illegal character after '?': " - + s.charAt(i), i); - } - if (type != DAY_OF_WEEK && type != DAY_OF_MONTH) { - throw new ParseException( - "'?' can only be specified for Day-of-Month or Day-of-Week.", - i); - } - if (type == DAY_OF_WEEK && !lastdayOfMonth) { - int val = daysOfMonth.last(); - if (val == NO_SPEC_INT) { - throw new ParseException( - "'?' can only be specified for Day-of-Month -OR- Day-of-Week.", - i); - } - } - - addToSet(NO_SPEC_INT, -1, 0, type); - return i; + return getType(s, type, i); } - if (c == '*' || c == '/') { - if (c == '*' && (i + 1) >= s.length()) { - addToSet(ALL_SPEC_INT, -1, incr, type); - return i + 1; - } else if (c == '/' - && ((i + 1) >= s.length() || s.charAt(i + 1) == ' ' || s - .charAt(i + 1) == '\t')) { - throw new ParseException("'/' must be followed by an integer.", i); - } else if (c == '*') { - i++; - } - c = s.charAt(i); - if (c == '/') { // is an increment specified? - i++; - if (i >= s.length()) { - throw new ParseException("Unexpected end of string.", i); - } + return getReturn(s, type, incr, i, c); + } - incr = getNumericValue(s, i); - - i++; - if (incr > 10) { - i++; - } - checkIncrementRange(incr, type, i); - } else { - incr = 1; - } - - addToSet(ALL_SPEC_INT, -1, incr, type); - return i; - } else if (c == 'L') { - i++; - if (type == DAY_OF_MONTH) { - lastdayOfMonth = true; - } - if (type == DAY_OF_WEEK) { - addToSet(7, 7, 0, type); - } - if (type == DAY_OF_MONTH && s.length() > i) { - c = s.charAt(i); + private int getReturn(String s, int type, int incr, int i, String sub) throws ParseException { + char c; + int sval = -1; + int eval = -1; + if (type == MONTH) { + sval = getMonthNumber(sub) + 1; + getSval(i, sub, sval <= 0, "Invalid Month value: '"); + if (s.length() > i + 3) { + c = s.charAt(i + 3); if (c == '-') { - ValueSet vs = getValue(0, s, i + 1); - lastdayOffset = vs.value; - if (lastdayOffset > 30) { - throw new ParseException("Offset from last day must be <= 30", i + 1); - } - i = vs.pos; - } - if (s.length() > i) { - c = s.charAt(i); - if (c == 'W') { - nearestWeekday = true; - i++; - } + i += 4; + sub = s.substring(i, i + 3); + eval = getMonthNumber(sub) + 1; + getSval(i, sub, eval <= 0, "Invalid Month value: '"); } } - return i; + } else if (type == DAY_OF_WEEK) { + sval = getDayOfWeekNumber(sub); + getSval(i, sub, sval < 0, "Invalid Day-of-Week value: '"); + if (s.length() > i + 3) { + c = s.charAt(i + 3); + if (c == '-') { + i += 4; + sub = s.substring(i, i + 3); + eval = getDayOfWeekNumber(sub); + } + i = getI(s, i, c, sub, eval); + } + + } else { + throw new ParseException( + "Illegal characters for this position: '" + sub + "'", + i); + } + incr = getIncr(incr, eval != -1, 1); + addToSet(sval, eval, incr, type); + return (i + 3); + } + + private int getI(String s, int i, char c, String sub, int eval) throws ParseException { + if (c == '-') { + getSval(i, sub, eval < 0, "Invalid Day-of-Week value: '"); + } else if (c == '#') { + i = getNthdayOfWeek(s, i); + } else if (c == 'L') { + lastdayOfWeek = true; + i++; + } + return i; + } + + private int getReturn(String s, int type, int incr, int i, char c) throws ParseException { + if (c == '*' || c == '/') { + return getC(s, type, incr, i, c); + } else if (c == 'L') { + return getType1(s, type, i); } else if (c >= '0' && c <= '9') { int val = Integer.parseInt(String.valueOf(c)); i++; @@ -693,6 +619,122 @@ public final class CronExpression implements Serializable, Cloneable { return i; } + private int getType1(String s, int type, int i) throws ParseException { + char c; + i++; + if (type == DAY_OF_MONTH) { + lastdayOfMonth = true; + } + if (type == DAY_OF_WEEK) { + addToSet(7, 7, 0, type); + } + if (type == DAY_OF_MONTH && s.length() > i) { + c = s.charAt(i); + if (c == '-') { + ValueSet vs = getValue(0, s, i + 1); + lastdayOffset = vs.value; + if (lastdayOffset > 30) { + throw new ParseException("Offset from last day must be <= 30", i + 1); + } + i = vs.pos; + } + if (s.length() > i) { + c = s.charAt(i); + if (c == 'W') { + nearestWeekday = true; + i++; + } + } + } + return i; + } + + private int getC(String s, int type, int incr, int i, char c) throws ParseException { + if (c == '*' && (i + 1) >= s.length()) { + addToSet(ALL_SPEC_INT, -1, incr, type); + return i + 1; + } else if (c == '/' + && ((i + 1) >= s.length() || s.charAt(i + 1) == ' ' || s + .charAt(i + 1) == '\t')) { + throw new ParseException("'/' must be followed by an integer.", i); + } else if (c == '*') { + i++; + } + c = s.charAt(i); + if (c == '/') { // is an increment specified? + i++; + if (i >= s.length()) { + throw new ParseException("Unexpected end of string.", i); + } + + incr = getNumericValue(s, i); + + i++; + if (incr > 10) { + i++; + } + checkIncrementRange(incr, type, i); + } else { + incr = 1; + } + + addToSet(ALL_SPEC_INT, -1, incr, type); + return i; + } + + private int getType(String s, int type, int i) throws ParseException { + i++; + if ((i + 1) < s.length() + && (s.charAt(i) != ' ' && s.charAt(i + 1) != '\t')) { + throw new ParseException("Illegal character after '?': " + + s.charAt(i), i); + } + if (type != DAY_OF_WEEK && type != DAY_OF_MONTH) { + throw new ParseException( + "'?' can only be specified for Day-of-Month or Day-of-Week.", + i); + } + if (type == DAY_OF_WEEK && !lastdayOfMonth) { + int val = daysOfMonth.last(); + if (val == NO_SPEC_INT) { + throw new ParseException( + "'?' can only be specified for Day-of-Month -OR- Day-of-Week.", + i); + } + } + + addToSet(NO_SPEC_INT, -1, 0, type); + return i; + } + + private int getIncr(int incr, boolean b, int i2) { + if (b) { + incr = i2; + } + return incr; + } + + private int getNthdayOfWeek(String s, int i) throws ParseException { + try { + i += 4; + nthdayOfWeek = Integer.parseInt(s.substring(i)); + if (nthdayOfWeek < 1 || nthdayOfWeek > 5) { + throw new IllegalArgumentException(); + } + } catch (Exception e) { + throw new ParseException( + "A numeric value between 1 and 5 must follow the '#' option", + i); + } + return i; + } + + private void getSval(int i, String sub, boolean b, String s2) throws ParseException { + if (b) { + throw new ParseException(s2 + sub + "'", i); + } + } + private void checkIncrementRange(int incr, int type, int idxPos) throws ParseException { if (incr > 59 && (type == SECOND || type == MINUTE)) { throw new ParseException("Increment > 60 : " + incr, idxPos); @@ -721,103 +763,23 @@ public final class CronExpression implements Serializable, Cloneable { char c = s.charAt(pos); if (c == 'L') { - if (type == DAY_OF_WEEK) { - if (val < 1 || val > 7) { - throw new ParseException("Day-of-Week values must be between 1 and 7", -1); - } - lastdayOfWeek = true; - } else { - throw new ParseException("'L' option is not valid here. (pos=" + i + ")", i); - } - TreeSet set = getSet(type); - set.add(val); - i++; - return i; + return getType(val, type, i); } if (c == 'W') { - if (type == DAY_OF_MONTH) { - nearestWeekday = true; - } else { - throw new ParseException("'W' option is not valid here. (pos=" + i + ")", i); - } - if (val > 31) { - throw new ParseException("The 'W' option does not make sense with values larger than 31 (max number of days in a month)", i); - } - TreeSet set = getSet(type); - set.add(val); - i++; - return i; + return getType1(val, type, i); } if (c == '#') { - if (type != DAY_OF_WEEK) { - throw new ParseException("'#' option is not valid here. (pos=" + i + ")", i); - } - i++; - try { - nthdayOfWeek = Integer.parseInt(s.substring(i)); - if (nthdayOfWeek < 1 || nthdayOfWeek > 5) { - throw new Exception(); - } - } catch (Exception e) { - throw new ParseException( - "A numeric value between 1 and 5 must follow the '#' option", - i); - } - - TreeSet set = getSet(type); - set.add(val); - i++; - return i; + return getType2(s, val, type, i); } if (c == '-') { - i++; - c = s.charAt(i); - int v = Integer.parseInt(String.valueOf(c)); - end = v; - i++; - if (i >= s.length()) { - addToSet(val, end, 1, type); - return i; - } - c = s.charAt(i); - if (c >= '0' && c <= '9') { - ValueSet vs = getValue(v, s, i); - end = vs.value; - i = vs.pos; - } - if (i < s.length() && ((s.charAt(i)) == '/')) { - i++; - c = s.charAt(i); - int v2 = Integer.parseInt(String.valueOf(c)); - i++; - if (i >= s.length()) { - addToSet(val, end, v2, type); - return i; - } - c = s.charAt(i); - if (c >= '0' && c <= '9') { - ValueSet vs = getValue(v2, s, i); - int v3 = vs.value; - addToSet(val, end, v3, type); - i = vs.pos; - return i; - } else { - addToSet(val, end, v2, type); - return i; - } - } else { - addToSet(val, end, 1, type); - return i; - } + return getC1(s, val, type, i); } if (c == '/') { - if ((i + 1) >= s.length() || s.charAt(i + 1) == ' ' || s.charAt(i + 1) == '\t') { - throw new ParseException("'/' must be followed by an integer.", i); - } + getParseException(s, i); i++; c = s.charAt(i); @@ -846,8 +808,115 @@ public final class CronExpression implements Serializable, Cloneable { return i; } + private int getC1(String s, int val, int type, int i) throws ParseException { + char c; + int end; + i++; + c = s.charAt(i); + int v = Integer.parseInt(String.valueOf(c)); + end = v; + i++; + if (i >= s.length()) { + addToSet(val, end, 1, type); + return i; + } + c = s.charAt(i); + if (c >= '0' && c <= '9') { + ValueSet vs = getValue(v, s, i); + end = vs.value; + i = vs.pos; + } + if (i < s.length() && ((s.charAt(i)) == '/')) { + return getC(s, val, type, end, i); + } else { + addToSet(val, end, 1, type); + return i; + } + } + + private int getC(String s, int val, int type, int end, int i) throws ParseException { + char c; + i++; + c = s.charAt(i); + int v2 = Integer.parseInt(String.valueOf(c)); + i++; + if (i >= s.length()) { + addToSet(val, end, v2, type); + return i; + } + c = s.charAt(i); + if (c >= '0' && c <= '9') { + ValueSet vs = getValue(v2, s, i); + int v3 = vs.value; + addToSet(val, end, v3, type); + i = vs.pos; + return i; + } else { + addToSet(val, end, v2, type); + return i; + } + } + + private void getParseException(String s, int i) throws ParseException { + if ((i + 1) >= s.length() || s.charAt(i + 1) == ' ' || s.charAt(i + 1) == '\t') { + throw new ParseException("'/' must be followed by an integer.", i); + } + } + + private int getType2(String s, int val, int type, int i) throws ParseException { + if (type != DAY_OF_WEEK) { + throw new ParseException("'#' option is not valid here. (pos=" + i + ")", i); + } + i++; + try { + nthdayOfWeek = Integer.parseInt(s.substring(i)); + if (nthdayOfWeek < 1 || nthdayOfWeek > 5) { + throw new IllegalArgumentException(); + } + } catch (Exception e) { + throw new ParseException( + "A numeric value between 1 and 5 must follow the '#' option", + i); + } + + TreeSet set = getSet(type); + set.add(val); + i++; + return i; + } + + private int getType1(int val, int type, int i) throws ParseException { + if (type == DAY_OF_MONTH) { + nearestWeekday = true; + } else { + throw new ParseException("'W' option is not valid here. (pos=" + i + ")", i); + } + if (val > 31) { + throw new ParseException("The 'W' option does not make sense with values larger than 31 (max number of days in a month)", i); + } + TreeSet set = getSet(type); + set.add(val); + i++; + return i; + } + + private int getType(int val, int type, int i) throws ParseException { + if (type == DAY_OF_WEEK) { + if (val < 1 || val > 7) { + throw new ParseException("Day-of-Week values must be between 1 and 7", -1); + } + lastdayOfWeek = true; + } else { + throw new ParseException("'L' option is not valid here. (pos=" + i + ")", i); + } + TreeSet set = getSet(type); + set.add(val); + i++; + return i; + } + public String getCronExpression() { - return cronExpression; + return cronExpressionData; } public String getExpressionSummary() { @@ -943,8 +1012,6 @@ public final class CronExpression implements Serializable, Cloneable { } protected int skipWhiteSpace(int i, String s) { - /*for (; i < s.length() && (s.charAt(i) == ' ' || s.charAt(i) == '\t'); i++) { - }*/ while (i < s.length() && (s.charAt(i) == ' ' || s.charAt(i) == '\t')){ i++; } @@ -952,8 +1019,6 @@ public final class CronExpression implements Serializable, Cloneable { } protected int findNextWhiteSpace(int i, String s) { - /*for (; i < s.length() && (s.charAt(i) != ' ' || s.charAt(i) != '\t'); i++) { - }*/ while (i < s.length() && (s.charAt(i) != ' ' || s.charAt(i) != '\t')){ i++; } @@ -966,43 +1031,16 @@ public final class CronExpression implements Serializable, Cloneable { TreeSet set = getSet(type); - if (type == SECOND || type == MINUTE) { - if ((val < 0 || val > 59 || end > 59) && (val != ALL_SPEC_INT)) { - throw new ParseException( - "Minute and Second values must be between 0 and 59", - -1); - } - } else if (type == HOUR) { - if ((val < 0 || val > 23 || end > 23) && (val != ALL_SPEC_INT)) { - throw new ParseException( - "Hour values must be between 0 and 23", -1); - } - } else if (type == DAY_OF_MONTH) { - if ((val < 1 || val > 31 || end > 31) && (val != ALL_SPEC_INT) - && (val != NO_SPEC_INT)) { - throw new ParseException( - "Day of month values must be between 1 and 31", -1); - } - } else if (type == MONTH) { - if ((val < 1 || val > 12 || end > 12) && (val != ALL_SPEC_INT)) { - throw new ParseException( - "Month values must be between 1 and 12", -1); - } - } else if ((type == DAY_OF_WEEK) && (val == 0 || val > 7 || end > 7) && (val != ALL_SPEC_INT) && (val != NO_SPEC_INT)) { - throw new ParseException( - "Day-of-Week values must be between 1 and 7", -1); - } - - if ((incr == 0 || incr == -1) && val != ALL_SPEC_INT) { - if (val != -1) { - set.add(val); - } else { - set.add(NO_SPEC); - } + getTypeVal(val, end, type); + if (setAdd(val, incr, set)) { return; } + getAt(val, end, incr, type, set); + } + + private void getAt(int val, int end, int incr, int type, TreeSet set) { int startAt = val; int stopAt = end; @@ -1012,47 +1050,23 @@ public final class CronExpression implements Serializable, Cloneable { } if (type == SECOND || type == MINUTE) { - if (stopAt == -1) { - stopAt = 59; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 0; - } + stopAt = getAt(stopAt, stopAt == -1, 59); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 0); } else if (type == HOUR) { - if (stopAt == -1) { - stopAt = 23; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 0; - } + stopAt = getAt(stopAt, stopAt == -1, 23); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 0); } else if (type == DAY_OF_MONTH) { - if (stopAt == -1) { - stopAt = 31; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 1; - } + stopAt = getAt(stopAt, stopAt == -1, 31); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 1); } else if (type == MONTH) { - if (stopAt == -1) { - stopAt = 12; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 1; - } + stopAt = getAt(stopAt, stopAt == -1, 12); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 1); } else if (type == DAY_OF_WEEK) { - if (stopAt == -1) { - stopAt = 7; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 1; - } + stopAt = getAt(stopAt, stopAt == -1, 7); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 1); } else if (type == YEAR) { - if (stopAt == -1) { - stopAt = MAX_YEAR; - } - if (startAt == -1 || startAt == ALL_SPEC_INT) { - startAt = 1970; - } + stopAt = getAt(stopAt, stopAt == -1, MAX_YEAR); + startAt = getAt(startAt, isFlag(startAt, -1, ALL_SPEC_INT), 1970); } // if the end of the range is before the start, then we need to overflow into @@ -1087,6 +1101,21 @@ public final class CronExpression implements Serializable, Cloneable { stopAt += max; } + setAdd1(incr, type, set, startAt, stopAt, max); + } + + private boolean isFlag(int startAt, int i, int allSpecInt) { + return startAt == i || startAt == allSpecInt; + } + + private int getAt(int stopAt, boolean b, int i) { + if (b) { + stopAt = i; + } + return stopAt; + } + + private void setAdd1(int incr, int type, TreeSet set, int startAt, int stopAt, int max) { for (int i = startAt; i <= stopAt; i += incr) { if (max == -1) { // ie: there's no max to overflow over @@ -1105,6 +1134,60 @@ public final class CronExpression implements Serializable, Cloneable { } } + private boolean setAdd(int val, int incr, TreeSet set) { + if ((incr == 0 || incr == -1) && val != ALL_SPEC_INT) { + if (val != -1) { + set.add(val); + } else { + set.add(NO_SPEC); + } + + return true; + } + return false; + } + + private void getTypeVal(int val, int end, int type) throws ParseException { + if (type == SECOND || type == MINUTE) { + if (isBoolean(val, end, 0, 59)) { + throw new ParseException( + "Minute and Second values must be between 0 and 59", + -1); + } + } else if (type == HOUR) { + if (isBoolean(val, end, 0, 23)) { + throw new ParseException( + "Hour values must be between 0 and 23", -1); + } + } else if (type == DAY_OF_MONTH) { + if (isBoolean1(val, end)) { + throw new ParseException( + "Day of month values must be between 1 and 31", -1); + } + } else if (type == MONTH) { + if (isBoolean(val, end, 1, 12)) { + throw new ParseException( + "Month values must be between 1 and 12", -1); + } + } else if (isBoolean2(val, end, type)) { + throw new ParseException( + "Day-of-Week values must be between 1 and 7", -1); + } + } + + private boolean isBoolean2(int val, int end, int type) { + return (type == DAY_OF_WEEK) && (val == 0 || val > 7 || end > 7) && (val != ALL_SPEC_INT) && (val != NO_SPEC_INT); + } + + private boolean isBoolean1(int val, int end) { + return (val < 1 || val > 31 || end > 31) && (val != ALL_SPEC_INT) + && (val != NO_SPEC_INT); + } + + private boolean isBoolean(int val, int end, int i, int i2) { + return (val < i || val > i2 || end > i2) && (val != ALL_SPEC_INT); + } + TreeSet getSet(int type) { switch (type) { case SECOND: @@ -1192,7 +1275,6 @@ public final class CronExpression implements Serializable, Cloneable { // loop until we've computed the next time, or we've past the endTime while (!gotOne) { - //if (endTime != null && cl.getTime().after(endTime)) return null; if (cl.get(Calendar.YEAR) > 2999) { // prevent endless loop... return null; } @@ -1204,15 +1286,7 @@ public final class CronExpression implements Serializable, Cloneable { int min = cl.get(Calendar.MINUTE); // get second................................................. - st = seconds.tailSet(sec); - if (st != null && !st.isEmpty()) { - sec = st.first(); - } else { - sec = seconds.first(); - min++; - cl.set(Calendar.MINUTE, min); - } - cl.set(Calendar.SECOND, sec); + sec = getSec(cl, sec, min); min = cl.get(Calendar.MINUTE); int hr = cl.get(Calendar.HOUR_OF_DAY); @@ -1299,15 +1373,7 @@ public final class CronExpression implements Serializable, Cloneable { int ldom = getLastDayOfMonth(mon, cl.get(Calendar.YEAR)); int dow = tcal.get(Calendar.DAY_OF_WEEK); - if (dow == Calendar.SATURDAY && day == 1) { - day += 2; - } else if (dow == Calendar.SATURDAY) { - day -= 1; - } else if (dow == Calendar.SUNDAY && day == ldom) { - day -= 2; - } else if (dow == Calendar.SUNDAY) { - day += 1; - } + day = getDay(day, ldom, dow); tcal.set(Calendar.SECOND, sec); tcal.set(Calendar.MINUTE, min); @@ -1335,15 +1401,7 @@ public final class CronExpression implements Serializable, Cloneable { int ldom = getLastDayOfMonth(mon, cl.get(Calendar.YEAR)); int dow = tcal.get(Calendar.DAY_OF_WEEK); - if (dow == Calendar.SATURDAY && day == 1) { - day += 2; - } else if (dow == Calendar.SATURDAY) { - day -= 1; - } else if (dow == Calendar.SUNDAY && day == ldom) { - day -= 2; - } else if (dow == Calendar.SUNDAY) { - day += 1; - } + day = getDay(day, ldom, dow); tcal.set(Calendar.SECOND, sec); @@ -1370,29 +1428,14 @@ public final class CronExpression implements Serializable, Cloneable { mon++; } - if (day != t || mon != tmon) { - cl.set(Calendar.SECOND, 0); - cl.set(Calendar.MINUTE, 0); - cl.set(Calendar.HOUR_OF_DAY, 0); - cl.set(Calendar.DAY_OF_MONTH, day); - cl.set(Calendar.MONTH, mon - 1); - // '- 1' because calendar is 0-based for this field, and we - // are 1-based - continue; - } + if (getContinue2(cl, t, day, mon, tmon)) continue; } else if (dayOfWSpec && !dayOfMSpec) { // get day by day of week rule if (lastdayOfWeek) { // are we looking for the last XXX day of // the month? int dow = daysOfWeek.first(); // desired // d-o-w int cDow = cl.get(Calendar.DAY_OF_WEEK); // current d-o-w - int daysToAdd = 0; - if (cDow < dow) { - daysToAdd = dow - cDow; - } - if (cDow > dow) { - daysToAdd = dow + (7 - cDow); - } + int daysToAdd = getDaysToAdd(cDow, dow); int lDay = getLastDayOfMonth(mon, cl.get(Calendar.YEAR)); @@ -1408,9 +1451,7 @@ public final class CronExpression implements Serializable, Cloneable { } // find date of last occurrence of this day in this month... - while ((day + daysToAdd + 7) <= lDay) { - daysToAdd += 7; - } + daysToAdd = getDaysToAdd(day, daysToAdd, lDay); day += daysToAdd; @@ -1429,43 +1470,16 @@ public final class CronExpression implements Serializable, Cloneable { int dow = daysOfWeek.first(); // desired // d-o-w int cDow = cl.get(Calendar.DAY_OF_WEEK); // current d-o-w - int daysToAdd = 0; - if (cDow < dow) { - daysToAdd = dow - cDow; - } else if (cDow > dow) { - daysToAdd = dow + (7 - cDow); - } + int daysToAdd = getDaysToAdd1(dow, cDow); - boolean dayShifted = false; - if (daysToAdd > 0) { - dayShifted = true; - } + boolean dayShifted = isDayShifted(daysToAdd); day += daysToAdd; - int weekOfMonth = day / 7; - if (day % 7 > 0) { - weekOfMonth++; - } + int weekOfMonth = getWeekOfMonth(day); daysToAdd = (nthdayOfWeek - weekOfMonth) * 7; day += daysToAdd; - if (daysToAdd < 0 - || day > getLastDayOfMonth(mon, cl - .get(Calendar.YEAR))) { - cl.set(Calendar.SECOND, 0); - cl.set(Calendar.MINUTE, 0); - cl.set(Calendar.HOUR_OF_DAY, 0); - cl.set(Calendar.DAY_OF_MONTH, 1); - cl.set(Calendar.MONTH, mon); - // no '- 1' here because we are promoting the month - continue; - } else if (daysToAdd > 0 || dayShifted) { - cl.set(Calendar.SECOND, 0); - cl.set(Calendar.MINUTE, 0); - cl.set(Calendar.HOUR_OF_DAY, 0); - cl.set(Calendar.DAY_OF_MONTH, day); - cl.set(Calendar.MONTH, mon - 1); - // '- 1' here because we are NOT promoting the month + if (getContinue(cl, day, mon, daysToAdd, dayShifted)) { continue; } } else { @@ -1473,37 +1487,13 @@ public final class CronExpression implements Serializable, Cloneable { int dow = daysOfWeek.first(); // desired // d-o-w st = daysOfWeek.tailSet(cDow); - if (st != null && !st.isEmpty()) { - dow = st.first(); - } + dow = getDow(dow, st != null && !st.isEmpty(), st.first()); - int daysToAdd = 0; - if (cDow < dow) { - daysToAdd = dow - cDow; - } - if (cDow > dow) { - daysToAdd = dow + (7 - cDow); - } + int daysToAdd = getDaysToAdd(cDow, dow); int lDay = getLastDayOfMonth(mon, cl.get(Calendar.YEAR)); - if (day + daysToAdd > lDay) { // will we pass the end of - // the month? - cl.set(Calendar.SECOND, 0); - cl.set(Calendar.MINUTE, 0); - cl.set(Calendar.HOUR_OF_DAY, 0); - cl.set(Calendar.DAY_OF_MONTH, 1); - cl.set(Calendar.MONTH, mon); - // no '- 1' here because we are promoting the month - continue; - } else if (daysToAdd > 0) { // are we swithing days? - cl.set(Calendar.SECOND, 0); - cl.set(Calendar.MINUTE, 0); - cl.set(Calendar.HOUR_OF_DAY, 0); - cl.set(Calendar.DAY_OF_MONTH, day + daysToAdd); - cl.set(Calendar.MONTH, mon - 1); - // '- 1' because calendar is 0-based for this field, - // and we are 1-based + if (getContinue(cl, day, mon, daysToAdd, lDay)) { continue; } } @@ -1550,7 +1540,6 @@ public final class CronExpression implements Serializable, Cloneable { // 1-based year = cl.get(Calendar.YEAR); - t = -1; // get year................................................... st = years.tailSet(year); @@ -1575,11 +1564,149 @@ public final class CronExpression implements Serializable, Cloneable { cl.set(Calendar.YEAR, year); gotOne = true; - } // while( !done ) + } return cl.getTime(); } + private boolean getContinue2(Calendar cl, int t, int day, int mon, int tmon) { + if (day != t || mon != tmon) { + cl.set(Calendar.SECOND, 0); + cl.set(Calendar.MINUTE, 0); + cl.set(Calendar.HOUR_OF_DAY, 0); + cl.set(Calendar.DAY_OF_MONTH, day); + cl.set(Calendar.MONTH, mon - 1); + // '- 1' because calendar is 0-based for this field, and we + // are 1-based + return true; + } + return false; + } + + private int getDow(int dow, boolean b, Integer first) { + if (b) { + dow = first; + } + return dow; + } + + private boolean getContinue(Calendar cl, int day, int mon, int daysToAdd, boolean dayShifted) { + if (daysToAdd < 0 + || day > getLastDayOfMonth(mon, cl + .get(Calendar.YEAR))) { + cl.set(Calendar.SECOND, 0); + cl.set(Calendar.MINUTE, 0); + cl.set(Calendar.HOUR_OF_DAY, 0); + cl.set(Calendar.DAY_OF_MONTH, 1); + cl.set(Calendar.MONTH, mon); + // no '- 1' here because we are promoting the month + return true; + } else if (daysToAdd > 0 || dayShifted) { + cl.set(Calendar.SECOND, 0); + cl.set(Calendar.MINUTE, 0); + cl.set(Calendar.HOUR_OF_DAY, 0); + cl.set(Calendar.DAY_OF_MONTH, day); + cl.set(Calendar.MONTH, mon - 1); + // '- 1' here because we are NOT promoting the month + return true; + } + return false; + } + + private int getWeekOfMonth(int day) { + int weekOfMonth = day / 7; + if (day % 7 > 0) { + weekOfMonth++; + } + return weekOfMonth; + } + + private boolean isDayShifted(int daysToAdd) { + boolean dayShifted = false; + if (daysToAdd > 0) { + dayShifted = true; + } + return dayShifted; + } + + private int getDaysToAdd1(int dow, int cDow) { + int daysToAdd = 0; + if (cDow < dow) { + daysToAdd = dow - cDow; + } else if (cDow > dow) { + daysToAdd = dow + (7 - cDow); + } + return daysToAdd; + } + + private int getDaysToAdd(int day, int daysToAdd, int lDay) { + while ((day + daysToAdd + 7) <= lDay) { + daysToAdd += 7; + } + return daysToAdd; + } + + private int getDay(int day, int ldom, int dow) { + if (dow == Calendar.SATURDAY && day == 1) { + day += 2; + } else if (dow == Calendar.SATURDAY) { + day -= 1; + } else if (dow == Calendar.SUNDAY && day == ldom) { + day -= 2; + } else if (dow == Calendar.SUNDAY) { + day += 1; + } + return day; + } + + private int getDaysToAdd(int cDow, int dow) { + int daysToAdd = 0; + if (cDow < dow) { + daysToAdd = dow - cDow; + } + if (cDow > dow) { + daysToAdd = dow + (7 - cDow); + } + return daysToAdd; + } + + private boolean getContinue(Calendar cl, int day, int mon, int daysToAdd, int lDay) { + if (day + daysToAdd > lDay) { // will we pass the end of + // the month? + cl.set(Calendar.SECOND, 0); + cl.set(Calendar.MINUTE, 0); + cl.set(Calendar.HOUR_OF_DAY, 0); + cl.set(Calendar.DAY_OF_MONTH, 1); + cl.set(Calendar.MONTH, mon); + // no '- 1' here because we are promoting the month + return true; + } else if (daysToAdd > 0) { // are we swithing days? + cl.set(Calendar.SECOND, 0); + cl.set(Calendar.MINUTE, 0); + cl.set(Calendar.HOUR_OF_DAY, 0); + cl.set(Calendar.DAY_OF_MONTH, day + daysToAdd); + cl.set(Calendar.MONTH, mon - 1); + // '- 1' because calendar is 0-based for this field, + // and we are 1-based + return true; + } + return false; + } + + private int getSec(Calendar cl, int sec, int min) { + SortedSet st; + st = seconds.tailSet(sec); + if (st != null && !st.isEmpty()) { + sec = st.first(); + } else { + sec = seconds.first(); + min++; + cl.set(Calendar.MINUTE, min); + } + cl.set(Calendar.SECOND, sec); + return sec; + } + /** * Advance the calendar to the particular hour paying particular attention * to daylight saving problems. @@ -1599,6 +1726,9 @@ public final class CronExpression implements Serializable, Cloneable { * that the CronExpression matches. */ public Date getTimeBefore(Date endTime) { + if(Objects.isNull(endTime)){ + log.info("未使用参数"); + } // FUTURE_TODO: implement QUARTZ-423 return null; } @@ -1655,9 +1785,10 @@ public final class CronExpression implements Serializable, Cloneable { stream.defaultReadObject(); try { - buildExpression(cronExpression); + buildExpression(cronExpressionData); } catch (Exception ignore) { - } // never happens + log.error("异常:", ignore); + } } /** diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/RemoteHttpJobBean.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/RemoteHttpJobBean.java index b2dd1515..e69de29b 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/RemoteHttpJobBean.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/RemoteHttpJobBean.java @@ -1,32 +0,0 @@ -//package com.xxl.job.admin.core.jobbean; -// -//import com.xxl.job.admin.core.thread.JobTriggerPoolHelper; -//import com.xxl.job.admin.core.trigger.TriggerTypeEnum; -//import org.quartz.JobExecutionContext; -//import org.quartz.JobExecutionException; -//import org.quartz.JobKey; -//import org.slf4j.Logger; -//import org.slf4j.LoggerFactory; -//import org.springframework.scheduling.quartz.QuartzJobBean; -// -///** -// * http job bean -// * “@DisallowConcurrentExecution” disable concurrent, thread size can not be only one, better given more -// * @author xuxueli 2015-12-17 18:20:34 -// */ -////@DisallowConcurrentExecution -//public class RemoteHttpJobBean extends QuartzJobBean { -// private static Logger logger = LoggerFactory.getLogger(RemoteHttpJobBean.class); -// -// @Override -// protected void executeInternal(JobExecutionContext context) -// throws JobExecutionException { -// -// // load jobId -// JobKey jobKey = context.getTrigger().getJobKey(); -// Integer jobId = Integer.valueOf(jobKey.getName()); -// -// -// } -// -//} \ No newline at end of file diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobDynamicScheduler.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobDynamicScheduler.java index 1e62aa19..e69de29b 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobDynamicScheduler.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobDynamicScheduler.java @@ -1,413 +0,0 @@ -//package com.xxl.job.admin.core.schedule; -// -//import com.xxl.job.admin.core.conf.XxlJobAdminConfig; -//import com.xxl.job.admin.core.jobbean.RemoteHttpJobBean; -//import com.xxl.job.admin.core.model.XxlJobInfo; -//import com.xxl.job.admin.core.thread.JobFailMonitorHelper; -//import com.xxl.job.admin.core.thread.JobRegistryMonitorHelper; -//import com.xxl.job.admin.core.thread.JobTriggerPoolHelper; -//import com.xxl.job.admin.core.util.I18nUtil; -//import com.xxl.job.core.biz.AdminBiz; -//import com.xxl.job.core.biz.ExecutorBiz; -//import com.xxl.job.core.enums.ExecutorBlockStrategyEnum; -//import com.xxl.rpc.remoting.invoker.XxlRpcInvokerFactory; -//import com.xxl.rpc.remoting.invoker.call.CallType; -//import com.xxl.rpc.remoting.invoker.reference.XxlRpcReferenceBean; -//import com.xxl.rpc.remoting.invoker.route.LoadBalance; -//import com.xxl.rpc.remoting.net.NetEnum; -//import com.xxl.rpc.remoting.net.impl.servlet.server.ServletServerHandler; -//import com.xxl.rpc.remoting.provider.XxlRpcProviderFactory; -//import com.xxl.rpc.serialize.Serializer; -//import org.quartz.*; -//import org.quartz.Trigger.TriggerState; -//import org.quartz.impl.triggers.CronTriggerImpl; -//import org.slf4j.Logger; -//import org.slf4j.LoggerFactory; -//import org.springframework.util.Assert; -// -//import javax.servlet.ServletException; -//import javax.servlet.http.HttpServletRequest; -//import javax.servlet.http.HttpServletResponse; -//import java.io.IOException; -//import java.util.Date; -//import java.util.concurrent.ConcurrentHashMap; -// -///** -// * base quartz scheduler util -// * @author xuxueli 2015-12-19 16:13:53 -// */ -//public final class XxlJobDynamicScheduler { -// private static final Logger logger = LoggerFactory.getLogger(XxlJobDynamicScheduler_old.class); -// -// // ---------------------- param ---------------------- -// -// // scheduler -// private static Scheduler scheduler; -// public void setScheduler(Scheduler scheduler) { -// XxlJobDynamicScheduler_old.scheduler = scheduler; -// } -// -// -// // ---------------------- init + destroy ---------------------- -// public void start() throws Exception { -// // valid -// Assert.notNull(scheduler, "quartz scheduler is null"); -// -// // init i18n -// initI18n(); -// -// // admin registry monitor run -// JobRegistryMonitorHelper.getInstance().start(); -// -// // admin monitor run -// JobFailMonitorHelper.getInstance().start(); -// -// // admin-server -// initRpcProvider(); -// -// logger.info(">>>>>>>>> init xxl-job admin success."); -// } -// -// -// public void destroy() throws Exception { -// // admin trigger pool stop -// JobTriggerPoolHelper.toStop(); -// -// // admin registry stop -// JobRegistryMonitorHelper.getInstance().toStop(); -// -// // admin monitor stop -// JobFailMonitorHelper.getInstance().toStop(); -// -// // admin-server -// stopRpcProvider(); -// } -// -// -// // ---------------------- I18n ---------------------- -// -// private void initI18n(){ -// for (ExecutorBlockStrategyEnum item:ExecutorBlockStrategyEnum.values()) { -// item.setTitle(I18nUtil.getString("jobconf_block_".concat(item.name()))); -// } -// } -// -// -// // ---------------------- admin rpc provider (no server version) ---------------------- -// private static ServletServerHandler servletServerHandler; -// private void initRpcProvider(){ -// // init -// XxlRpcProviderFactory xxlRpcProviderFactory = new XxlRpcProviderFactory(); -// xxlRpcProviderFactory.initConfig( -// NetEnum.NETTY_HTTP, -// Serializer.SerializeEnum.HESSIAN.getSerializer(), -// null, -// 0, -// XxlJobAdminConfig.getAdminConfig().getAccessToken(), -// null, -// null); -// -// // add services -// xxlRpcProviderFactory.addService(AdminBiz.class.getName(), null, XxlJobAdminConfig.getAdminConfig().getAdminBiz()); -// -// // servlet handler -// servletServerHandler = new ServletServerHandler(xxlRpcProviderFactory); -// } -// private void stopRpcProvider() throws Exception { -// XxlRpcInvokerFactory.getInstance().stop(); -// } -// public static void invokeAdminService(HttpServletRequest request, HttpServletResponse response) throws IOException, ServletException { -// servletServerHandler.handle(null, request, response); -// } -// -// -// // ---------------------- executor-client ---------------------- -// private static ConcurrentHashMap executorBizRepository = new ConcurrentHashMap(); -// public static ExecutorBiz getExecutorBiz(String address) throws Exception { -// // valid -// if (address==null || address.trim().length()==0) { -// return null; -// } -// -// // load-cache -// address = address.trim(); -// ExecutorBiz executorBiz = executorBizRepository.get(address); -// if (executorBiz != null) { -// return executorBiz; -// } -// -// // set-cache -// executorBiz = (ExecutorBiz) new XxlRpcReferenceBean( -// NetEnum.NETTY_HTTP, -// Serializer.SerializeEnum.HESSIAN.getSerializer(), -// CallType.SYNC, -// LoadBalance.ROUND, -// ExecutorBiz.class, -// null, -// 5000, -// address, -// XxlJobAdminConfig.getAdminConfig().getAccessToken(), -// null, -// null).getObject(); -// -// executorBizRepository.put(address, executorBiz); -// return executorBiz; -// } -// -// -// // ---------------------- schedule util ---------------------- -// -// /** -// * fill job info -// * -// * @param jobInfo -// */ -// public static void fillJobInfo(XxlJobInfo jobInfo) { -// -// String name = String.valueOf(jobInfo.getId()); -// -// // trigger key -// TriggerKey triggerKey = TriggerKey.triggerKey(name); -// try { -// -// // trigger cron -// Trigger trigger = scheduler.getTrigger(triggerKey); -// if (trigger!=null && trigger instanceof CronTriggerImpl) { -// String cronExpression = ((CronTriggerImpl) trigger).getCronExpression(); -// jobInfo.setJobCron(cronExpression); -// } -// -// // trigger state -// TriggerState triggerState = scheduler.getTriggerState(triggerKey); -// if (triggerState!=null) { -// jobInfo.setJobStatus(triggerState.name()); -// } -// -// //JobKey jobKey = new JobKey(jobInfo.getJobName(), String.valueOf(jobInfo.getJobGroup())); -// //JobDetail jobDetail = scheduler.getJobDetail(jobKey); -// //String jobClass = jobDetail.getJobClass().getName(); -// -// } catch (SchedulerException e) { -// logger.error(e.getMessage(), e); -// } -// } -// -// -// /** -// * add trigger + job -// * -// * @param jobName -// * @param cronExpression -// * @return -// * @throws SchedulerException -// */ -// public static boolean addJob(String jobName, String cronExpression) throws SchedulerException { -// // 1、job key -// TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// JobKey jobKey = new JobKey(jobName); -// -// // 2、valid -// if (scheduler.checkExists(triggerKey)) { -// return true; // PASS -// } -// -// // 3、corn trigger -// CronScheduleBuilder cronScheduleBuilder = CronScheduleBuilder.cronSchedule(cronExpression).withMisfireHandlingInstructionDoNothing(); // withMisfireHandlingInstructionDoNothing 忽略掉调度终止过程中忽略的调度 -// CronTrigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).withSchedule(cronScheduleBuilder).build(); -// -// // 4、job detail -// Class jobClass_ = RemoteHttpJobBean.class; // Class.forName(jobInfo.getJobClass()); -// JobDetail jobDetail = JobBuilder.newJob(jobClass_).withIdentity(jobKey).build(); -// -// /*if (jobInfo.getJobData()!=null) { -// JobDataMap jobDataMap = jobDetail.getJobDataMap(); -// jobDataMap.putAll(JacksonUtil.readValue(jobInfo.getJobData(), Map.class)); -// // JobExecutionContext context.getMergedJobDataMap().get("mailGuid"); -// }*/ -// -// // 5、schedule job -// Date date = scheduler.scheduleJob(jobDetail, cronTrigger); -// -// logger.info(">>>>>>>>>>> addJob success(quartz), jobDetail:{}, cronTrigger:{}, date:{}", jobDetail, cronTrigger, date); -// return true; -// } -// -// -// /** -// * remove trigger + job -// * -// * @param jobName -// * @return -// * @throws SchedulerException -// */ -// public static boolean removeJob(String jobName) throws SchedulerException { -// -// JobKey jobKey = new JobKey(jobName); -// scheduler.deleteJob(jobKey); -// -// /*TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// if (scheduler.checkExists(triggerKey)) { -// scheduler.unscheduleJob(triggerKey); // trigger + job -// }*/ -// -// logger.info(">>>>>>>>>>> removeJob success(quartz), jobKey:{}", jobKey); -// return true; -// } -// -// -// /** -// * updateJobCron -// * -// * @param jobName -// * @param cronExpression -// * @return -// * @throws SchedulerException -// */ -// public static boolean updateJobCron(String jobName, String cronExpression) throws SchedulerException { -// -// // 1、job key -// TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// -// // 2、valid -// if (!scheduler.checkExists(triggerKey)) { -// return true; // PASS -// } -// -// CronTrigger oldTrigger = (CronTrigger) scheduler.getTrigger(triggerKey); -// -// // 3、avoid repeat cron -// String oldCron = oldTrigger.getCronExpression(); -// if (oldCron.equals(cronExpression)){ -// return true; // PASS -// } -// -// // 4、new cron trigger -// CronScheduleBuilder cronScheduleBuilder = CronScheduleBuilder.cronSchedule(cronExpression).withMisfireHandlingInstructionDoNothing(); -// oldTrigger = oldTrigger.getTriggerBuilder().withIdentity(triggerKey).withSchedule(cronScheduleBuilder).build(); -// -// // 5、rescheduleJob -// scheduler.rescheduleJob(triggerKey, oldTrigger); -// -// /* -// JobKey jobKey = new JobKey(jobName); -// -// // old job detail -// JobDetail jobDetail = scheduler.getJobDetail(jobKey); -// -// // new trigger -// HashSet triggerSet = new HashSet(); -// triggerSet.add(cronTrigger); -// // cover trigger of job detail -// scheduler.scheduleJob(jobDetail, triggerSet, true);*/ -// -// logger.info(">>>>>>>>>>> resumeJob success, JobName:{}", jobName); -// return true; -// } -// -// -// /** -// * pause -// * -// * @param jobName -// * @return -// * @throws SchedulerException -// */ -// /*public static boolean pauseJob(String jobName) throws SchedulerException { -// -// TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// -// boolean result = false; -// if (scheduler.checkExists(triggerKey)) { -// scheduler.pauseTrigger(triggerKey); -// result = true; -// } -// -// logger.info(">>>>>>>>>>> pauseJob {}, triggerKey:{}", (result?"success":"fail"),triggerKey); -// return result; -// }*/ -// -// -// /** -// * resume -// * -// * @param jobName -// * @return -// * @throws SchedulerException -// */ -// /*public static boolean resumeJob(String jobName) throws SchedulerException { -// -// TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// -// boolean result = false; -// if (scheduler.checkExists(triggerKey)) { -// scheduler.resumeTrigger(triggerKey); -// result = true; -// } -// -// logger.info(">>>>>>>>>>> resumeJob {}, triggerKey:{}", (result?"success":"fail"), triggerKey); -// return result; -// }*/ -// -// -// /** -// * run -// * -// * @param jobName -// * @return -// * @throws SchedulerException -// */ -// /*public static boolean triggerJob(String jobName) throws SchedulerException { -// // TriggerKey : name + group -// JobKey jobKey = new JobKey(jobName); -// TriggerKey triggerKey = TriggerKey.triggerKey(jobName); -// -// boolean result = false; -// if (scheduler.checkExists(triggerKey)) { -// scheduler.triggerJob(jobKey); -// result = true; -// logger.info(">>>>>>>>>>> runJob success, jobKey:{}", jobKey); -// } else { -// logger.info(">>>>>>>>>>> runJob fail, jobKey:{}", jobKey); -// } -// return result; -// }*/ -// -// -// /** -// * finaAllJobList -// * -// * @return -// *//* -// @Deprecated -// public static List> finaAllJobList(){ -// List> jobList = new ArrayList>(); -// -// try { -// if (scheduler.getJobGroupNames()==null || scheduler.getJobGroupNames().size()==0) { -// return null; -// } -// String groupName = scheduler.getJobGroupNames().get(0); -// Set jobKeys = scheduler.getJobKeys(GroupMatcher.jobGroupEquals(groupName)); -// if (jobKeys!=null && jobKeys.size()>0) { -// for (JobKey jobKey : jobKeys) { -// TriggerKey triggerKey = TriggerKey.triggerKey(jobKey.getName(), Scheduler.DEFAULT_GROUP); -// Trigger trigger = scheduler.getTrigger(triggerKey); -// JobDetail jobDetail = scheduler.getJobDetail(jobKey); -// TriggerState triggerState = scheduler.getTriggerState(triggerKey); -// Map jobMap = new HashMap(); -// jobMap.put("TriggerKey", triggerKey); -// jobMap.put("Trigger", trigger); -// jobMap.put("JobDetail", jobDetail); -// jobMap.put("TriggerState", triggerState); -// jobList.add(jobMap); -// } -// } -// -// } catch (SchedulerException e) { -// logger.error(e.getMessage(), e); -// return null; -// } -// return jobList; -// }*/ -// -//} \ No newline at end of file diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobThreadPool.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobThreadPool.java index ad074307..e69de29b 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobThreadPool.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/old/XxlJobThreadPool.java @@ -1,58 +0,0 @@ -//package com.xxl.job.admin.core.quartz; -// -//import org.quartz.SchedulerConfigException; -//import org.quartz.spi.ThreadPool; -// -///** -// * single thread pool, for async trigger -// * -// * @author xuxueli 2019-03-06 -// */ -//public class XxlJobThreadPool implements ThreadPool { -// -// @Override -// public boolean runInThread(Runnable runnable) { -// -// // async run -// runnable.run(); -// return true; -// -// //return false; -// } -// -// @Override -// public int blockForAvailableThreads() { -// return 1; -// } -// -// @Override -// public void initialize() throws SchedulerConfigException { -// -// } -// -// @Override -// public void shutdown(boolean waitForJobsToComplete) { -// -// } -// -// @Override -// public int getPoolSize() { -// return 1; -// } -// -// @Override -// public void setInstanceId(String schedInstId) { -// -// } -// -// @Override -// public void setInstanceName(String schedName) { -// -// } -// -// // support -// public void setThreadCount(int count) { -// // -// } -// -//} diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteConsistentHash.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteConsistentHash.java index df5c8d9b..227f83cf 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteConsistentHash.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteConsistentHash.java @@ -33,7 +33,7 @@ public class ExecutorRouteConsistentHash extends ExecutorRouter { try { md5 = MessageDigest.getInstance("MD5"); } catch (NoSuchAlgorithmException e) { - throw new RuntimeException("MD5 not supported", e); + throw new IllegalArgumentException("MD5 not supported", e); } md5.reset(); byte[] keyBytes = null; diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLFU.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLFU.java index ffdc3c55..64f5088d 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLFU.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLFU.java @@ -3,7 +3,10 @@ package com.xxl.job.admin.core.route.strategy; import com.xxl.job.admin.core.route.ExecutorRouter; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.biz.model.TriggerParam; +import lombok.extern.slf4j.Slf4j; +import java.security.NoSuchAlgorithmException; +import java.security.SecureRandom; import java.util.*; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -15,11 +18,21 @@ import java.util.concurrent.ConcurrentMap; * * Created by xuxueli on 17/3/10. */ +@Slf4j public class ExecutorRouteLFU extends ExecutorRouter { private static ConcurrentMap> jobLfuMap = new ConcurrentHashMap<>(); private long cacheValidTime = 0; + private static Random rand; + static { + try { + rand = SecureRandom.getInstanceStrong(); + } catch (NoSuchAlgorithmException e) { + log.info("Exception:{}", e.getMessage()); + } + } + public String route(int jobId, List addressList) { // cache clear @@ -38,7 +51,8 @@ public class ExecutorRouteLFU extends ExecutorRouter { // put new for (String address: addressList) { if (!lfuItemMap.containsKey(address) || lfuItemMap.get(address) >1000000 ) { - lfuItemMap.put(address, new Random().nextInt(addressList.size())); // 初始化时主动Random一次,缓解首次压力 + // 初始化时主动Random一次,缓解首次压力 + lfuItemMap.put(address, rand.nextInt(addressList.size())); } } // remove old @@ -56,16 +70,9 @@ public class ExecutorRouteLFU extends ExecutorRouter { // load least userd count address List> lfuItemList = new ArrayList<>(lfuItemMap.entrySet()); - /* Collections.sort(lfuItemList, new Comparator>() { - @Override - public int compare(Map.Entry o1, Map.Entry o2) { - return o1.getValue().compareTo(o2.getValue()); - } - });*/ - Collections.sort(lfuItemList, (o1,o2)-> o1.getValue().compareTo(o2.getValue())); + Collections.sort(lfuItemList, Comparator.comparing(Map.Entry::getValue)); Map.Entry addressItem = lfuItemList.get(0); -// String minAddress = addressItem.getKey(); addressItem.setValue(addressItem.getValue() + 1); return addressItem.getKey(); diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLRU.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLRU.java index a01f5c8e..19874da3 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLRU.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteLRU.java @@ -44,9 +44,6 @@ public class ExecutorRouteLRU extends ExecutorRouter { // put new for (String address: addressList) { - /*if (!lruItem.containsKey(address)) { - lruItem.put(address, address); - }*/ lruItem.computeIfAbsent(address,key->address); } // remove old diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteRound.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteRound.java index d43183df..6533925e 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteRound.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/route/strategy/ExecutorRouteRound.java @@ -3,7 +3,10 @@ package com.xxl.job.admin.core.route.strategy; import com.xxl.job.admin.core.route.ExecutorRouter; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.biz.model.TriggerParam; +import lombok.extern.slf4j.Slf4j; +import java.security.NoSuchAlgorithmException; +import java.security.SecureRandom; import java.util.List; import java.util.Random; import java.util.concurrent.ConcurrentHashMap; @@ -12,10 +15,21 @@ import java.util.concurrent.ConcurrentMap; /** * Created by xuxueli on 17/3/10. */ +@Slf4j public class ExecutorRouteRound extends ExecutorRouter { private static ConcurrentMap routeCountEachJob = new ConcurrentHashMap<>(); private static long cacheValidTime = 0; + + private static Random rand; + static { + try { + rand = SecureRandom.getInstanceStrong(); + } catch (NoSuchAlgorithmException e) { + log.info("Exception:{}", e.getMessage()); + } + } + private static int count(int jobId) { // cache clear if (System.currentTimeMillis() > cacheValidTime) { @@ -25,7 +39,7 @@ public class ExecutorRouteRound extends ExecutorRouter { // count++ Integer count = routeCountEachJob.get(jobId); - count = (count==null || count>1000000)?(new Random().nextInt(100)):++count; // 初始化时主动Random一次,缓解首次压力 + count = (count==null || count>1000000)?(rand.nextInt(100)):++count; // 初始化时主动Random一次,缓解首次压力 routeCountEachJob.put(jobId, count); return count; } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobFailMonitorHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobFailMonitorHelper.java index 3598a46f..7a220c56 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobFailMonitorHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobFailMonitorHelper.java @@ -5,6 +5,7 @@ import com.xxl.job.admin.core.model.XxlJobInfo; import com.xxl.job.admin.core.model.XxlJobLog; import com.xxl.job.admin.core.trigger.TriggerTypeEnum; import com.xxl.job.admin.core.util.I18nUtil; +import lombok.extern.slf4j.Slf4j; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -16,93 +17,102 @@ import java.util.concurrent.TimeUnit; * * @author xuxueli 2015-9-1 18:05:56 */ +@Slf4j public class JobFailMonitorHelper { - private static Logger logger = LoggerFactory.getLogger(JobFailMonitorHelper.class); - private static JobFailMonitorHelper instance = new JobFailMonitorHelper(); - public static JobFailMonitorHelper getInstance(){ - return instance; - } + private static JobFailMonitorHelper instance = new JobFailMonitorHelper(); - // ---------------------- monitor ---------------------- + public static JobFailMonitorHelper getInstance() { + return instance; + } - private Thread monitorThread; - private volatile boolean toStop = false; - public void start(){ - monitorThread = new Thread(()-> { + // ---------------------- monitor ---------------------- + + private Thread monitorThread; + private volatile boolean toStop = false; + + public void start() { + monitorThread = new Thread(() -> { + // monitor + getWhile(); + + log.info(">>>>>>>>>>> xxl-job, job fail monitor thread stop"); - // monitor - while (!toStop) { - try { - - List failLogIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findFailJobLogIds(1000); - if (failLogIds!=null && !failLogIds.isEmpty()) { - for (long failLogId: failLogIds) { - - // lock log - int lockRet = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, 0, -1); - if (lockRet < 1) { - continue; - } - XxlJobLog log = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().load(failLogId); - XxlJobInfo info = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().loadById(log.getJobId()); - - // 1、fail retry monitor - if (log.getExecutorFailRetryCount() > 0) { - JobTriggerPoolHelper.trigger(log.getJobId(), TriggerTypeEnum.RETRY, (log.getExecutorFailRetryCount()-1), log.getExecutorShardingParam(), log.getExecutorParam(), null); - String retryMsg = "

>>>>>>>>>>>"+ I18nUtil.getString("jobconf_trigger_type_retry") +"<<<<<<<<<<<
"; - log.setTriggerMsg(log.getTriggerMsg() + retryMsg); - XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(log); - } - - // 2、fail alarm monitor - int newAlarmStatus = 0; // 告警状态:0-默认、-1=锁定状态、1-无需告警、2-告警成功、3-告警失败 - if (info!=null && info.getAlarmEmail()!=null && info.getAlarmEmail().trim().length()>0) { - boolean alarmResult = XxlJobAdminConfig.getAdminConfig().getJobAlarmer().alarm(info, log); - newAlarmStatus = alarmResult?2:3; - } else { - newAlarmStatus = 1; - } - - XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, -1, newAlarmStatus); - } - } - - } catch (Exception e) { - if (!toStop) { - logger.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e); - } - } - - try { - TimeUnit.SECONDS.sleep(10); - } catch (Exception e) { - if (!toStop) { - logger.error(e.getMessage(), e); - } - } + }); + monitorThread.setDaemon(true); + monitorThread.setName("xxl-job, admin JobFailMonitorHelper"); + monitorThread.start(); + } + private void getWhile() { + while (!toStop) { + try { + List failLogIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findFailJobLogIds(1000); + if (failLogIds != null && !failLogIds.isEmpty()) { + getFor(failLogIds); } - logger.info(">>>>>>>>>>> xxl-job, job fail monitor thread stop"); + } catch (Exception e) { + if (!toStop) { + log.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e); + } + } + try { + TimeUnit.SECONDS.sleep(10); + } catch (Exception e) { + if (!toStop) { + log.error(e.getMessage(), e); + } + Thread.currentThread().interrupt(); + } - }); - monitorThread.setDaemon(true); - monitorThread.setName("xxl-job, admin JobFailMonitorHelper"); - monitorThread.start(); - } + } + } - public void toStop(){ - toStop = true; - // interrupt and wait - monitorThread.interrupt(); - try { - monitorThread.join(); - } catch (InterruptedException e) { - logger.error(e.getMessage(), e); - } - } + private void getFor(List failLogIds) { + for (long failLogId : failLogIds) { + + // lock log + int lockRet = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, 0, -1); + if (lockRet < 1) { + continue; + } + XxlJobLog log = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().load(failLogId); + XxlJobInfo info = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().loadById(log.getJobId()); + + // 1、fail retry monitor + if (log.getExecutorFailRetryCount() > 0) { + JobTriggerPoolHelper.trigger(log.getJobId(), TriggerTypeEnum.RETRY, (log.getExecutorFailRetryCount() - 1), log.getExecutorShardingParam(), log.getExecutorParam(), null); + String retryMsg = "

>>>>>>>>>>>" + I18nUtil.getString("jobconf_trigger_type_retry") + "<<<<<<<<<<<
"; + log.setTriggerMsg(log.getTriggerMsg() + retryMsg); + XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(log); + } + + // 2、fail alarm monitor + int newAlarmStatus = 0; // 告警状态:0-默认、-1=锁定状态、1-无需告警、2-告警成功、3-告警失败 + if (info != null && info.getAlarmEmail() != null && info.getAlarmEmail().trim().length() > 0) { + boolean alarmResult = XxlJobAdminConfig.getAdminConfig().getJobAlarmer().alarm(info, log); + newAlarmStatus = alarmResult ? 2 : 3; + } else { + newAlarmStatus = 1; + } + + XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, -1, newAlarmStatus); + } + } + + public void toStop() { + toStop = true; + // interrupt and wait + monitorThread.interrupt(); + try { + monitorThread.join(); + } catch (InterruptedException e) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLogReportHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLogReportHelper.java index 7dcce347..fb38c307 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLogReportHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLogReportHelper.java @@ -2,8 +2,7 @@ package com.xxl.job.admin.core.thread; import com.xxl.job.admin.core.conf.XxlJobAdminConfig; import com.xxl.job.admin.core.model.XxlJobLogReport; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import lombok.extern.slf4j.Slf4j; import java.util.Calendar; import java.util.Date; @@ -16,131 +15,142 @@ import java.util.concurrent.TimeUnit; * * @author xuxueli 2019-11-22 */ +@Slf4j public class JobLogReportHelper { - private static Logger logger = LoggerFactory.getLogger(JobLogReportHelper.class); private static JobLogReportHelper instance = new JobLogReportHelper(); - public static JobLogReportHelper getInstance(){ + + public static JobLogReportHelper getInstance() { return instance; } private Thread logrThread; private volatile boolean toStop = false; - public void start(){ - logrThread = new Thread(()-> { - // last clean log time - long lastCleanLogTime = 0; + public void start() { + logrThread = new Thread(() -> { + // last clean log time + long lastCleanLogTime = 0; + getWhile(lastCleanLogTime); - - while (!toStop) { - - // 1、log-report refresh: refresh log report in 3 days - try { - - for (int i = 0; i < 3; i++) { - - // today - Calendar itemDay = Calendar.getInstance(); - itemDay.add(Calendar.DAY_OF_MONTH, -i); - itemDay.set(Calendar.HOUR_OF_DAY, 0); - itemDay.set(Calendar.MINUTE, 0); - itemDay.set(Calendar.SECOND, 0); - itemDay.set(Calendar.MILLISECOND, 0); - - Date todayFrom = itemDay.getTime(); - - itemDay.set(Calendar.HOUR_OF_DAY, 23); - itemDay.set(Calendar.MINUTE, 59); - itemDay.set(Calendar.SECOND, 59); - itemDay.set(Calendar.MILLISECOND, 999); - - Date todayTo = itemDay.getTime(); - - // refresh log-report every minute - XxlJobLogReport xxlJobLogReport = new XxlJobLogReport(); - xxlJobLogReport.setTriggerDay(todayFrom); - xxlJobLogReport.setRunningCount(0); - xxlJobLogReport.setSucCount(0); - xxlJobLogReport.setFailCount(0); - - Map triggerCountMap = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findLogReport(todayFrom, todayTo); - if (triggerCountMap!=null && triggerCountMap.size()>0) { - int triggerDayCount = triggerCountMap.containsKey("triggerDayCount")?Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCount"))):0; - int triggerDayCountRunning = triggerCountMap.containsKey("triggerDayCountRunning")?Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCountRunning"))):0; - int triggerDayCountSuc = triggerCountMap.containsKey("triggerDayCountSuc")?Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCountSuc"))):0; - int triggerDayCountFail = triggerDayCount - triggerDayCountRunning - triggerDayCountSuc; - - xxlJobLogReport.setRunningCount(triggerDayCountRunning); - xxlJobLogReport.setSucCount(triggerDayCountSuc); - xxlJobLogReport.setFailCount(triggerDayCountFail); - } - - // do refresh - int ret = XxlJobAdminConfig.getAdminConfig().getXxlJobLogReportDao().update(xxlJobLogReport); - if (ret < 1) { - XxlJobAdminConfig.getAdminConfig().getXxlJobLogReportDao().save(xxlJobLogReport); - } - } - - } catch (Exception e) { - if (!toStop) { - logger.error(">>>>>>>>>>> xxl-job, job log report thread error:{}", e); - } - } - - // 2、log-clean: switch open & once each day - if (XxlJobAdminConfig.getAdminConfig().getLogretentiondays()>0 - && System.currentTimeMillis() - lastCleanLogTime > 24*60*60*1000) { - - // expire-time - Calendar expiredDay = Calendar.getInstance(); - expiredDay.add(Calendar.DAY_OF_MONTH, -1 * XxlJobAdminConfig.getAdminConfig().getLogretentiondays()); - expiredDay.set(Calendar.HOUR_OF_DAY, 0); - expiredDay.set(Calendar.MINUTE, 0); - expiredDay.set(Calendar.SECOND, 0); - expiredDay.set(Calendar.MILLISECOND, 0); - Date clearBeforeTime = expiredDay.getTime(); - - // clean expired log - List logIds = null; - do { - logIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findClearLogIds(0, 0, clearBeforeTime, 0, 1000); - if (logIds!=null && !logIds.isEmpty()) { - XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().clearLog(logIds); - } - } while (logIds!=null && !logIds.isEmpty()); - - // update clean time - lastCleanLogTime = System.currentTimeMillis(); - } - - try { - TimeUnit.MINUTES.sleep(1); - } catch (Exception e) { - if (!toStop) { - logger.error(e.getMessage(), e); - } - } - - } - - logger.info(">>>>>>>>>>> xxl-job, job log report thread stop"); + log.info(">>>>>>>>>>> xxl-job, job log report thread stop"); }); logrThread.setDaemon(true); logrThread.setName("xxl-job, admin JobLogReportHelper"); logrThread.start(); } - public void toStop(){ + private void getWhile(long lastCleanLogTime) { + while (!toStop) { + // 1、log-report refresh: refresh log report in 3 days + try { + getFor(); + + } catch (Exception e) { + if (!toStop) { + log.error(">>>>>>>>>>> xxl-job, job log report thread error:{}", e); + } + } + + // 2、log-clean: switch open & once each day + lastCleanLogTime = getLastCleanLogTime(lastCleanLogTime); + + try { + TimeUnit.MINUTES.sleep(1); + } catch (Exception e) { + if (!toStop) { + log.error(e.getMessage(), e); + } + Thread.currentThread().interrupt(); + } + + } + } + + private long getLastCleanLogTime(long lastCleanLogTime) { + if (XxlJobAdminConfig.getAdminConfig().getLogretentiondays() > 0 + && System.currentTimeMillis() - lastCleanLogTime > 24 * 60 * 60 * 1000) { + + // expire-time + Calendar expiredDay = Calendar.getInstance(); + expiredDay.add(Calendar.DAY_OF_MONTH, -1 * XxlJobAdminConfig.getAdminConfig().getLogretentiondays()); + expiredDay.set(Calendar.HOUR_OF_DAY, 0); + expiredDay.set(Calendar.MINUTE, 0); + expiredDay.set(Calendar.SECOND, 0); + expiredDay.set(Calendar.MILLISECOND, 0); + Date clearBeforeTime = expiredDay.getTime(); + + // clean expired log + List logIds = null; + do { + logIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findClearLogIds(0, 0, clearBeforeTime, 0, 1000); + if (logIds != null && !logIds.isEmpty()) { + XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().clearLog(logIds); + } + } while (logIds != null && !logIds.isEmpty()); + + // update clean time + lastCleanLogTime = System.currentTimeMillis(); + } + return lastCleanLogTime; + } + + private void getFor() { + for (int i = 0; i < 3; i++) { + // today + Calendar itemDay = Calendar.getInstance(); + itemDay.add(Calendar.DAY_OF_MONTH, -i); + itemDay.set(Calendar.HOUR_OF_DAY, 0); + itemDay.set(Calendar.MINUTE, 0); + itemDay.set(Calendar.SECOND, 0); + itemDay.set(Calendar.MILLISECOND, 0); + + Date todayFrom = itemDay.getTime(); + + itemDay.set(Calendar.HOUR_OF_DAY, 23); + itemDay.set(Calendar.MINUTE, 59); + itemDay.set(Calendar.SECOND, 59); + itemDay.set(Calendar.MILLISECOND, 999); + + Date todayTo = itemDay.getTime(); + + // refresh log-report every minute + XxlJobLogReport xxlJobLogReport = new XxlJobLogReport(); + xxlJobLogReport.setTriggerDay(todayFrom); + xxlJobLogReport.setRunningCount(0); + xxlJobLogReport.setSucCount(0); + xxlJobLogReport.setFailCount(0); + + Map triggerCountMap = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findLogReport(todayFrom, todayTo); + if (triggerCountMap != null && triggerCountMap.size() > 0) { + int triggerDayCount = triggerCountMap.containsKey("triggerDayCount") ? Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCount"))) : 0; + int triggerDayCountRunning = triggerCountMap.containsKey("triggerDayCountRunning") ? Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCountRunning"))) : 0; + int triggerDayCountSuc = triggerCountMap.containsKey("triggerDayCountSuc") ? Integer.valueOf(String.valueOf(triggerCountMap.get("triggerDayCountSuc"))) : 0; + int triggerDayCountFail = triggerDayCount - triggerDayCountRunning - triggerDayCountSuc; + + xxlJobLogReport.setRunningCount(triggerDayCountRunning); + xxlJobLogReport.setSucCount(triggerDayCountSuc); + xxlJobLogReport.setFailCount(triggerDayCountFail); + } + + // do refresh + int ret = XxlJobAdminConfig.getAdminConfig().getXxlJobLogReportDao().update(xxlJobLogReport); + if (ret < 1) { + XxlJobAdminConfig.getAdminConfig().getXxlJobLogReportDao().save(xxlJobLogReport); + } + } + } + + public void toStop() { toStop = true; // interrupt and wait logrThread.interrupt(); try { logrThread.join(); } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); } } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLosedMonitorHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLosedMonitorHelper.java index b5575253..5a75d4e3 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLosedMonitorHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobLosedMonitorHelper.java @@ -5,8 +5,7 @@ import com.xxl.job.admin.core.model.XxlJobLog; import com.xxl.job.admin.core.util.I18nUtil; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.util.DateUtil; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import lombok.extern.slf4j.Slf4j; import java.util.Date; import java.util.List; @@ -17,70 +16,81 @@ import java.util.concurrent.TimeUnit; * * @author xuxueli 2015-9-1 18:05:56 */ +@Slf4j public class JobLosedMonitorHelper { - private static Logger logger = LoggerFactory.getLogger(JobLosedMonitorHelper.class); - private static JobLosedMonitorHelper instance = new JobLosedMonitorHelper(); - public static JobLosedMonitorHelper getInstance(){ - return instance; - } + private static JobLosedMonitorHelper instance = new JobLosedMonitorHelper(); - // ---------------------- monitor ---------------------- + public static JobLosedMonitorHelper getInstance() { + return instance; + } - private Thread monitorThread; - private volatile boolean toStop = false; - public void start(){ - monitorThread = new Thread(()->{ - // monitor - while (!toStop) { - try { - // 任务结果丢失处理:调度记录停留在 "运行中" 状态超过10min,且对应执行器心跳注册失败不在线,则将本地调度主动标记失败; - Date losedTime = DateUtil.addMinutes(new Date(), -10); - List losedJobIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findLostJobIds(losedTime); + // ---------------------- monitor ---------------------- - if (losedJobIds!=null && !losedJobIds.isEmpty()) { - for (Long logId: losedJobIds) { + private Thread monitorThread; + private volatile boolean toStop = false; - XxlJobLog jobLog = new XxlJobLog(); - jobLog.setId(logId); + public void start() { + monitorThread = new Thread(() -> { + // monitor + getWhile(); + log.info(">>>>>>>>>>> xxl-job, JobLosedMonitorHelper stop"); + }); + monitorThread.setDaemon(true); + monitorThread.setName("xxl-job, admin JobLosedMonitorHelper"); + monitorThread.start(); + } - jobLog.setHandleTime(new Date()); - jobLog.setHandleCode(ReturnT.FAIL_CODE); - jobLog.setHandleMsg( I18nUtil.getString("joblog_lost_fail") ); + private void getWhile() { + while (!toStop) { + try { + // 任务结果丢失处理:调度记录停留在 "运行中" 状态超过10min,且对应执行器心跳注册失败不在线,则将本地调度主动标记失败; + Date losedTime = DateUtil.addMinutes(new Date(), -10); + List losedJobIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findLostJobIds(losedTime); - XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateHandleInfo(jobLog); - } - - } - } catch (Exception e) { - if (!toStop) { - logger.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e); - } - } - try { - TimeUnit.SECONDS.sleep(60); - } catch (Exception e) { - if (!toStop) { - logger.error(e.getMessage(), e); - } - } + if (losedJobIds != null && !losedJobIds.isEmpty()) { + getFor(losedJobIds); } - logger.info(">>>>>>>>>>> xxl-job, JobLosedMonitorHelper stop"); - }); - monitorThread.setDaemon(true); - monitorThread.setName("xxl-job, admin JobLosedMonitorHelper"); - monitorThread.start(); - } + } catch (Exception e) { + if (!toStop) { + log.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e); + } + } + try { + TimeUnit.SECONDS.sleep(60); + } catch (Exception e) { + if (!toStop) { + log.error(e.getMessage(), e); + } + Thread.currentThread().interrupt(); + } + } + } - public void toStop(){ - toStop = true; - // interrupt and wait - monitorThread.interrupt(); - try { - monitorThread.join(); - } catch (InterruptedException e) { - logger.error(e.getMessage(), e); - } - } + private void getFor(List losedJobIds) { + for (Long logId : losedJobIds) { + + XxlJobLog jobLog = new XxlJobLog(); + jobLog.setId(logId); + + jobLog.setHandleTime(new Date()); + jobLog.setHandleCode(ReturnT.FAIL_CODE); + jobLog.setHandleMsg(I18nUtil.getString("joblog_lost_fail")); + + XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateHandleInfo(jobLog); + } + } + + public void toStop() { + toStop = true; + // interrupt and wait + monitorThread.interrupt(); + try { + monitorThread.join(); + } catch (InterruptedException e) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobRegistryMonitorHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobRegistryMonitorHelper.java index 804f6774..01f12ec1 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobRegistryMonitorHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobRegistryMonitorHelper.java @@ -4,105 +4,124 @@ import com.xxl.job.admin.core.conf.XxlJobAdminConfig; import com.xxl.job.admin.core.model.XxlJobGroup; import com.xxl.job.admin.core.model.XxlJobRegistry; import com.xxl.job.core.enums.RegistryConfig; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import lombok.extern.slf4j.Slf4j; import java.util.*; import java.util.concurrent.TimeUnit; /** * job registry instance + * * @author xuxueli 2016-10-02 19:10:24 */ +@Slf4j public class JobRegistryMonitorHelper { - private static Logger logger = LoggerFactory.getLogger(JobRegistryMonitorHelper.class); - private static JobRegistryMonitorHelper instance = new JobRegistryMonitorHelper(); - public static JobRegistryMonitorHelper getInstance(){ - return instance; - } + private static JobRegistryMonitorHelper instance = new JobRegistryMonitorHelper(); - private Thread registryThread; - private volatile boolean toStop = false; - public void start(){ - registryThread = new Thread(()-> { - while (!toStop) { - try { - // auto registry group - List groupList = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().findByAddressType(0); - if (groupList!=null && !groupList.isEmpty()) { + public static JobRegistryMonitorHelper getInstance() { + return instance; + } - // remove dead address (admin/executor) - List ids = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findDead(RegistryConfig.DEAD_TIMEOUT, new Date()); - if (ids!=null && !ids.isEmpty()) { - XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().removeDead(ids); - } + private Thread registryThread; + private volatile boolean toStop = false; - // fresh online address (admin/executor) - HashMap> appAddressMap = new HashMap<>(); - List list = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findAll(RegistryConfig.DEAD_TIMEOUT, new Date()); - if (list != null) { - for (XxlJobRegistry item: list) { - if (RegistryConfig.RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) { - String appname = item.getRegistryKey(); - List registryList = appAddressMap.get(appname); - if (registryList == null) { - registryList = new ArrayList<>(); - } + public void start() { + registryThread = new Thread(() -> { + getWhile(); + log.info(">>>>>>>>>>> xxl-job, job registry monitor thread stop"); + }); + registryThread.setDaemon(true); + registryThread.setName("xxl-job, admin JobRegistryMonitorHelper"); + registryThread.start(); + } - if (!registryList.contains(item.getRegistryValue())) { - registryList.add(item.getRegistryValue()); - } - appAddressMap.put(appname, registryList); - } - } - } + private void getWhile() { + while (!toStop) { + try { + // auto registry group + List groupList = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().findByAddressType(0); + getGroupList(groupList); + } catch (Exception e) { + if (!toStop) { + log.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e); + } + } + try { + TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT); + } catch (InterruptedException e) { + if (!toStop) { + log.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e); + } + Thread.currentThread().interrupt(); + } + } + } - // fresh group address - for (XxlJobGroup group: groupList) { - List registryList = appAddressMap.get(group.getAppname()); - StringBuilder addressListStr = null; - if (registryList!=null && !registryList.isEmpty()) { - Collections.sort(registryList); - addressListStr = new StringBuilder(); - for (String item:registryList) { - addressListStr.append(item).append(","); - } - addressListStr = new StringBuilder(addressListStr.substring(0, addressListStr.length() - 1)); - } - group.setAddressList(addressListStr.toString()); - XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().update(group); - } - } - } catch (Exception e) { - if (!toStop) { - logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e); - } - } - try { - TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT); - } catch (InterruptedException e) { - if (!toStop) { - logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e); - } - } - } - logger.info(">>>>>>>>>>> xxl-job, job registry monitor thread stop"); - }); - registryThread.setDaemon(true); - registryThread.setName("xxl-job, admin JobRegistryMonitorHelper"); - registryThread.start(); - } + private void getGroupList(List groupList) { + if (groupList != null && !groupList.isEmpty()) { - public void toStop(){ - toStop = true; - // interrupt and wait - registryThread.interrupt(); - try { - registryThread.join(); - } catch (InterruptedException e) { - logger.error(e.getMessage(), e); - } - } + // remove dead address (admin/executor) + List ids = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findDead(RegistryConfig.DEAD_TIMEOUT, new Date()); + if (ids != null && !ids.isEmpty()) { + XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().removeDead(ids); + } + + // fresh online address (admin/executor) + HashMap> appAddressMap = new HashMap<>(); + List list = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findAll(RegistryConfig.DEAD_TIMEOUT, new Date()); + if (list != null) { + getFor(appAddressMap, list); + } + + // fresh group address + getFreshGroupAddress(groupList, appAddressMap); + } + } + + private void getFreshGroupAddress(List groupList, HashMap> appAddressMap) { + for (XxlJobGroup group : groupList) { + List registryList = appAddressMap.get(group.getAppname()); + StringBuilder addressListStr = new StringBuilder(); + if (registryList != null && !registryList.isEmpty()) { + Collections.sort(registryList); + for (String item : registryList) { + addressListStr.append(item).append(","); + } + addressListStr = new StringBuilder(addressListStr.substring(0, addressListStr.length() - 1)); + } + group.setAddressList(addressListStr.toString()); + XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().update(group); + } + } + + private void getFor(HashMap> appAddressMap, List list) { + for (XxlJobRegistry item : list) { + if (RegistryConfig.RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) { + String appname = item.getRegistryKey(); + List registryList = appAddressMap.get(appname); + if (registryList == null) { + registryList = new ArrayList<>(); + } + + if (!registryList.contains(item.getRegistryValue())) { + registryList.add(item.getRegistryValue()); + } + appAddressMap.put(appname, registryList); + } + } + } + + public void toStop() { + toStop = true; + // interrupt and wait + registryThread.interrupt(); + try { + registryThread.join(); + } catch (InterruptedException e) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } } 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 ed664337..813fd9da 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 @@ -4,8 +4,7 @@ import com.xxl.job.admin.core.conf.XxlJobAdminConfig; import com.xxl.job.admin.core.cron.CronExpression; import com.xxl.job.admin.core.model.XxlJobInfo; import com.xxl.job.admin.core.trigger.TriggerTypeEnum; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import lombok.extern.slf4j.Slf4j; import java.sql.Connection; import java.sql.PreparedStatement; @@ -18,11 +17,12 @@ import java.util.concurrent.TimeUnit; /** * @author xuxueli 2019-05-21 */ +@Slf4j public class JobScheduleHelper { - private static Logger logger = LoggerFactory.getLogger(JobScheduleHelper.class); private static JobScheduleHelper instance = new JobScheduleHelper(); - public static JobScheduleHelper getInstance(){ + + public static JobScheduleHelper getInstance() { return instance; } @@ -32,169 +32,14 @@ public class JobScheduleHelper { private Thread ringThread; private volatile boolean scheduleThreadToStop = false; private volatile boolean ringThreadToStop = false; - private static volatile Map> ringData = new ConcurrentHashMap<>(); + private static Map> ringData = new ConcurrentHashMap<>(); - public void start(){ + public void start() { // schedule thread - 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; - - boolean preReadSuc = true; - try { - - conn = XxlJobAdminConfig.getAdminConfig().getDataSource().getConnection(); - connAutoCommit = conn.getAutoCommit(); - conn.setAutoCommit(false); - - preparedStatement = conn.prepareStatement( "select * from xxl_job_lock where lock_name = 'schedule_lock' for update" ); - preparedStatement.execute(); - - // tx start - - // 1、pre read - long nowTime = System.currentTimeMillis(); - List scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount); - if (scheduleList!=null && !scheduleList.isEmpty()) { - // 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 - logger.warn(">>>>>>>>>>> xxl-job, schedule misfire, jobId = " + jobInfo.getId()); - - // fresh next - refreshNextValidTime(jobInfo, new Date()); - - } else if (nowTime > jobInfo.getTriggerNextTime()) { - // 2.2、trigger-expire < 5s:direct-trigger && make next-trigger-time - - // 1、trigger - JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null); - logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = {}" , jobInfo.getId() ); - - // 2、fresh next - refreshNextValidTime(jobInfo, new Date()); - - // next-trigger-time in 5s, pre-read again - if (jobInfo.getTriggerStatus()==1 && nowTime + PRE_READ_MS > jobInfo.getTriggerNextTime()) { - - // 1、make ring second - int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60); - - // 2、push time ring - pushTimeRing(ringSecond, jobInfo.getId()); - - // 3、fresh next - refreshNextValidTime(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()); - - // 3、fresh next - refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime())); - - } - - } - - // 3、update trigger info - for (XxlJobInfo jobInfo: scheduleList) { - XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleUpdate(jobInfo); - } - - } else { - preReadSuc = false; - } - - // tx stop - - - } catch (Exception e) { - if (!scheduleThreadToStop) { - logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread error:{}", e); - } - } finally { - - // commit - if (conn != null) { - try { - conn.commit(); - } catch (SQLException e) { - if (!scheduleThreadToStop) { - logger.error(e.getMessage(), e); - } - } - try { - conn.setAutoCommit(connAutoCommit); - } catch (SQLException e) { - if (!scheduleThreadToStop) { - logger.error(e.getMessage(), e); - } - } - try { - conn.close(); - } catch (SQLException e) { - if (!scheduleThreadToStop) { - logger.error(e.getMessage(), e); - } - } - } - - // close PreparedStatement - if (null != preparedStatement) { - try { - preparedStatement.close(); - } catch (SQLException e) { - if (!scheduleThreadToStop) { - logger.error(e.getMessage(), e); - } - } - } - } - 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 (InterruptedException e) { - if (!scheduleThreadToStop) { - logger.error(e.getMessage(), e); - Thread.currentThread().interrupt(); - } - } - } - } - logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop"); + scheduleThread = new Thread(() -> { + getThread(); + log.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop"); }); scheduleThread.setDaemon(true); scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread"); @@ -202,65 +47,245 @@ public class JobScheduleHelper { // ring thread - 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(); - } - } - - while (!ringThreadToStop) { - - try { - // second data - List ringItemData = new ArrayList<>(); - int nowSecond = Calendar.getInstance().get(Calendar.SECOND); // 避免处理耗时太长,跨过刻度,向前校验一个刻度; - for (int i = 0; i < 2; i++) { - List tmpData = ringData.remove( (nowSecond+60-i)%60 ); - if (tmpData != null) { - ringItemData.addAll(tmpData); - } - } - - // ring trigger - logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) ); - if (!ringItemData.isEmpty()) { - // do trigger - for (int jobId: ringItemData) { - // do trigger - JobTriggerPoolHelper.trigger(jobId, TriggerTypeEnum.CRON, -1, null, null, null); - } - // clear - ringItemData.clear(); - } - } catch (Exception e) { - if (!ringThreadToStop) { - logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e); - } - } - - // next second, align second - try { - TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000); - } catch (InterruptedException e) { - if (!ringThreadToStop) { - logger.error(e.getMessage(), e); - Thread.currentThread().interrupt(); - } - - } - } - logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop"); + ringThread = new Thread(() -> { + // align second + getThread1(); + log.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop"); }); ringThread.setDaemon(true); ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread"); ringThread.start(); } + private void getThread1() { + try { + TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis() % 1000); + } catch (InterruptedException e) { + if (!ringThreadToStop) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } + + while (!ringThreadToStop) { + + try { + // second data + getSecondData(); + } catch (Exception e) { + if (!ringThreadToStop) { + log.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e); + } + } + + // next second, align second + try { + TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis() % 1000); + } catch (InterruptedException e) { + if (!ringThreadToStop) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + + } + } + } + + private void getSecondData() { + List ringItemData = new ArrayList<>(); + int nowSecond = Calendar.getInstance().get(Calendar.SECOND); // 避免处理耗时太长,跨过刻度,向前校验一个刻度; + for (int i = 0; i < 2; i++) { + List tmpData = ringData.remove((nowSecond + 60 - i) % 60); + if (tmpData != null) { + ringItemData.addAll(tmpData); + } + } + + // ring trigger + log.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData)); + if (!ringItemData.isEmpty()) { + // do trigger + for (int jobId : ringItemData) { + // do trigger + JobTriggerPoolHelper.trigger(jobId, TriggerTypeEnum.CRON, -1, null, null, null); + } + // clear + ringItemData.clear(); + } + } + + private void getThread() { + try { + TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis() % 1000); + } catch (InterruptedException e) { + if (!scheduleThreadToStop) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } + log.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 + getWhile(preReadCount); + } + } + + private void getWhile(int preReadCount) { + long start = System.currentTimeMillis(); + Connection conn = null; + Boolean connAutoCommit = null; + PreparedStatement preparedStatement = null; + + boolean preReadSuc = true; + try { + + conn = XxlJobAdminConfig.getAdminConfig().getDataSource().getConnection(); + connAutoCommit = conn.getAutoCommit(); + conn.setAutoCommit(false); + + preparedStatement = conn.prepareStatement("select * from xxl_job_lock where lock_name = 'schedule_lock' for update"); + preparedStatement.execute(); + + // tx start + + // 1、pre read + long nowTime = System.currentTimeMillis(); + List scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount); + if (scheduleList != null && !scheduleList.isEmpty()) { + // 2、push time-ring + getFor(nowTime, scheduleList); + + // 3、update trigger info + for (XxlJobInfo jobInfo : scheduleList) { + XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleUpdate(jobInfo); + } + + } else { + preReadSuc = false; + } + + // tx stop + + + } catch (Exception e) { + if (!scheduleThreadToStop) { + log.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread error:{}", e); + } + } finally { + + // commit + getFinally(conn, connAutoCommit, preparedStatement); + } + long cost = System.currentTimeMillis() - start; + + + // Wait seconds, align second + getCost(preReadSuc, cost); + } + + private void getCost(boolean preReadSuc, long cost) { + if (cost < 1000) { // scan-overtime, not wait + try { + TimeUnit.MILLISECONDS.sleep((preReadSuc ? 1000 : PRE_READ_MS) - System.currentTimeMillis() % 1000); + } catch (InterruptedException e) { + if (!scheduleThreadToStop) { + log.error(e.getMessage(), e); + Thread.currentThread().interrupt(); + } + } + } + } + + private void getFinally(Connection conn, Boolean connAutoCommit, PreparedStatement preparedStatement) { + if (conn != null) { + try { + conn.commit(); + } catch (SQLException e) { + getScheduleThreadToStop(e); + } + try { + conn.setAutoCommit(connAutoCommit); + } catch (SQLException e) { + getScheduleThreadToStop(e); + } + try { + conn.close(); + } catch (SQLException e) { + getScheduleThreadToStop(e); + } + } + + // close PreparedStatement + if (null != preparedStatement) { + try { + preparedStatement.close(); + } catch (SQLException e) { + getScheduleThreadToStop(e); + } + } + } + + private void getScheduleThreadToStop(SQLException e) { + if (!scheduleThreadToStop) { + log.error(e.getMessage(), e); + } + } + + private void getFor(long nowTime, List scheduleList) throws ParseException { + for (XxlJobInfo jobInfo : scheduleList) { + + // time-ring jump + if (nowTime > jobInfo.getTriggerNextTime() + PRE_READ_MS) { + // 2.1、trigger-expire > 5s:pass && make next-trigger-time + log.warn(">>>>>>>>>>> xxl-job, schedule misfire, jobId = " + jobInfo.getId()); + + // fresh next + refreshNextValidTime(jobInfo, new Date()); + + } else if (nowTime > jobInfo.getTriggerNextTime()) { + // 2.2、trigger-expire < 5s:direct-trigger && make next-trigger-time + + // 1、trigger + JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null); + log.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = {}", jobInfo.getId()); + + // 2、fresh next + refreshNextValidTime(jobInfo, new Date()); + + // next-trigger-time in 5s, pre-read again + if (jobInfo.getTriggerStatus() == 1 && nowTime + PRE_READ_MS > jobInfo.getTriggerNextTime()) { + + // 1、make ring second + int ringSecond = (int) ((jobInfo.getTriggerNextTime() / 1000) % 60); + + // 2、push time ring + pushTimeRing(ringSecond, jobInfo.getId()); + + // 3、fresh next + refreshNextValidTime(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()); + + // 3、fresh next + refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime())); + + } + + } + } + private void refreshNextValidTime(XxlJobInfo jobInfo, Date fromTime) throws ParseException { Date nextValidTime = new CronExpression(jobInfo.getJobCron()).getNextValidTimeAfter(fromTime); if (nextValidTime != null) { @@ -273,56 +298,42 @@ public class JobScheduleHelper { } } - private void pushTimeRing(int ringSecond, int jobId){ + private void pushTimeRing(int ringSecond, int jobId) { // push async ring - List ringItemData = ringData.get(ringSecond); - if (ringItemData == null) { - ringItemData = new ArrayList<>(); - ringData.put(ringSecond, ringItemData); - } + List ringItemData = ringData.computeIfAbsent(ringSecond, k -> new ArrayList<>()); ringItemData.add(jobId); - logger.debug(">>>>>>>>>>> xxl-job, schedule push time-ring : " + ringSecond + " = " + Arrays.asList(ringItemData) ); + log.debug(">>>>>>>>>>> xxl-job, schedule push time-ring : " + ringSecond + " = " + Arrays.asList(ringItemData)); } - public void toStop(){ + public void toStop() { // 1、stop schedule scheduleThreadToStop = true; try { TimeUnit.SECONDS.sleep(1); // wait } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); Thread.currentThread().interrupt(); } - if (scheduleThread.getState() != Thread.State.TERMINATED){ + if (scheduleThread.getState() != Thread.State.TERMINATED) { // interrupt and wait scheduleThread.interrupt(); try { scheduleThread.join(); } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); Thread.currentThread().interrupt(); } } // if has ring data - boolean hasRingData = false; - if (!ringData.isEmpty()) { - Set>> entries = ringData.entrySet(); - for (Map.Entry> e : entries) { - List tmpData = e.getValue(); - if (tmpData!=null && !tmpData.isEmpty()) { - hasRingData = true; - break; - } - } - } + boolean hasRingData = isHasRingData(); if (hasRingData) { try { TimeUnit.SECONDS.sleep(8); } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); Thread.currentThread().interrupt(); } } @@ -332,21 +343,36 @@ public class JobScheduleHelper { try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); Thread.currentThread().interrupt(); } - if (ringThread.getState() != Thread.State.TERMINATED){ + if (ringThread.getState() != Thread.State.TERMINATED) { // interrupt and wait ringThread.interrupt(); try { ringThread.join(); } catch (InterruptedException e) { - logger.error(e.getMessage(), e); + log.error(e.getMessage(), e); Thread.currentThread().interrupt(); } } - logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper stop"); + log.info(">>>>>>>>>>> xxl-job, JobScheduleHelper stop"); + } + + private boolean isHasRingData() { + boolean hasRingData = false; + if (!ringData.isEmpty()) { + Set>> entries = ringData.entrySet(); + for (Map.Entry> e : entries) { + List tmpData = e.getValue(); + if (tmpData != null && !tmpData.isEmpty()) { + hasRingData = true; + break; + } + } + } + return hasRingData; } } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobTriggerPoolHelper.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobTriggerPoolHelper.java index 8213ddfb..17412946 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobTriggerPoolHelper.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/thread/JobTriggerPoolHelper.java @@ -32,12 +32,6 @@ public class JobTriggerPoolHelper { TimeUnit.SECONDS, new LinkedBlockingQueue<>(1000), r-> new Thread(r, "xxl-job, admin JobTriggerPoolHelper-fastTriggerPool-" + r.hashCode())); - /*new ThreadFactory() { - @Override - public Thread newThread(Runnable r) { - return new Thread(r, "xxl-job, admin JobTriggerPoolHelper-fastTriggerPool-" + r.hashCode()); - } - });*/ slowTriggerPool = new ThreadPoolExecutor( 10, @@ -46,17 +40,10 @@ public class JobTriggerPoolHelper { TimeUnit.SECONDS, new LinkedBlockingQueue<>(2000), r-> new Thread(r, "xxl-job, admin JobTriggerPoolHelper-slowTriggerPool-" + r.hashCode())); - /*new ThreadFactory() { - @Override - public Thread newThread(Runnable r) { - return new Thread(r, "xxl-job, admin JobTriggerPoolHelper-slowTriggerPool-" + r.hashCode()); - } - });*/ } public void stop() { - //triggerPool.shutdown(); fastTriggerPool.shutdownNow(); slowTriggerPool.shutdownNow(); logger.info(">>>>>>>>> xxl-job trigger thread pool shutdown success."); @@ -65,7 +52,7 @@ public class JobTriggerPoolHelper { // job timeout count private volatile long minTim = System.currentTimeMillis()/60000; // ms > min - private volatile ConcurrentMap jobTimeoutCountMap = new ConcurrentHashMap<>(); + private ConcurrentMap jobTimeoutCountMap = new ConcurrentHashMap<>(); /** diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/trigger/XxlJobTrigger.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/trigger/XxlJobTrigger.java index 7e857061..04dafa95 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/trigger/XxlJobTrigger.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/trigger/XxlJobTrigger.java @@ -71,14 +71,7 @@ public class XxlJobTrigger { // sharding param int[] shardingParam = null; - if (executorShardingParam!=null){ - String[] shardingArr = executorShardingParam.split("/"); - if (shardingArr.length==2 && isNumeric(shardingArr[0]) && isNumeric(shardingArr[1])) { - shardingParam = new int[2]; - shardingParam[0] = Integer.valueOf(shardingArr[0]); - shardingParam[1] = Integer.valueOf(shardingArr[1]); - } - } + shardingParam = getInts(executorShardingParam, shardingParam); if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST==ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null) && group.getRegistryList()!=null && !group.getRegistryList().isEmpty() && shardingParam==null) { @@ -94,6 +87,18 @@ public class XxlJobTrigger { } + private static int[] getInts(String executorShardingParam, int[] shardingParam) { + if (executorShardingParam!=null){ + String[] shardingArr = executorShardingParam.split("/"); + if (shardingArr.length==2 && isNumeric(shardingArr[0]) && isNumeric(shardingArr[1])) { + shardingParam = new int[2]; + shardingParam[0] = Integer.valueOf(shardingArr[0]); + shardingParam[1] = Integer.valueOf(shardingArr[1]); + } + } + return shardingParam; + } + private static boolean isNumeric(String str){ try { Integer.valueOf(str); @@ -146,11 +151,7 @@ public class XxlJobTrigger { ReturnT routeAddressResult = null; if (group.getRegistryList()!=null && !group.getRegistryList().isEmpty()) { if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST == executorRouteStrategyEnum) { - if (index < group.getRegistryList().size()) { - address = group.getRegistryList().get(index); - } else { - address = group.getRegistryList().get(0); - } + address = getString(group, index); } else { routeAddressResult = executorRouteStrategyEnum.getRouter().route(triggerParam, group.getRegistryList()); if (routeAddressResult.getCode() == ReturnT.SUCCESS_CODE) { @@ -163,11 +164,7 @@ public class XxlJobTrigger { // 4、trigger remote executor ReturnT triggerResult = null; - if (address != null) { - triggerResult = runExecutor(triggerParam, address); - } else { - triggerResult = new ReturnT<>(ReturnT.FAIL_CODE, null); - } + triggerResult = getTriggerResult(triggerParam, address); // 5、collection trigger info StringBuilder triggerMsgSb = new StringBuilder(); @@ -177,9 +174,7 @@ public class XxlJobTrigger { .append( (group.getAddressType() == 0)?I18nUtil.getString("jobgroup_field_addressType_0"):I18nUtil.getString("jobgroup_field_addressType_1") ); triggerMsgSb.append("
").append(I18nUtil.getString("jobconf_trigger_exe_regaddress")).append(":").append(group.getRegistryList()); triggerMsgSb.append("
").append(I18nUtil.getString("jobinfo_field_executorRouteStrategy")).append(":").append(executorRouteStrategyEnum.getTitle()); - if (shardingParam != null) { - triggerMsgSb.append("("+shardingParam+")"); - } + triggerMsgSbAppend(shardingParam, triggerMsgSb); triggerMsgSb.append("
").append(I18nUtil.getString("jobinfo_field_executorBlockStrategy")).append(":").append(blockStrategy.getTitle()); triggerMsgSb.append("
").append(I18nUtil.getString("jobinfo_field_timeout")).append(":").append(jobInfo.getExecutorTimeout()); triggerMsgSb.append("
").append(I18nUtil.getString("jobinfo_field_executorFailRetryCount")).append(":").append(finalFailRetryCount); @@ -193,7 +188,6 @@ public class XxlJobTrigger { jobLog.setExecutorParam(jobInfo.getExecutorParam()); jobLog.setExecutorShardingParam(shardingParam); jobLog.setExecutorFailRetryCount(finalFailRetryCount); - //jobLog.setTriggerTime(); jobLog.setTriggerCode(triggerResult.getCode()); jobLog.setTriggerMsg(triggerMsgSb.toString()); XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(jobLog); @@ -201,6 +195,32 @@ public class XxlJobTrigger { logger.debug(">>>>>>>>>>> xxl-job trigger end, jobId:{}", jobLog.getId()); } + private static ReturnT getTriggerResult(TriggerParam triggerParam, String address) { + ReturnT triggerResult; + if (address != null) { + triggerResult = runExecutor(triggerParam, address); + } else { + triggerResult = new ReturnT<>(ReturnT.FAIL_CODE, null); + } + return triggerResult; + } + + private static String getString(XxlJobGroup group, int index) { + String address; + if (index < group.getRegistryList().size()) { + address = group.getRegistryList().get(index); + } else { + address = group.getRegistryList().get(0); + } + return address; + } + + private static void triggerMsgSbAppend(String shardingParam, StringBuilder triggerMsgSb) { + if (shardingParam != null) { + triggerMsgSb.append("("+shardingParam+")"); + } + } + /** * run executor * @param triggerParam diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/util/FtlUtil.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/util/FtlUtil.java index b86d4100..c772be05 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/util/FtlUtil.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/core/util/FtlUtil.java @@ -18,7 +18,7 @@ public class FtlUtil { private static Logger logger = LoggerFactory.getLogger(FtlUtil.class); - private static BeansWrapper wrapper = new BeansWrapperBuilder(Configuration.DEFAULT_INCOMPATIBLE_IMPROVEMENTS).build(); //BeansWrapper.getDefaultInstance(); + private static BeansWrapper wrapper = new BeansWrapperBuilder(Configuration.DEFAULT_INCOMPATIBLE_IMPROVEMENTS).build(); public static TemplateHashModel generateStaticModel(String packageName) { try { diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/LoginService.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/LoginService.java index 333e18c9..b90ca989 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/LoginService.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/LoginService.java @@ -6,6 +6,7 @@ import com.xxl.job.admin.core.util.I18nUtil; import com.xxl.job.admin.core.util.JacksonUtil; import com.xxl.job.admin.dao.XxlJobUserDao; import com.xxl.job.core.biz.model.ReturnT; +import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Configuration; import org.springframework.util.DigestUtils; @@ -13,10 +14,12 @@ import javax.annotation.Resource; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; import java.math.BigInteger; +import java.util.Objects; /** * @author xuxueli 2019-05-04 22:13:264 */ +@Slf4j @Configuration public class LoginService { @@ -41,7 +44,9 @@ public class LoginService { public ReturnT login(HttpServletRequest request, HttpServletResponse response, String username, String password, boolean ifRemember){ - + if(Objects.isNull(request)){ + log.info("未使用参数"); + } // param if (username==null || username.trim().length()==0 || password==null || password.trim().length()==0){ return new ReturnT<>(500, I18nUtil.getString("login_param_empty")); diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/AdminBizImpl.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/AdminBizImpl.java index c158adb7..0887a644 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/AdminBizImpl.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/AdminBizImpl.java @@ -14,6 +14,7 @@ import com.xxl.job.core.biz.model.HandleCallbackParam; import com.xxl.job.core.biz.model.RegistryParam; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.handler.IJobHandler; +import lombok.extern.slf4j.Slf4j; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; @@ -23,10 +24,12 @@ import javax.annotation.Resource; import java.text.MessageFormat; import java.util.Date; import java.util.List; +import java.util.Objects; /** * @author xuxueli 2017-07-27 21:54:20 */ +@Slf4j @Service public class AdminBizImpl implements AdminBiz { private static final Logger logger = LoggerFactory.getLogger(AdminBizImpl.class); @@ -68,30 +71,7 @@ public class AdminBizImpl implements AdminBiz { XxlJobInfo xxlJobInfo = xxlJobInfoDao.loadById(log.getJobId()); if (xxlJobInfo!=null && xxlJobInfo.getChildJobId()!=null && xxlJobInfo.getChildJobId().trim().length()>0) { callbackMsg = new StringBuilder("

>>>>>>>>>>>" + I18nUtil.getString("jobconf_trigger_child_run") + "<<<<<<<<<<<
"); - - String[] childJobIds = xxlJobInfo.getChildJobId().split(","); - for (int i = 0; i < childJobIds.length; i++) { - int childJobId = (childJobIds[i]!=null && childJobIds[i].trim().length()>0 && isNumeric(childJobIds[i]))?Integer.valueOf(childJobIds[i]):-1; - if (childJobId > 0) { - - JobTriggerPoolHelper.trigger(childJobId, TriggerTypeEnum.PARENT, -1, null, null, null); - ReturnT triggerChildResult = ReturnT.SUCCESS; - - // add msg - callbackMsg.append(MessageFormat.format(I18nUtil.getString("jobconf_callback_child_msg1"), - (i + 1), - childJobIds.length, - childJobIds[i], - (triggerChildResult.getCode() == ReturnT.SUCCESS_CODE ? I18nUtil.getString("system_success") : I18nUtil.getString("system_fail")), - triggerChildResult.getMsg())); - } else { - callbackMsg.append(MessageFormat.format(I18nUtil.getString("jobconf_callback_child_msg2"), - (i + 1), - childJobIds.length, - childJobIds[i])); - } - } - + appendCallbackMsg(callbackMsg, xxlJobInfo); } } @@ -120,6 +100,31 @@ public class AdminBizImpl implements AdminBiz { return ReturnT.SUCCESS; } + private void appendCallbackMsg(StringBuilder callbackMsg, XxlJobInfo xxlJobInfo) { + String[] childJobIds = xxlJobInfo.getChildJobId().split(","); + for (int i = 0; i < childJobIds.length; i++) { + int childJobId = (childJobIds[i]!=null && childJobIds[i].trim().length()>0 && isNumeric(childJobIds[i]))?Integer.valueOf(childJobIds[i]):-1; + if (childJobId > 0) { + + JobTriggerPoolHelper.trigger(childJobId, TriggerTypeEnum.PARENT, -1, null, null, null); + ReturnT triggerChildResult = ReturnT.SUCCESS; + + // add msg + callbackMsg.append(MessageFormat.format(I18nUtil.getString("jobconf_callback_child_msg1"), + (i + 1), + childJobIds.length, + childJobIds[i], + (triggerChildResult.getCode() == ReturnT.SUCCESS_CODE ? I18nUtil.getString("system_success") : I18nUtil.getString("system_fail")), + triggerChildResult.getMsg())); + } else { + callbackMsg.append(MessageFormat.format(I18nUtil.getString("jobconf_callback_child_msg2"), + (i + 1), + childJobIds.length, + childJobIds[i])); + } + } + } + private boolean isNumeric(String str){ try { Integer.valueOf(str); @@ -169,6 +174,9 @@ public class AdminBizImpl implements AdminBiz { } private void freshGroupRegistryInfo(RegistryParam registryParam){ + if(Objects.isNull(registryParam)){ + log.info("未使用参数"); + } // Under consideration, prevent affecting core tables } diff --git a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/XxlJobServiceImpl.java b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/XxlJobServiceImpl.java index 011518ad..2608b43b 100644 --- a/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/XxlJobServiceImpl.java +++ b/jero-boot/jero-cloud-module/jero-cloud-xxljob/src/main/java/com/xxl/job/admin/service/impl/XxlJobServiceImpl.java @@ -91,14 +91,35 @@ public class XxlJobServiceImpl implements XxlJobService { if (GlueTypeEnum.BEAN==GlueTypeEnum.match(jobInfo.getGlueType()) && (jobInfo.getExecutorHandler()==null || jobInfo.getExecutorHandler().trim().length()==0) ) { return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString(systemPleaseInput)+"JobHandler") ); } + getGlueTypeEnum(jobInfo); + // ChildJobId valid + ReturnT childJobIdItem = getChildJobId(jobInfo); + if (childJobIdItem != null) { + return childJobIdItem; + } + + // add in db + jobInfo.setAddTime(new Date()); + jobInfo.setUpdateTime(new Date()); + jobInfo.setGlueUpdatetime(new Date()); + xxlJobInfoDao.save(jobInfo); + if (jobInfo.getId() < 1) { + return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_add")+I18nUtil.getString("system_fail")) ); + } + + return new ReturnT<>(String.valueOf(jobInfo.getId())); + } + + private void getGlueTypeEnum(XxlJobInfo jobInfo) { // fix "\r" in shell if (GlueTypeEnum.GLUE_SHELL==GlueTypeEnum.match(jobInfo.getGlueType()) && jobInfo.getGlueSource()!=null) { jobInfo.setGlueSource(jobInfo.getGlueSource().replace("\r", "")); } + } - // ChildJobId valid - if (jobInfo.getChildJobId()!=null && jobInfo.getChildJobId().trim().length()>0) { + private ReturnT getChildJobId(XxlJobInfo jobInfo) { + if (jobInfo.getChildJobId()!=null && jobInfo.getChildJobId().trim().length()>0) { String[] childJobIds = jobInfo.getChildJobId().split(","); for (String childJobIdItem: childJobIds) { if (childJobIdItem!=null && childJobIdItem.trim().length()>0 && isNumeric(childJobIdItem)) { @@ -122,17 +143,7 @@ public class XxlJobServiceImpl implements XxlJobService { jobInfo.setChildJobId(temp.toString()); } - - // add in db - jobInfo.setAddTime(new Date()); - jobInfo.setUpdateTime(new Date()); - jobInfo.setGlueUpdatetime(new Date()); - xxlJobInfoDao.save(jobInfo); - if (jobInfo.getId() < 1) { - return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_add")+I18nUtil.getString("system_fail")) ); - } - - return new ReturnT<>(String.valueOf(jobInfo.getId())); + return null; } private boolean isNumeric(String str){ @@ -146,39 +157,16 @@ public class XxlJobServiceImpl implements XxlJobService { @Override public ReturnT update(XxlJobInfo jobInfo) { - - // valid - if (!CronExpression.isValidExpression(jobInfo.getJobCron())) { - return new ReturnT<>(ReturnT.FAIL_CODE, I18nUtil.getString(jobinfoFieldCronUnvalid) ); - } - if (jobInfo.getJobDesc()==null || jobInfo.getJobDesc().trim().length()==0) { - return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString(systemPleaseInput)+I18nUtil.getString("jobinfo_field_jobdesc")) ); - } - if (jobInfo.getAuthor()==null || jobInfo.getAuthor().trim().length()==0) { - return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString(systemPleaseInput)+I18nUtil.getString("jobinfo_field_author")) ); - } - if (ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null) == null) { - return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_executorRouteStrategy")+I18nUtil.getString(systemUnvalid)) ); - } - if (ExecutorBlockStrategyEnum.match(jobInfo.getExecutorBlockStrategy(), null) == null) { - return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_executorBlockStrategy")+I18nUtil.getString(systemUnvalid)) ); + ReturnT x = getStringReturnT(jobInfo); + if (x != null) { + return x; } // ChildJobId valid if (jobInfo.getChildJobId()!=null && jobInfo.getChildJobId().trim().length()>0) { String[] childJobIds = jobInfo.getChildJobId().split(","); - for (String childJobIdItem: childJobIds) { - if (childJobIdItem!=null && childJobIdItem.trim().length()>0 && isNumeric(childJobIdItem)) { - XxlJobInfo childJobInfo = xxlJobInfoDao.loadById(Integer.parseInt(childJobIdItem)); - if (childJobInfo==null) { - return new ReturnT<>(ReturnT.FAIL_CODE, - MessageFormat.format((I18nUtil.getString(jobinfoFieldChildJobId)+zero+I18nUtil.getString(systemNotFound)), childJobIdItem)); - } - } else { - return new ReturnT<>(ReturnT.FAIL_CODE, - MessageFormat.format((I18nUtil.getString(jobinfoFieldChildJobId)+zero+I18nUtil.getString(systemUnvalid)), childJobIdItem)); - } - } + ReturnT childJobIdItem = getStringReturnT(childJobIds); + if (childJobIdItem != null) return childJobIdItem; // join , avoid "xxx,," StringBuilder temp = new StringBuilder(); @@ -238,6 +226,42 @@ public class XxlJobServiceImpl implements XxlJobService { return ReturnT.SUCCESS; } + private ReturnT getStringReturnT(String[] childJobIds) { + for (String childJobIdItem: childJobIds) { + if (childJobIdItem!=null && childJobIdItem.trim().length()>0 && isNumeric(childJobIdItem)) { + XxlJobInfo childJobInfo = xxlJobInfoDao.loadById(Integer.parseInt(childJobIdItem)); + if (childJobInfo==null) { + return new ReturnT<>(ReturnT.FAIL_CODE, + MessageFormat.format((I18nUtil.getString(jobinfoFieldChildJobId)+zero+I18nUtil.getString(systemNotFound)), childJobIdItem)); + } + } else { + return new ReturnT<>(ReturnT.FAIL_CODE, + MessageFormat.format((I18nUtil.getString(jobinfoFieldChildJobId)+zero+I18nUtil.getString(systemUnvalid)), childJobIdItem)); + } + } + return null; + } + + private ReturnT getStringReturnT(XxlJobInfo jobInfo) { + // valid + if (!CronExpression.isValidExpression(jobInfo.getJobCron())) { + return new ReturnT<>(ReturnT.FAIL_CODE, I18nUtil.getString(jobinfoFieldCronUnvalid) ); + } + if (jobInfo.getJobDesc()==null || jobInfo.getJobDesc().trim().length()==0) { + return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString(systemPleaseInput)+I18nUtil.getString("jobinfo_field_jobdesc")) ); + } + if (jobInfo.getAuthor()==null || jobInfo.getAuthor().trim().length()==0) { + return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString(systemPleaseInput)+I18nUtil.getString("jobinfo_field_author")) ); + } + if (ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null) == null) { + return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_executorRouteStrategy")+I18nUtil.getString(systemUnvalid)) ); + } + if (ExecutorBlockStrategyEnum.match(jobInfo.getExecutorBlockStrategy(), null) == null) { + return new ReturnT<>(ReturnT.FAIL_CODE, (I18nUtil.getString("jobinfo_field_executorBlockStrategy")+I18nUtil.getString(systemUnvalid)) ); + } + return null; + } + @Override public ReturnT remove(int id) { XxlJobInfo xxlJobInfo = xxlJobInfoDao.loadById(id);

"+ I18nUtil.getString("jobinfo_field_jobgroup") +""+ I18nUtil.getString("jobinfo_field_id") +""+ I18nUtil.getString("jobinfo_field_jobdesc") +""+ I18nUtil.getString("jobconf_monitor_alarm_title") +""+ I18nUtil.getString("jobconf_monitor_alarm_content") +""+ I18nUtil.getString("jobinfo_field_jobgroup") +TD + + " "+ I18nUtil.getString("jobinfo_field_id") +TD + + " "+ I18nUtil.getString("jobinfo_field_jobdesc") +TD + + " "+ I18nUtil.getString("jobconf_monitor_alarm_title") +TD + + " "+ I18nUtil.getString("jobconf_monitor_alarm_content") +TD + "
{0}{1}{2}"+ I18nUtil.getString("jobconf_monitor_alarm_type") +""+ I18nUtil.getString("jobconf_monitor_alarm_type") +TD + " {3}