最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

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
        }

        // 限定集群大?。?代表不限制)
        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如何動態(tài)導出excel表頭與數據

    jxls2.4.5如何動態(tài)導出excel表頭與數據

    這篇文章主要介紹了jxls2.4.5如何動態(tài)導出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實現二叉樹的創(chuàng)建及5種遍歷方法(總結)

    java實現二叉樹的創(chuàng)建及5種遍歷方法(總結)

    下面小編就為大家?guī)硪黄猨ava實現二叉樹的創(chuàng)建及5種遍歷方法(總結)。小編覺得挺不錯的,現在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-04-04
  • JWT整合Springboot的方法步驟

    JWT整合Springboot的方法步驟

    本文主要介紹了JWT整合Springboot的方法步驟,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • 解決Error:Java:無效的源發(fā)行版:14問題

    解決Error:Java:無效的源發(fā)行版:14問題

    在項目開發(fā)中,版本不一致常見問題,首先,應檢查本地JDK版本,使用命令java-version,其次,核對項目及模塊版本,若有不一致,通過修改pom.xml文件同步版本,重新下載依賴即可解決問題,這種方法簡單有效,適用于多種開發(fā)環(huán)境
    2024-10-10
  • 簡單了解SpringMVC與Struts2的區(qū)別

    簡單了解SpringMVC與Struts2的區(qū)別

    這篇文章主要介紹了簡單了解SpringMVC與Struts2的區(qū)別,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下
    2019-11-11
  • Servlet的5種方式實現表單提交(注冊小功能),后臺獲取表單數據實例

    Servlet的5種方式實現表單提交(注冊小功能),后臺獲取表單數據實例

    這篇文章主要介紹了Servlet的5種方式實現表單提交(注冊小功能),后臺獲取表單數據實例,非常具有實用價值,需要的朋友可以參考下
    2017-05-05

最新評論

买车| 江陵县| 扎赉特旗| 广丰县| 隆林| 昌邑市| 兴化市| 客服| 岐山县| 梧州市| 岑巩县| 万安县| 尉氏县| 巩义市| 湘潭市| 郁南县| 祁阳县| 绿春县| 双鸭山市| 宜宾市| 三门峡市| 高青县| 东兰县| 于都县| 东辽县| 抚远县| 扎赉特旗| 炉霍县| 岳阳县| 邻水| 宜黄县| 湾仔区| 张掖市| 宜君县| 绥江县| 兴城市| 耒阳市| 吉隆县| 嘉善县| 光山县| 九龙坡区|