PowerJob的DispatchStrategy方法工作流程源码解读

 更新时间:2024年01月12日 09:34:39   作者:codecraft  
这篇文章主要为大家介绍了PowerJob的DispatchStrategy方法工作流程源码解读,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪

本文主要研究一下PowerJob的DispatchStrategy

DispatchStrategy

tech/powerjob/common/enums/DispatchStrategy.java

@Getter
@AllArgsConstructor
public enum DispatchStrategy {
    HEALTH_FIRST(1),
    RANDOM(2);
    private final int v;
    public static DispatchStrategy of(Integer v) {
        if (v == null) {
            return HEALTH_FIRST;
        }
        for (DispatchStrategy ds : values()) {
            if (v.equals(ds.v)) {
                return ds;
            }
        }
        throw new IllegalArgumentException("unknown DispatchStrategy of " + v);
    }
}
DispatchStrategy定义了HEALTH_FIRST、RANDOM两个枚举值

getSuitableWorkers

tech/powerjob/server/remote/worker/WorkerClusterQueryService.java

public List<WorkerInfo> getSuitableWorkers(JobInfoDO jobInfo) {

        List<WorkerInfo> workers = Lists.newLinkedList(getWorkerInfosByAppId(jobInfo.getAppId()).values());

        workers.removeIf(workerInfo -> filterWorker(workerInfo, jobInfo));

        DispatchStrategy dispatchStrategy = DispatchStrategy.of(jobInfo.getDispatchStrategy());
        switch (dispatchStrategy) {
            case RANDOM:
                Collections.shuffle(workers);
                break;
            case HEALTH_FIRST:
                workers.sort((o1, o2) -> o2.getSystemMetrics().calculateScore() - o1.getSystemMetrics().calculateScore());
                break;
            default:
                // do nothing
        }

        // 限定集群大小(0代表不限制)
        if (!workers.isEmpty() && jobInfo.getMaxWorkerCount() > 0 && workers.size() > jobInfo.getMaxWorkerCount()) {
            workers = workers.subList(0, jobInfo.getMaxWorkerCount());
        }
        return workers;
    }
WorkerClusterQueryService的getSuitableWorkers方法先通过getWorkerInfosByAppId获取指定appId的WorkerInfo,然后通过filterWorker进行一次过滤,最后根据dispatchStrategy来对workers进行排序,如果是RANDOM则通过Collections.shuffle(workers)随机化,如果是HEALTH_FIRST则根据systemMetrics的calculateScore结果进行排序,如果有限定maxWorkerCount则对workers进行subList,没有则返回排序后的workers

getWorkerInfosByAppId

private Map<String, WorkerInfo> getWorkerInfosByAppId(Long appId) {
        ClusterStatusHolder clusterStatusHolder = getAppId2ClusterStatus().get(appId);
        if (clusterStatusHolder == null) {
            log.warn("[WorkerManagerService] can't find any worker for app(appId={}) yet.", appId);
            return Collections.emptyMap();
        }
        return clusterStatusHolder.getAllWorkers();
    }

    public Map<Long, ClusterStatusHolder> getAppId2ClusterStatus() {
        return WorkerClusterManagerService.getAppId2ClusterStatus();
    }
getWorkerInfosByAppId通过WorkerClusterManagerService.getAppId2ClusterStatus()获取ClusterStatusHolder,在返回ClusterStatusHolder的getAllWorkers

filterWorker

private boolean filterWorker(WorkerInfo workerInfo, JobInfoDO jobInfo) {
        for (WorkerFilter filter : workerFilters) {
            if (filter.filter(workerInfo, jobInfo)) {
                return true;
            }
        }
        return false;
    }
filterWorker方法则是遍历workerFilters直接filter

calculateScore

tech/powerjob/common/model/SystemMetrics.java

public int calculateScore() {
        if (score > 0) {
            return score;
        }
        // Memory is vital to TaskTracker, so we set the multiplier factor as 2.
        double memScore = (jvmMaxMemory - jvmUsedMemory) * 2;
        // Calculate the remaining load of CPU. Multiplier is set as 1.
        double cpuScore = cpuProcessors - cpuLoad;
        // Windows can not fetch CPU load, set cpuScore as 1.
        if (cpuScore > cpuProcessors) {
            cpuScore = 1;
        }
        score = (int) (memScore + cpuScore);
        return score;
    }
SystemMetrics的calculateScore方法则是基于memScore与cpuScore来计算

WorkerFilter

public interface WorkerFilter {

    /**
     *
     * @param workerInfo worker info, maybe you need to use your customized info in SystemMetrics#extra
     * @param jobInfoDO job info
     * @return true will remove the worker in process list
     */
    boolean filter(WorkerInfo workerInfo, JobInfoDO jobInfoDO);
}
WorkerFilter定义了filter接口用于过滤worker,它有3个实现类,分别是DesignatedWorkerFilter、DisconnectedWorkerFilter、SystemMetricsWorkerFilter

DesignatedWorkerFilter

tech/powerjob/server/extension/defaultimpl/workerfilter/DesignatedWorkerFilter.java

@Slf4j
@Component
public class DesignatedWorkerFilter implements WorkerFilter {
    @Override
    public boolean filter(WorkerInfo workerInfo, JobInfoDO jobInfo) {
        String designatedWorkers = jobInfo.getDesignatedWorkers();
        // no worker is specified, no filter of any
        if (StringUtils.isEmpty(designatedWorkers)) {
            return false;
        }
        Set<String> designatedWorkersSet = Sets.newHashSet(SJ.COMMA_SPLITTER.splitToList(designatedWorkers));
        for (String tagOrAddress : designatedWorkersSet) {
            if (tagOrAddress.equals(workerInfo.getTag()) || tagOrAddress.equals(workerInfo.getAddress())) {
                return false;
            }
        }
        return true;
    }
}
DesignatedWorkerFilter的filter方法遍历jobInfo的designatedWorkers信息,判断workerInfo的tag或者address是否匹配

