xxl-job这个定时任务框架应该大部分人都用过,非常强大,而且它拥有完善的管理后台可以创建执行器、管理任务、管理用户等。

在大多数需求中我们只需要在项目中创建一个标注@XxlJob的方法,然后实现定时任务要执行的逻辑,最后通过管理后台创建一个定时任务并且关联创建好的@XxlJob方法,然后固定时间去执行就好了,比如:定时对账、异常接口扫描、定时处理日志等。

但是近期我遇到一个其他的场景,就是要在创建订单的时候去添加一个定时任务,然后固定时间去执行某些逻辑。那这种场景下通过xxl-job的管理后台就无法实现了,我们不可能提前创建这个定时任务,并且也不可能提前知道订单id。

对于这种实时性比较强的或者需要绑定业务id(无法提前预知)的这类定时任务,我们就可以通过API的方式来创建,下面我结合了一些真实的业务场景封装了一个xxl-job的API模块,可以在项目的任意业务逻辑中调用。

1.下载并部署xxl-job源码

这个就不多说了,大家自行到官网下载源码,https://gitee.com/xuxueli0323/xxl-job

我用的版本是:3.3.2(应该是目前最新的)

然后修改一下 application.properties 配置文件把它跑起来就行了

在这里插入图片描述

2.构建xxl-job的API工具包模块

在这里插入图片描述

2.1 引入xxl-job核心包依赖

<!-- xxl-job-core -->
        <dependency>
            <groupId>com.xuxueli</groupId>
            <artifactId>xxl-job-core</artifactId>
            <version>3.3.2</version>
        </dependency>

2.2 创建配置文件

创建下面这三个类

/**
 * 定时任务配置信息
 */
@Data
public class ScheduledConfig {

    private String type;

    private String name;

    private String entityId;

    private String cron;

    private String handler;
}

/**
 * 自己的定时任务的配置信息
 */
@Component
@ConfigurationProperties("schedule")
@Data
public class ScheduledConfigProperties {

    private Boolean enabled;

    private List<ScheduledConfig> config;
}

/**
 * xxl-job config
 *
 * @author xuxueli 2017-04-28
 */
@Configuration
public class XxlJobConfig {
    private static final Logger logger = LoggerFactory.getLogger(XxlJobConfig.class);

    @Value("${xxl.job.admin.addresses}")
    private String adminAddresses;

    @Value("${xxl.job.admin.accessToken}")
    private String accessToken;

    @Value("${xxl.job.admin.timeout}")
    private int timeout;

    @Value("${xxl.job.executor.enabled}")
    private Boolean enabled;

    @Value("${xxl.job.executor.appname}")
    private String appname;

    @Value("${xxl.job.executor.address}")
    private String address;

    @Value("${xxl.job.executor.ip}")
    private String ip;

    @Value("${xxl.job.executor.port}")
    private int port;

    @Value("${xxl.job.executor.logpath}")
    private String logPath;

    @Value("${xxl.job.executor.logretentiondays}")
    private int logRetentionDays;

    @Value("${xxl.job.executor.excludedpackage}")
    private String excludedPackage;


    @Bean
    public XxlJobSpringExecutor xxlJobExecutor() {
        logger.info(">>>>>>>>>>> xxl-job config init.");
        XxlJobSpringExecutor xxlJobSpringExecutor = new XxlJobSpringExecutor();
        xxlJobSpringExecutor.setAdminAddresses(adminAddresses);
        xxlJobSpringExecutor.setAccessToken(accessToken);
        xxlJobSpringExecutor.setTimeout(timeout);
        xxlJobSpringExecutor.setEnabled(enabled);
        xxlJobSpringExecutor.setAppname(appname);
        xxlJobSpringExecutor.setAddress(address);
        xxlJobSpringExecutor.setIp(ip);
        xxlJobSpringExecutor.setPort(port);
        xxlJobSpringExecutor.setLogPath(logPath);
        xxlJobSpringExecutor.setLogRetentionDays(logRetentionDays);
        xxlJobSpringExecutor.setExcludedPackage(excludedPackage);

        return xxlJobSpringExecutor;
    }

}

