xieb
2023-05-12 ae7ffbc02d0ceb98c6d230bc99a1ba3316627b95
调度执行器
8 files modified
3 files added
429 ■■■■■ changed files
skjcmanager/skjcmanager-ops/skjcmanager-xxljob-admin/src/main/resources/application.yml 19 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/pom.xml 5 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/entity/PatrolResponsiblePerson.java 18 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/mapper/AttResManagePersonMapper.java 3 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/mapper/AttResManagePersonMapper.xml 6 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/service/IAttResManagePersonService.java 8 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/service/impl/AttResManagePersonServiceImpl.java 7 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/xxljob/config/XxlJobConfig.java 74 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/xxljob/jobhandler/SampleXxlJob.java 250 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/resources/application-dev.yml 20 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/resources/application-prod.yml 19 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-ops/skjcmanager-xxljob-admin/src/main/resources/application.yml
@@ -19,7 +19,7 @@
    templateLoaderPath: classpath:/templates/
  mail:
    host: smtp.qq.com
    password: xxx
    password: fmgmlbflwgkkbjgi
    port: 25
    properties:
      mail:
@@ -30,7 +30,7 @@
          starttls:
            enable: true
            required: true
    username: xxx@qq.com
    username: 524270321@qq.com
  mvc:
    servlet:
      load-on-startup: 0
@@ -38,13 +38,18 @@
  resources:
    static-locations: classpath:/static/
#对外暴露端口
management:
  health:
    mail:
  endpoints:
    web:
      exposure:
        include: "*"
        exclude: health
  endpoint:
    env:
      enabled: false
  server:
    servlet:
      context-path: /actuator
    health:
      show-details: always
mybatis:
  mapper-locations: classpath:/mybatis-mapper/*Mapper.xml
skjcmanager/skjcmanager-service/skjcmanager-sm/pom.xml
@@ -52,6 +52,11 @@
            <artifactId>skjcmanager-user-api</artifactId>
            <version>3.0.1.RELEASE</version>
        </dependency>
        <!--Job-->
        <dependency>
            <groupId>com.xuxueli</groupId>
            <artifactId>xxl-job-core</artifactId>
        </dependency>
    </dependencies>
    <build>
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/entity/PatrolResponsiblePerson.java
New file
@@ -0,0 +1,18 @@
package cn.gistack.sm.sjztmd.entity;
import com.baomidou.mybatisplus.annotation.TableField;
import io.swagger.annotations.ApiModel;
import lombok.Data;
@Data
@ApiModel(value = "巡查责任人对象", description = "巡查责任人对象")
public class PatrolResponsiblePerson {
    @TableField("\"res_guid\"")
    private String resGuid;
    @TableField("\"name\"")
    private String name;
    @TableField("\"user_id\"")
    private String userId;
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/mapper/AttResManagePersonMapper.java
@@ -1,6 +1,7 @@
package cn.gistack.sm.sjztmd.mapper;
import cn.gistack.sm.sjztmd.entity.AttResManagePerson;
import cn.gistack.sm.sjztmd.entity.PatrolResponsiblePerson;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import java.util.List;
@@ -16,4 +17,6 @@
    List<AttResManagePerson> attResManagePersonListByResId(List<Long> resId);
    List<PatrolResponsiblePerson> getPatrolResponsiblePersonAll();
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/mapper/AttResManagePersonMapper.xml
@@ -23,4 +23,10 @@
        </foreach>
    </select>
    <select id="getPatrolResponsiblePersonAll" resultType="cn.gistack.sm.sjztmd.entity.PatrolResponsiblePerson">
        select DISTINCT a."res_guid",c."name",b.id as user_id
        from "SJZT_MD"."att_res_manage_person" a,"BLADE_USER" b ,"SJZT_MD"."att_res_base" c
        WHERE a."user_phone" = b.PHONE and a."res_guid" = c."guid" and a."type" = 5
    </select>
</mapper>
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/service/IAttResManagePersonService.java
@@ -1,6 +1,7 @@
package cn.gistack.sm.sjztmd.service;
import cn.gistack.sm.sjztmd.entity.AttResManagePerson;
import cn.gistack.sm.sjztmd.entity.PatrolResponsiblePerson;
import com.baomidou.mybatisplus.extension.service.IService;
import java.util.List;
@@ -20,4 +21,11 @@
     * @return
     */
    List<AttResManagePerson> attResManagePersonListByResIds(String ResId);
    /**
     * 获取巡查责任人 水库 用户id
     * @return
     */
    List<PatrolResponsiblePerson> getPatrolResponsiblePersonAll();
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/service/impl/AttResManagePersonServiceImpl.java
@@ -1,6 +1,7 @@
package cn.gistack.sm.sjztmd.service.impl;
import cn.gistack.sm.sjztmd.entity.AttResManagePerson;
import cn.gistack.sm.sjztmd.entity.PatrolResponsiblePerson;
import cn.gistack.sm.sjztmd.mapper.AttResManagePersonMapper;
import cn.gistack.sm.sjztmd.service.IAttResManagePersonService;
import com.baomidou.dynamic.datasource.annotation.DS;
@@ -24,4 +25,10 @@
    public List<AttResManagePerson> attResManagePersonListByResIds(String resId) {
        return baseMapper.attResManagePersonListByResId(Func.toLongList(resId));
    }
    @Override
    @DS("master")
    public List<PatrolResponsiblePerson> getPatrolResponsiblePersonAll() {
        return baseMapper.getPatrolResponsiblePersonAll();
    }
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/xxljob/config/XxlJobConfig.java
New file
@@ -0,0 +1,74 @@
package cn.gistack.sm.xxljob.config;
import com.xxl.job.core.executor.impl.XxlJobSpringExecutor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
 * xxl-job config
 *
 * @author xuxueli 2017-04-28
 */