DisconnectedWorkerFilter

tech/powerjob/server/extension/defaultimpl/workerfilter/DisconnectedWorkerFilter.java

@Slf4j
@Component
public class DisconnectedWorkerFilter implements WorkerFilter {

    @Override
    public boolean filter(WorkerInfo workerInfo, JobInfoDO jobInfo) {
        boolean timeout = workerInfo.timeout();
        if (timeout) {
            log.info("[Job-{}] filter worker[{}] due to timeout(lastActiveTime={})", jobInfo.getId(), workerInfo.getAddress(), workerInfo.getLastActiveTime());
        }
        return timeout;
    }
}
DisconnectedWorkerFilter的filter方法则通过WorkerInfo的timeout方法来判断,它主要是判断当前时间与lastActiveTime的时间差是否大于WORKER_TIMEOUT_MS(60s)

SystemMetricsWorkerFilter

tech/powerjob/server/extension/defaultimpl/workerfilter/SystemMetricsWorkerFilter.java

@Slf4j
@Component
public class SystemMetricsWorkerFilter implements WorkerFilter {

    @Override
    public boolean filter(WorkerInfo workerInfo, JobInfoDO jobInfo) {
        SystemMetrics metrics = workerInfo.getSystemMetrics();
        boolean filter = !metrics.available(jobInfo.getMinCpuCores(), jobInfo.getMinMemorySpace(), jobInfo.getMinDiskSpace());
        if (filter) {
            log.info("[Job-{}] filter worker[{}] because the {} do not meet the requirements", jobInfo.getId(), workerInfo.getAddress(), workerInfo.getSystemMetrics());
        }
        return filter;
    }
}
SystemMetricsWorkerFilter的filter方法则根据workerInfo的SystemMetrics判断可用cpu核数、内存、磁盘空间是否大于阈值

小结

DispatchStrategy定义了HEALTH_FIRST、RANDOM两个枚举值;WorkerClusterQueryService的getSuitableWorkers方法先通过getWorkerInfosByAppId获取指定appId的WorkerInfo,然后通过filterWorker进行一次过滤,最后根据dispatchStrategy来对workers进行排序,如果是RANDOM则通过Collections.shuffle(workers)随机化,如果是HEALTH_FIRST则根据systemMetrics的calculateScore结果进行排序,如果有限定maxWorkerCount则对workers进行subList,没有则返回排序后的workers。

以上就是PowerJob的DispatchStrategy方法工作流程源码解读的详细内容,更多关于PowerJob DispatchStrategy的资料请关注脚本之家其它相关文章!

相关文章

  • jxls2.4.5如何动态导出excel表头与数据

    jxls2.4.5如何动态导出excel表头与数据

    这篇文章主要介绍了jxls2.4.5如何动态导出excel表头与数据问题,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教
    2024-08-08
  • java中归并排序和Master公式详解

    java中归并排序和Master公式详解

    大家好,本篇文章主要讲的是java中归并排序和Master公式详解,感兴趣的同学赶快来看一看吧,对你有帮助的话记得收藏一下,方便下次浏览
    2022-01-01
  • 通过实例解析Spring argNames属性

    通过实例解析Spring argNames属性

    这篇文章主要介绍了通过实例解析Spring argNames属性,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下
    2020-11-11
  • 解决mybatis-plus新增数据自增ID变无序问题

    解决mybatis-plus新增数据自增ID变无序问题

    这篇文章主要介绍了解决mybatis-plus新增数据自增ID变无序问题,具有很好的参考价值,希望对大家有所帮助。
    2023-07-07
  • maven项目无法解析插件的解决方案

    maven项目无法解析插件的解决方案

    这篇文章主要介绍了maven项目无法解析插件的解决,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教
    2025-06-06
  • java实现二叉树的创建及5种遍历方法(总结)

    java实现二叉树的创建及5种遍历方法(总结)

    下面小编就为大家带来一篇java实现二叉树的创建及5种遍历方法(总结)。小编觉得挺不错的,现在就分享给大家,也给大家做个参考。一起跟随小编过来看看吧
    2017-04-04
  • JWT整合Springboot的方法步骤

    JWT整合Springboot的方法步骤

    本文主要介绍了JWT整合Springboot的方法步骤,文中通过示例代码介绍的非常详细,具有一定的参考价值,感兴趣的小伙伴们可以参考一下
    2021-11-11
  • 解决Error:Java:无效的源发行版:14问题

    解决Error:Java:无效的源发行版:14问题

    在项目开发中,版本不一致常见问题,首先,应检查本地JDK版本,使用命令java-version,其次,核对项目及模块版本,若有不一致,通过修改pom.xml文件同步版本,重新下载依赖即可解决问题,这种方法简单有效,适用于多种开发环境
    2024-10-10
  • 简单了解SpringMVC与Struts2的区别

    简单了解SpringMVC与Struts2的区别

    这篇文章主要介绍了简单了解SpringMVC与Struts2的区别,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下
    2019-11-11
  • Servlet的5种方式实现表单提交(注册小功能),后台获取表单数据实例

    Servlet的5种方式实现表单提交(注册小功能),后台获取表单数据实例

    这篇文章主要介绍了Servlet的5种方式实现表单提交(注册小功能),后台获取表单数据实例,非常具有实用价值,需要的朋友可以参考下
    2017-05-05

最新评论