2.3 创建实体类

@Data
public class XxlJobInfo {

    private int id;				// 主键ID

    private int jobGroup;		// 执行器主键ID
    private String jobDesc;

    private Date addTime;
    private Date updateTime;

    private String author;		// 负责人
    private String alarmEmail;	// 报警邮件

    private String scheduleType;			// 调度类型
    private String scheduleConf;			// 调度配置,值含义取决于调度类型
    private String misfireStrategy;			// 调度过期策略

    private String executorRouteStrategy;	// 执行器路由策略
    private String executorHandler;		    // 执行器,任务Handler名称
    private String executorParam;		    // 执行器,任务参数
    private String executorBlockStrategy;	// 阻塞处理策略
    private int executorTimeout;     		// 任务执行超时时间,单位秒
    private int executorFailRetryCount;		// 失败重试次数

    private String glueType;		// GLUE类型	#com.xxl.job.core.glue.GlueTypeEnum
    private String glueSource;		// GLUE源代码
    private String glueRemark;		// GLUE备注
    private Date glueUpdatetime;	// GLUE更新时间

    private String childJobId;		// 子任务ID,多个逗号分隔

    private int triggerStatus;		// 调度状态:0-停止,1-运行
    private long triggerLastTime;	// 上次调度时间
    private long triggerNextTime;	// 下次调度时间
}

@Data
public class XxlJobInfoDTO {
    private String scheduleConf;			// 调度配置,值含义取决于调度类型

    private String jobDesc;

    private String executorParam;		    // 执行器,任务参数

    private String author;		// 负责人

    private String executorHandler;		    // 执行器,任务Handler名称


    public XxlJobInfoDTO(String scheduleConf, String jobDesc, String executorParam, String author, String executorHandler) {
        this.scheduleConf = scheduleConf;
        this.jobDesc = jobDesc;
        this.executorParam = executorParam;
        this.author = author;
        this.executorHandler = executorHandler;
    }
}

2.4 创建接口和实现类

这部分接口主要用于在业务中进行调用,其中在获取Cookie方法 ensureValidCookie(); 中使用了双重检测锁,可以重点看看。

// 接口
public interface MyJobApiService {

    Response addJob(XxlJobInfoDTO jobDto);

    Response removeJob(Integer jobId);

    Response startJob(Integer jobId);

    Response stopJob(Integer jobId);

    Response triggerJob(Integer jobId, String executorParam);
}

// 实现类
@Service
@RequiredArgsConstructor
public class MyJobApiServiceImpl implements MyJobApiService {

    private static final Logger logger = LoggerFactory.getLogger(MyJobApiServiceImpl.class);

    private final RestTemplate restTemplate;

    private static String xxlJobCookie;

    @Value("${xxl.job.admin.addresses}")
    private String adminAddresses;

    @Value("${xxl.job.admin.accessToken}")
    private String accessToken;

    @Value("${xxl.job.username}")
    private String username;

    @Value("${xxl.job.password}")
    private String password;

    /**
     * 添加任务
     */
    public Response addJob(XxlJobInfoDTO jobDto) {
        XxlJobInfo xxlJobInfo = JobUtils.addJobCron(jobDto.getScheduleConf(), jobDto.getJobDesc(), jobDto.getExecutorParam(), jobDto.getAuthor(), jobDto.getExecutorHandler());
        // 调用Admin API
        String url = adminAddresses + "/jobinfo/insert";
        return postFormToAdmin(url, JSONObject.parseObject(JSONObject.toJSONString(xxlJobInfo), Map.class));
    }

    /**
     * 删除任务
     */
    public Response removeJob(Integer jobId) {
        String url = adminAddresses + "/jobinfo/delete";
        Map<String, Object> params = new HashMap<>();
        params.put("ids[]", jobId);
        return postFormToAdmin(url, params);
    }

    /**
     * 启动任务
     */
    public Response startJob(Integer jobId) {
        String url = adminAddresses + "/jobinfo/start";
        Map<String, Object> params = new HashMap<>();
        params.put("ids[]", jobId);
        return postFormToAdmin(url, params);
    }