@Configuration(proxyBeanMethods = false)
public class XxlJobConfig {
    private final Logger logger = LoggerFactory.getLogger(XxlJobConfig.class);
    @Value("${xxl.job.admin.addresses}")
    private String adminAddresses;
    @Value("${xxl.job.executor.appname}")
    private String appName;
    @Value("${xxl.job.executor.ip}")
    private String ip;
    @Value("${xxl.job.executor.port}")
    private int port;
    @Value("${xxl.job.accessToken}")
    private String accessToken;
    @Value("${xxl.job.executor.logpath}")
    private String logPath;
    @Value("${xxl.job.executor.logretentiondays}")
    private int logRetentionDays;
    @Bean
    public XxlJobSpringExecutor xxlJobExecutor() {
        logger.info(">>>>>>>>>>> xxl-job config init.");
        XxlJobSpringExecutor xxlJobSpringExecutor = new XxlJobSpringExecutor();
        xxlJobSpringExecutor.setAdminAddresses(adminAddresses);
        xxlJobSpringExecutor.setAppName(appName);
        xxlJobSpringExecutor.setIp(ip);
        xxlJobSpringExecutor.setPort(port);
        xxlJobSpringExecutor.setAccessToken(accessToken);
        xxlJobSpringExecutor.setLogPath(logPath);
        xxlJobSpringExecutor.setLogRetentionDays(logRetentionDays);
        return xxlJobSpringExecutor;
    }
    /**
     * 针对多网卡、容器内部署等情况,可借助 "spring-cloud-commons" 提供的 "InetUtils" 组件灵活定制注册IP;
     *
     *      1、引入依赖:
     *          <dependency>
     *             <groupId>org.springframework.cloud</groupId>
     *             <artifactId>spring-cloud-commons</artifactId>
     *             <version>${version}</version>
     *         </dependency>
     *
     *      2、配置文件,或者容器启动变量
     *          spring.cloud.inetutils.preferred-networks: 'xxx.xxx.xxx.'
     *
     *      3、获取IP
     *          String ip_ = inetUtils.findFirstNonLoopbackHostInfo().getIpAddress();
     */
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/xxljob/jobhandler/SampleXxlJob.java
New file
@@ -0,0 +1,250 @@
package cn.gistack.sm.xxljob.jobhandler;
import cn.gistack.sm.patrol.entity.PatrolTask;
import cn.gistack.sm.patrol.service.IPatrolTaskService;
import cn.gistack.sm.sjztmd.entity.PatrolResponsiblePerson;
import cn.gistack.sm.sjztmd.service.IAttResManagePersonService;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.xxl.job.core.biz.model.ReturnT;
import com.xxl.job.core.handler.IJobHandler;
import com.xxl.job.core.handler.annotation.XxlJob;
import com.xxl.job.core.log.XxlJobLogger;
import com.xxl.job.core.util.ShardingUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springblade.core.tool.utils.DateUtil;
import org.springframework.stereotype.Component;
import java.io.BufferedInputStream;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.net.HttpURLConnection;
import java.net.URL;
import java.util.Date;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
/**
 * XxlJob开发示例(Bean模式)
 * <p>
 * 开发步骤:
 * 1、在Spring Bean实例中,开发Job方法,方式格式要求为 "public ReturnT<String> execute(String param)"
 * 2、为Job方法添加注解 "@XxlJob(value="自定义jobhandler名称", init = "JobHandler初始化方法", destroy = "JobHandler销毁方法")",注解value值对应的是调度中心新建任务的JobHandler属性的值。
 * 3、执行日志:需要通过 "XxlJobLogger.log" 打印执行日志;
 *
 * @author xuxueli
 */
@Component
public class SampleXxlJob {
    private static final Logger logger = LoggerFactory.getLogger(SampleXxlJob.class);
    private final IPatrolTaskService patrolTaskService;
    private final IAttResManagePersonService attResManagePersonService;
    public SampleXxlJob(IPatrolTaskService patrolTaskService, IAttResManagePersonService attResManagePersonService) {
        this.patrolTaskService = patrolTaskService;
        this.attResManagePersonService = attResManagePersonService;
    }
    /**
     * 1、简单任务示例(Bean模式)
     */
    @XxlJob("demoJobHandler")
    public ReturnT<String> demoJobHandler(String param) throws Exception {
        XxlJobLogger.log("<span style='color:red'>XXL-JOB-BLADE_SM, Hello World.</span>");
        for (int i = 0; i < 5; i++) {
            XxlJobLogger.log("beat at:" + i);
            TimeUnit.SECONDS.sleep(2);
        }
        return ReturnT.SUCCESS;
    }
    @XxlJob("createTaskJobHandler")
    public ReturnT<String> createTaskJobHandler(String param) throws Exception {
        XxlJobLogger.log("开始自动创建任务...");
        JSONObject jsonParam = JSON.parseObject(param);
        String processDefinitionId = jsonParam.getString("processDefinitionId");
        final String taskType = jsonParam.getString("taskType");
        final String title = jsonParam.getString("title");
        final String content = jsonParam.getString("content");
        List<PatrolResponsiblePerson> patrolResponsiblePersonList =  attResManagePersonService.getPatrolResponsiblePersonAll();
        AtomicInteger failureCount = new AtomicInteger();
        int successCount = 0;
        patrolResponsiblePersonList.forEach(item -> {
            XxlJobLogger.log("当前创建任务对象:" + JSON.toJSONString(item));
            try {
                PatrolTask patrolTask = new PatrolTask();
                patrolTask.setProcessDefinitionId(processDefinitionId);
                patrolTask.setToUserId(item.getUserId());
                patrolTask.setProjectId(item.getResGuid());
                patrolTask.setReservoirName(item.getName());
                patrolTask.setTitle(item.getName() + "-" + title);
                patrolTask.setContent(DateUtil.format(new Date(),"yyyyMMdd") +content);
                patrolTask.setTaskType(taskType);
                // 创建任务
                patrolTaskService.startProcess(patrolTask);
                TimeUnit.SECONDS.sleep(2);
            } catch (Exception e) {
                XxlJobLogger.log("创建任务失败对象:"  + JSON.toJSONString(item));
                failureCount.getAndIncrement();
            }
        });
        XxlJobLogger.log("结束自动创建任务...");
        XxlJobLogger.log("总数量:" + patrolResponsiblePersonList.size() + " ... 失败数量:" + failureCount);
        return ReturnT.SUCCESS;
    }
    /**
     * 2、分片广播任务
     */
    @XxlJob("shardingJobHandler")
    public ReturnT<String> shardingJobHandler(String param) throws Exception {
        // 分片参数
        ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo();
        XxlJobLogger.log("分片参数:当前分片序号 = {}, 总分片数 = {}", shardingVO.getIndex(), shardingVO.getTotal());
        // 业务逻辑
        for (int i = 0; i < shardingVO.getTotal(); i++) {
            if (i == shardingVO.getIndex()) {
                XxlJobLogger.log("第 {} 片, 命中分片开始处理", i);
            } else {
                XxlJobLogger.log("第 {} 片, 忽略", i);
            }
        }
        return ReturnT.SUCCESS;
    }
    /**
     * 3、命令行任务
     */
    @XxlJob("commandJobHandler")
    public ReturnT<String> commandJobHandler(String param) throws Exception {
        String command = param;
        int exitValue = -1;
        BufferedReader bufferedReader = null;
        try {
            // command process
            Process process = Runtime.getRuntime().exec(command);
            BufferedInputStream bufferedInputStream = new BufferedInputStream(process.getInputStream());
            bufferedReader = new BufferedReader(new InputStreamReader(bufferedInputStream));
            // command log
            String line;
            while ((line = bufferedReader.readLine()) != null) {
                XxlJobLogger.log(line);
            }
            // command exit
            process.waitFor();
            exitValue = process.exitValue();
        } catch (Exception e) {
            XxlJobLogger.log(e);
        } finally {
            if (bufferedReader != null) {
                bufferedReader.close();
            }
        }
        if (exitValue == 0) {
            return IJobHandler.SUCCESS;
        } else {
            return new ReturnT<String>(IJobHandler.FAIL.getCode(), "command exit value(" + exitValue + ") is failed");
        }
    }
    /**
     * 4、跨平台Http任务
     */
    @XxlJob("httpJobHandler")
    public ReturnT<String> httpJobHandler(String param) throws Exception {
        // request
        HttpURLConnection connection = null;
        BufferedReader bufferedReader = null;
        try {
            // connection
            URL realUrl = new URL(param);
            connection = (HttpURLConnection) realUrl.openConnection();
            // connection setting
            connection.setRequestMethod("GET");
            connection.setDoOutput(true);
            connection.setDoInput(true);
            connection.setUseCaches(false);
            connection.setReadTimeout(5 * 1000);
            connection.setConnectTimeout(3 * 1000);
            connection.setRequestProperty("connection", "Keep-Alive");
            connection.setRequestProperty("Content-Type", "application/json;charset=UTF-8");
            connection.setRequestProperty("Accept-Charset", "application/json;charset=UTF-8");
            // do connection
            connection.connect();
            //Map<String, List<String>> map = connection.getHeaderFields();
            // valid StatusCode
            int statusCode = connection.getResponseCode();
            if (statusCode != 200) {
                throw new RuntimeException("Http Request StatusCode(" + statusCode + ") Invalid.");
            }
            // result
            bufferedReader = new BufferedReader(new InputStreamReader(connection.getInputStream(), "UTF-8"));
            StringBuilder result = new StringBuilder();
            String line;
            while ((line = bufferedReader.readLine()) != null) {
                result.append(line);
            }
            String responseMsg = result.toString();
            XxlJobLogger.log(responseMsg);
            return ReturnT.SUCCESS;
        } catch (Exception e) {
            XxlJobLogger.log(e);
            return ReturnT.FAIL;
        } finally {
            try {
                if (bufferedReader != null) {
                    bufferedReader.close();
                }
                if (connection != null) {
                    connection.disconnect();
                }
            } catch (Exception e2) {
                XxlJobLogger.log(e2);
            }
        }
    }
    /**
     * 5、生命周期任务示例:任务初始化与销毁时,支持自定义相关逻辑;
     */
    @XxlJob(value = "demoJobHandler2", init = "init", destroy = "destroy")
    public ReturnT<String> demoJobHandler2(String param) throws Exception {
        XxlJobLogger.log("XXL-JOB, Hello World.");
        return ReturnT.SUCCESS;
    }
    public void init() {
        logger.info("init");
    }
    public void destroy() {
        logger.info("destory");
    }
}
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/resources/application-dev.yml
@@ -32,4 +32,22 @@
          url: ${blade.datasource.dev.ztznwh.url}
          username: ${blade.datasource.dev.ztznwh.username}
          password: ${blade.datasource.dev.ztznwh.password}
  main:
    allow-circular-references: true
xxl:
  job:
    accessToken: ''
    admin:
      ### 调度中心部署根地址 [选填]:如调度中心集群部署存在多个地址则用逗号分隔。执行器将会使用该地址进行"执行器心跳注册"和"任务结果回调";为空则关闭自动注册;
      addresses: http://127.0.0.1:7009/xxl-job-admin
    executor:
      ### 执行器AppName [选填]:执行器心跳注册分组依据;为空则关闭自动注册
      appname: blade-xxljob
      ### 执行器IP [选填]:默认为空表示自动获取IP,多网卡时可手动设置指定IP,该IP不会绑定Host仅作为通讯实用;地址信息用于 "执行器注册" 和 "调度中心请求并触发任务";
      ip: 127.0.0.1
      ### 执行器运行日志文件存储磁盘路径 [选填] :需要对该路径拥有读写权限;为空则使用默认路径;
      logpath: ../data/applogs/xxl-job/jobhandler
      ### 执行器日志文件保存天数 [选填] : 过期日志自动清理, 限制值大于等于3时生效; 否则, 如-1, 关闭自动清理功能;
      logretentiondays: -1
      ### 执行器端口号 [选填]:小于等于0则自动获取;默认端口为9999,单机部署多个执行器时,注意要配置不同执行器端口;
      port: 7018
skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/resources/application-prod.yml
@@ -32,3 +32,22 @@
          url: ${blade.datasource.prod.ztznwh.url}
          username: ${blade.datasource.prod.ztznwh.username}
          password: ${blade.datasource.prod.ztznwh.password}
  main:
    allow-circular-references: true
xxl:
  job:
    accessToken: ''
    admin:
      ### 调度中心部署根地址 [选填]:如调度中心集群部署存在多个地址则用逗号分隔。执行器将会使用该地址进行"执行器心跳注册"和"任务结果回调";为空则关闭自动注册;
      addresses: http://127.0.0.1:7009/xxl-job-admin
    executor:
      ### 执行器AppName [选填]:执行器心跳注册分组依据;为空则关闭自动注册
      appname: blade-xxljob
      ### 执行器IP [选填]:默认为空表示自动获取IP,多网卡时可手动设置指定IP,该IP不会绑定Host仅作为通讯实用;地址信息用于 "执行器注册" 和 "调度中心请求并触发任务";
      ip: 127.0.0.1
      ### 执行器运行日志文件存储磁盘路径 [选填] :需要对该路径拥有读写权限;为空则使用默认路径;
      logpath: ../data/applogs/xxl-job/jobhandler
      ### 执行器日志文件保存天数 [选填] : 过期日志自动清理, 限制值大于等于3时生效; 否则, 如-1, 关闭自动清理功能;
      logretentiondays: -1
      ### 执行器端口号 [选填]:小于等于0则自动获取;默认端口为9999,单机部署多个执行器时,注意要配置不同执行器端口;
      port: 7018