    /**
     * 停止任务
     */
    public Response stopJob(Integer jobId) {
        String url = adminAddresses + "/jobinfo/stop";
        Map<String, Object> params = new HashMap<>();
        params.put("ids[]", jobId);
        return postFormToAdmin(url, params);
    }

    /**
     * 触发执行一次任务
     */
    public Response triggerJob(Integer jobId, String executorParam) {
        String url = adminAddresses + "/jobinfo/trigger";
        Map<String, Object> params = new HashMap<>();
        params.put("ids[]", jobId);
        params.put("executorParam", executorParam);
        return postFormToAdmin(url, params);
    }


    private Response postFormToAdmin(String url, Map<String, Object> params) {
        try {
            HttpHeaders headers = new HttpHeaders();
            headers.setContentType(MediaType.APPLICATION_FORM_URLENCODED);
            if (accessToken != null && !accessToken.isEmpty()) {
                headers.set("XXL-JOB-ACCESS-TOKEN", accessToken);
            }

            // 添加登录 cookie
            ensureValidCookie();
            headers.set("Cookie", xxlJobCookie);

            MultiValueMap<String, Object> body = new LinkedMultiValueMap<>();
            params.forEach((key, value) -> body.add(key, value.toString()));

            HttpEntity<MultiValueMap<String, Object>> request = new HttpEntity<>(body, headers);
            ResponseEntity<Response> response = restTemplate.postForEntity(url, request, Response.class);

            if (response.getStatusCode().is2xxSuccessful()) {
                return response.getBody();
            } else {
                return Response.ofFail("【MY-JOB】HTTP请求失败: " + response.getStatusCode());
            }
        } catch (Exception e) {
            return Response.ofFail("【MY-JOB】请求异常: " + e.getMessage());
        }
    }

    private void loginAndGetCookie() {
        try {
            String loginUrl = adminAddresses + "/auth/doLogin";

            HttpHeaders headers = new HttpHeaders();
            headers.setContentType(MediaType.APPLICATION_FORM_URLENCODED);

            MultiValueMap<String, String> params = new LinkedMultiValueMap<>();
            params.add("userName", username); // 从配置读取
            params.add("password", password); // 从配置读取
            params.add("ifRemember", "on"); // ifRemember=on 代表记住登录状态,所以我们可以把Cookie存起来

            HttpEntity<MultiValueMap<String, String>> request = new HttpEntity<>(params, headers);
            ResponseEntity<String> response = restTemplate.postForEntity(loginUrl, request, String.class);

            // 获取 Cookie
            List<String> cookies = response.getHeaders().get("Set-Cookie");
            if (cookies != null && !cookies.isEmpty()) {
                xxlJobCookie = cookies.get(0);
            }
        } catch (Exception e) {
            logger.error("登录XXL-Job Admin失败", e);
        }
    }

    private static final Object COOKIE_LOCK = new Object();

    // 双重检测锁
    private void ensureValidCookie() {
        if (StringUtils.isEmpty(xxlJobCookie)) {
            synchronized (COOKIE_LOCK) {
                if (StringUtils.isEmpty(xxlJobCookie)) {
                    this.loginAndGetCookie();
                }
            }
        }
    }
}

2.5 创建工具类

@Component
public class JobUtils implements InitializingBean {

    @Value("${xxl.job.email}")
    private String jobEmail;

    @Value("${xxl.job.group}")
    private Integer jobGroup;

    private static String email;
    private static Integer group;

    /**
     * 创建任务corn表达式
     *
     * @param schedule
     * @param desc
     * @param param
     * @param author
     * @param jobHandler
     * @return
     */
    public static XxlJobInfo addJobCron(String schedule, String desc,
                                        String param, String author, String jobHandler) {
        XxlJobInfo xxlJobInfo = addJob(schedule, desc, param, author, jobHandler);
        //设置调度类型
        xxlJobInfo.setScheduleType("CRON");
        return xxlJobInfo;
    }

    /**
     * 创建任务固定秒数second(单位秒)
     *
     * @param schedule
     * @param desc
     * @param param
     * @param author
     * @param jobHandler
     * @return
     */
    public static XxlJobInfo addJobSecond(String schedule, String desc,
                                   String param, String author, String jobHandler) {
        XxlJobInfo xxlJobInfo = addJob(schedule, desc, param, author, jobHandler);
        //设置调度类型
        xxlJobInfo.setScheduleType("FIX_RATE");
        return xxlJobInfo;
    }

    /**
     * 通用job添加方法
     *
     * @param desc       描述
     * @param param      参数
     * @param author     负责人
     * @param jobHandler 任务执行器
     */
    public static XxlJobInfo addJob(String schedule, String desc,
                             String param, String author, String jobHandler) {
        //构建job对象
        XxlJobInfo xxlJobInfo = new XxlJobInfo();
        xxlJobInfo.setJobGroup(group);
        //设置desc
        xxlJobInfo.setJobDesc(desc);
        //设置job定时器
        xxlJobInfo.setScheduleConf(schedule);
        //设置运行模式  GLUE_GROOVY,JAVA
        xxlJobInfo.setGlueType(GlueTypeEnum.BEAN.getDesc());
        //设置job处理器
        xxlJobInfo.setExecutorHandler(jobHandler);
        //设置执行参数
        xxlJobInfo.setExecutorParam(param);
        //设置路由策略
        xxlJobInfo.setExecutorRouteStrategy("FIRST");
        //调度过期策略    FIRE_ONCE_NOW,立即执行一次
        xxlJobInfo.setMisfireStrategy("DO_NOTHING");
        //设置阻塞处理策略  COVER_EARLY,覆盖之前调度
        xxlJobInfo.setExecutorBlockStrategy("SERIAL_EXECUTION");
        //设置失败重置次数
        xxlJobInfo.setExecutorFailRetryCount(3);
        //设置负责人
        xxlJobInfo.setAuthor(author);
        //设置警告邮箱
        xxlJobInfo.setAlarmEmail(email);
        //设置启动状态
        xxlJobInfo.setTriggerStatus(1);
        return xxlJobInfo;
    }

    @Override
    public void afterPropertiesSet() throws Exception {
        email = jobEmail;
        group = jobGroup;
    }
}

至此我们的API工具包就构建好了

3.如何在业务中使用

引入以上创建好的工具包依赖

<dependency>
            <groupId>com.xxx</groupId>
            <artifactId>healthnow-xxl-job</artifactId>
            <version>${project.version}</version>
        </dependency>

application配置文件中添加定时任务相关配置

xxl:
  job:
    email: 
    group: 4 # 执行器主键ID
    username: your_username
    password: your_password
    admin:
      addresses: http://xxx:8092/xxl-job-admin
      accessToken: your_access_token
      timeout: 3
    executor:
      enabled: true
      appname: your_appname
      address:
      ip:
      port: 9999
      logpath: ./test/xxl-job
      logretentiondays: 30
      excludedpackage:


schedule:
  # 定时任务设置开关 开启(true)、关闭(false)
  enabled: true
  config:
    # 派单未接单通知
    - type: order_notification
      name: 订单未接单通知
      cron: 0 0/2 * * * ? # 默认值用于测试,实际由程序计算
      handler: orderNotAcceptedNotificationJobHandler # 对应JobHandler中的@XxlJob("orderNotAcceptedNotificationJobHandler")

创建定时任务id与业务id绑定关系表,并且创建对应实体和curd(这部分代码自行实现)

-- xxlJob映射表
CREATE TABLE `sys_xxl_job_mapping`  (
  `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT COMMENT '自增id',
  `entity_id` varchar(20) DEFAULT NULL COMMENT '业务id,通过此id找到JobId',
  `job_id` varchar(20) DEFAULT NULL COMMENT 'JobId',
  `task_type` varchar(100) DEFAULT NULL COMMENT '任务类型',
  `is_terminate` tinyint(1) unsigned DEFAULT '0' COMMENT '是否终止:1终止 0未终止',
  `created_user` varchar(32) DEFAULT NULL COMMENT '记录创建人',
  `updated_user` varchar(32) DEFAULT NULL COMMENT '最后修改人',
  `created_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT "创建时间",
  `updated_time` datetime DEFAULT NULL ON UPDATE CURRENT_TIMESTAMP COMMENT "更新时间",
  `is_enabled` tinyint(1) unsigned DEFAULT '1' COMMENT '是否启用:1启用 0禁用',
  `is_deleted` tinyint(1) unsigned DEFAULT '0' COMMENT '逻辑删除:1是 0否',
  PRIMARY KEY (`id`) USING BTREE,
  KEY `idx_entity_id_task_type` (`entity_id`,`task_type`) USING BTREE
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='xxlJob映射表';

创建订单未接单通知任务handler

@Slf4j
@Component
public class NotificationJobHandler {


    /**
     * 订单未接单通知任务handler
     * @throws Exception 异常
     */
    @XxlJob("orderNotAcceptedNotificationJobHandler") // 这个名称就是后续创建任务时指定的 executorHandler
    public void orderNotAcceptedNotificationJobHandler() throws Exception {
        // 1. 记录日志
        String param = XxlJobHelper.getJobParam();
        log.info("XXL-JOB, 接收到参数: {}, handler: orderNotAcceptedNotificationJobHandler", param);

        // TODO 业务待开发
    }
}

模拟业务中创建定时任务

@RestController
@RequestMapping("/job")
@Slf4j
@RequiredArgsConstructor
public class JobApiTestController {

    private final MyJobApiService myJobApiService; // API工具包中的MyJobApiService

    private final ISysXxlJobMappingService xxlJobMappingService;

    private final ScheduledConfigProperties configProperties;

    /**
     * 通过订单id添加并启动一个任务
     * @param orderNo
     * @return
     */
    @GetMapping("/start-test")
    public Response startJob(String orderNo){
      	// 获取application配置文件中的定时任务信息
        List<ScheduledConfig> scheduledConfig = configProperties.getConfig();
        ScheduledConfig config = scheduledConfig.stream().filter(it -> it.getType().equals("order_notification")).findFirst().orElse(null);

        Map<String, String> param = new HashMap<>();
        param.put("orderNo", orderNo);
        assert config != null;
      	// config.getCron()这个触发时间可以自行根据需求自定义,不一定要从配置中获取
        XxlJobInfoDTO dto = new XxlJobInfoDTO(config.getCron(), config.getName(), JSONObject.toJSONString(param), "admin", config.getHandler());
        Response response = myJobApiService.addJob(dto);
        log.info("添加任务结果: {}", JSONObject.toJSONString(response));

        // 保存业务映射 这一步主要为了后续删除任务时能通过业务id找到这个任务
        SysXxlJobMapping mapping = new SysXxlJobMapping();
        mapping.setEntityId(orderNo);
        mapping.setJobId(response.getData().toString());
        mapping.setTaskType(config.getType());
        xxlJobMappingService.saveJob(mapping);
        return response;
    }

    /**
     * 通过订单id删除一个任务
     * @param orderNo
     * @return
     */
    @GetMapping("/remove-test")
    public Response removeJob(String orderNo){

        SysXxlJobMapping jobMapping = xxlJobMappingService.getJob(orderNo, "order_notification");

        // 删除任务
        Response response = myJobApiService.removeJob(Integer.parseInt(jobMapping.getJobId()));
        log.info("删除任务结果: {}", JSONObject.toJSONString(response));

        // 更新业务映射
        xxlJobMappingService.terminateJob(orderNo, "order_notification");
        return response;
    }
}

最近看到一个很扎心的现象:企业越来越关注开发效率,而 AI 正在成为新的生产力工具。同样的需求,会使用 AI 的工程师往往能够更快完成设计、编码和测试工作。与其担心被 AI 替代,不如尽早学会驾驭 AI。最近我不仅在学习 Java 底层,还在学习一些人工智能的知识,发现了一个不错的 AI 学习网站,内容通俗易懂,比较适合程序员快速上手,感兴趣的话也可以看看:人工智能学习网

One more thing

当你开始爱自己,你就会睡得越来越早,也越来越喜欢锻炼,也不再纠结和焦虑,变得自信满满,去追求有意义的人和事,并为之燃烧自己的热情。你会发现,人生才真正开始。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