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

PowerJob的WorkerHealthReporter工作流程源碼解讀

 更新時(shí)間:2023年12月25日 09:59:23   作者:codecraft  
這篇文章主要為大家介紹了PowerJob的WorkerHealthReporter工作流程源碼解讀,

本文主要研究一下PowerJob的WorkerHealthReporter

WorkerHealthReporter

tech/powerjob/worker/background/WorkerHealthReporter.java

@Slf4j
@RequiredArgsConstructor
public class WorkerHealthReporter implements Runnable {
    private final WorkerRuntime workerRuntime;
    @Override
    public void run() {
        // 沒(méi)有可用Server,無(wú)法上報(bào)
        String currentServer = workerRuntime.getServerDiscoveryService().getCurrentServerAddress();
        if (StringUtils.isEmpty(currentServer)) {
            log.warn("[WorkerHealthReporter] no available server,fail to report health info!");
            return;
        }
        SystemMetrics systemMetrics;
        if (workerRuntime.getWorkerConfig().getSystemMetricsCollector() == null) {
            systemMetrics = SystemInfoUtils.getSystemMetrics();
        } else {
            systemMetrics = workerRuntime.getWorkerConfig().getSystemMetricsCollector().collect();
        }
        WorkerHeartbeat heartbeat = new WorkerHeartbeat();
        heartbeat.setSystemMetrics(systemMetrics);
        heartbeat.setWorkerAddress(workerRuntime.getWorkerAddress());
        heartbeat.setAppName(workerRuntime.getWorkerConfig().getAppName());
        heartbeat.setAppId(workerRuntime.getAppId());
        heartbeat.setHeartbeatTime(System.currentTimeMillis());
        heartbeat.setVersion(PowerJobWorkerVersion.getVersion());
        heartbeat.setProtocol(workerRuntime.getWorkerConfig().getProtocol().name());
        heartbeat.setClient("KingPenguin");
        heartbeat.setTag(workerRuntime.getWorkerConfig().getTag());
        // 上報(bào) Tracker 數(shù)量
        heartbeat.setLightTaskTrackerNum(LightTaskTrackerManager.currentTaskTrackerSize());
        heartbeat.setHeavyTaskTrackerNum(HeavyTaskTrackerManager.currentTaskTrackerSize());
        // 是否超載
        if (workerRuntime.getWorkerConfig().getMaxLightweightTaskNum() <= LightTaskTrackerManager.currentTaskTrackerSize() || workerRuntime.getWorkerConfig().getMaxHeavyweightTaskNum() <= HeavyTaskTrackerManager.currentTaskTrackerSize()){
            heartbeat.setOverload(true);
        }
        // 獲取當(dāng)前加載的容器列表
        heartbeat.setContainerInfos(OmsContainerFactory.getDeployedContainerInfos());
        // 發(fā)送請(qǐng)求
        if (StringUtils.isEmpty(currentServer)) {
            return;
        }
        // log
        log.info("[WorkerHealthReporter] report health status,appId:{},appName:{},isOverload:{},maxLightweightTaskNum:{},currentLightweightTaskNum:{},maxHeavyweightTaskNum:{},currentHeavyweightTaskNum:{}" ,
                heartbeat.getAppId(),
                heartbeat.getAppName(),
                heartbeat.isOverload(),
                workerRuntime.getWorkerConfig().getMaxLightweightTaskNum(),
                heartbeat.getLightTaskTrackerNum(),
                workerRuntime.getWorkerConfig().getMaxHeavyweightTaskNum(),
                heartbeat.getHeavyTaskTrackerNum()
        );
        TransportUtils.reportWorkerHeartbeat(heartbeat, currentServer, workerRuntime.getTransporter());
    }
}
WorkerHealthReporter實(shí)現(xiàn)了Runnable接口,其run方法先獲取currentServer,再獲取systemMetrics,接著構(gòu)建WorkerHeartbeat,最后通過(guò)TransportUtils.reportWorkerHeartbeat上報(bào)

reportWorkerHeartbeat

tech/powerjob/worker/common/utils/TransportUtils.java

public static void reportWorkerHeartbeat(WorkerHeartbeat req, String address, Transporter transporter) {
        final URL url = easyBuildUrl(ServerType.SERVER, S4W_PATH, S4W_HANDLER_WORKER_HEARTBEAT, address);
        transporter.tell(url, req);
    }
    public static URL easyBuildUrl(ServerType serverType, String rootPath, String handlerPath, String address) {
        HandlerLocation handlerLocation = new HandlerLocation()
                .setRootPath(rootPath)
                .setMethodPath(handlerPath);
        return new URL()
                .setServerType(serverType)
                .setAddress(Address.fromIpv4(address))
                .setLocation(handlerLocation);
    }
reportWorkerHeartbeat通過(guò)transporter.tell發(fā)送請(qǐng)求,其rootPath為server,其handlerPath為workerHeartbeat

processWorkerHeartbeat

tech/powerjob/server/core/handler/AbWorkerRequestHandler.java

@Handler(path = S4W_HANDLER_WORKER_HEARTBEAT, processType = ProcessType.NO_BLOCKING)
    public void processWorkerHeartbeat(WorkerHeartbeat heartbeat) {
        long startMs = System.currentTimeMillis();
        WorkerHeartbeatEvent event = new WorkerHeartbeatEvent()
                .setAppName(heartbeat.getAppName())
                .setAppId(heartbeat.getAppId())
                .setVersion(heartbeat.getVersion())
                .setProtocol(heartbeat.getProtocol())
                .setTag(heartbeat.getTag())
                .setWorkerAddress(heartbeat.getWorkerAddress())
                .setDelayMs(startMs - heartbeat.getHeartbeatTime())
                .setScore(heartbeat.getSystemMetrics().getScore());
        processWorkerHeartbeat0(heartbeat, event);
        monitorService.monitor(event);
    }
processWorkerHeartbeat方法將heartbeat轉(zhuǎn)換為WorkerHeartbeatEvent,然后執(zhí)行processWorkerHeartbeat0及monitorService.monitor(event)

processWorkerHeartbeat0

tech/powerjob/server/core/handler/WorkerRequestHandlerImpl.java

protected void processWorkerHeartbeat0(WorkerHeartbeat heartbeat, WorkerHeartbeatEvent event) {
        WorkerClusterManagerService.updateStatus(heartbeat);
    }
processWorkerHeartbeat0通過(guò)WorkerClusterManagerService.updateStatus(heartbeat)來(lái)更新?tīng)顟B(tài)

WorkerClusterManagerService.updateStatus

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

public static void updateStatus(WorkerHeartbeat heartbeat) {
        Long appId = heartbeat.getAppId();
        String appName = heartbeat.getAppName();
        ClusterStatusHolder clusterStatusHolder = APP_ID_2_CLUSTER_STATUS.computeIfAbsent(appId, ignore -> new ClusterStatusHolder(appName));
        clusterStatusHolder.updateStatus(heartbeat);
    }
updateStatus先獲取appId對(duì)應(yīng)的clusterStatusHolder,然后更新status

ClusterStatusHolder.updateStatus

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

public void updateStatus(WorkerHeartbeat heartbeat) {

        String workerAddress = heartbeat.getWorkerAddress();
        long heartbeatTime = heartbeat.getHeartbeatTime();

        WorkerInfo workerInfo = address2WorkerInfo.computeIfAbsent(workerAddress, ignore -> {
            WorkerInfo wf = new WorkerInfo();
            wf.refresh(heartbeat);
            return wf;
        });
        long oldTime = workerInfo.getLastActiveTime();
        if (heartbeatTime < oldTime) {
            log.warn("[ClusterStatusHolder-{}] receive the expired heartbeat from {}, serverTime: {}, heartTime: {}", appName, heartbeat.getWorkerAddress(), System.currentTimeMillis(), heartbeat.getHeartbeatTime());
            return;
        }

        workerInfo.refresh(heartbeat);

        List<DeployedContainerInfo> containerInfos = heartbeat.getContainerInfos();
        if (!CollectionUtils.isEmpty(containerInfos)) {
            containerInfos.forEach(containerInfo -> {
                Map<String, DeployedContainerInfo> infos = containerId2Infos.computeIfAbsent(containerInfo.getContainerId(), ignore -> Maps.newConcurrentMap());
                infos.put(workerAddress, containerInfo);
            });
        }
    }
ClusterStatusHolder的updateStatus方法先獲取workerInfo,判斷其heartbeatTime是否小于lastActiveTime,是則返回,否則執(zhí)行workerInfo.refresh(heartbeat),最后更新一下heartbeat.getContainerInfos()

refresh

tech/powerjob/server/common/module/WorkerInfo.java

public void refresh(WorkerHeartbeat workerHeartbeat) {
        address = workerHeartbeat.getWorkerAddress();
        lastActiveTime = workerHeartbeat.getHeartbeatTime();
        protocol = workerHeartbeat.getProtocol();
        client = workerHeartbeat.getClient();
        tag = workerHeartbeat.getTag();
        systemMetrics = workerHeartbeat.getSystemMetrics();
        containerInfos = workerHeartbeat.getContainerInfos();

        lightTaskTrackerNum = workerHeartbeat.getLightTaskTrackerNum();
        heavyTaskTrackerNum = workerHeartbeat.getHeavyTaskTrackerNum();

        if (workerHeartbeat.isOverload()) {
            overloading = true;
            lastOverloadTime = workerHeartbeat.getHeartbeatTime();
            log.warn("[WorkerInfo] worker {} is overload!", getAddress());
        } else {
            overloading = false;
        }
    }
WorkerInfo的refresh方法根據(jù)workerHeartbeat更新lastActiveTime及overloading等信息

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實(shí)現(xiàn)了WorkerFilter接口,其filter方法返回workerInfo.timeout()

timeout

tech/powerjob/server/common/module/WorkerInfo.java

private static final long WORKER_TIMEOUT_MS = 60000;

    public boolean timeout() {
        long timeout = System.currentTimeMillis() - lastActiveTime;
        return timeout > WORKER_TIMEOUT_MS;
    }
timeout方法判斷當(dāng)前時(shí)間與lastActiveTime的時(shí)間差,之后與默認(rèn)的WORKER_TIMEOUT_MS(60s)對(duì)比

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;
    }

    private boolean filterWorker(WorkerInfo workerInfo, JobInfoDO jobInfo) {
        for (WorkerFilter filter : workerFilters) {
            if (filter.filter(workerInfo, jobInfo)) {
                return true;
            }
        }
        return false;
    }
getSuitableWorkers方法會(huì)remove掉filterWorker(workerInfo, jobInfo)為true的worker

小結(jié)

PowerJob的WorkerHealthReporter實(shí)現(xiàn)了Runnable接口,其run方法先獲取currentServer,再獲取systemMetrics,接著構(gòu)建WorkerHeartbeat,最后通過(guò)TransportUtils.reportWorkerHeartbeat上報(bào);

reportWorkerHeartbeat通過(guò)transporter.tell發(fā)送請(qǐng)求,其rootPath為server,其handlerPath為workerHeartbeat;

服務(wù)端通過(guò)WorkerClusterManagerService.updateStatus(heartbeat)來(lái)更新?tīng)顟B(tài),主要是執(zhí)行WorkerInfo的refresh方法,它根據(jù)workerHeartbeat更新lastActiveTime及overloading等信息;

而DisconnectedWorkerFilter實(shí)現(xiàn)了WorkerFilter接口,其filter方法返回workerInfo.timeout(),它會(huì)將心跳超時(shí)的worker給排除掉。

以上就是PowerJob的WorkerHealthReporter工作流程源碼解讀的詳細(xì)內(nèi)容,更多關(guān)于PowerJob WorkerHealthReporter工作流程的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • java普通項(xiàng)目讀取不到resources目錄下資源文件的解決辦法

    java普通項(xiàng)目讀取不到resources目錄下資源文件的解決辦法

    這篇文章主要給大家介紹了關(guān)于java普通項(xiàng)目讀取不到resources目錄下資源文件的解決辦法,Web項(xiàng)目中應(yīng)該經(jīng)常有這樣的需求,在maven項(xiàng)目的resources目錄下放一些文件,比如一些配置文件,資源文件等,需要的朋友可以參考下
    2023-09-09
  • Elasticsearch算分優(yōu)化方案之rescore_query示例詳解

    Elasticsearch算分優(yōu)化方案之rescore_query示例詳解

    這篇文章主要為大家介紹了Elasticsearch算分優(yōu)化方案之rescore_query示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-08-08
  • Java行為型設(shè)計(jì)模式之外觀設(shè)計(jì)模式詳解

    Java行為型設(shè)計(jì)模式之外觀設(shè)計(jì)模式詳解

    外觀模式為多個(gè)復(fù)雜的子系統(tǒng),提供了一個(gè)一致的界面,使得調(diào)用端只和這個(gè)接口發(fā)生調(diào)用,而無(wú)須關(guān)系這個(gè)子系統(tǒng)內(nèi)部的細(xì)節(jié)。本文將通過(guò)示例詳細(xì)為大家講解一下外觀模式,需要的可以參考一下
    2022-11-11
  • MyBatis Plus大數(shù)據(jù)量查詢慢原因分析及解決

    MyBatis Plus大數(shù)據(jù)量查詢慢原因分析及解決

    大數(shù)據(jù)量查詢慢常因全表掃描、分頁(yè)不當(dāng)、索引缺失、內(nèi)存占用高及ORM開(kāi)銷(xiāo),優(yōu)化措施包括分頁(yè)查詢、流式讀取、SQL優(yōu)化、批處理、多數(shù)據(jù)源、結(jié)果集二次處理及配置調(diào)優(yōu)
    2025-09-09
  • Java Scala數(shù)據(jù)類(lèi)型與變量常量及類(lèi)和對(duì)象超詳細(xì)講解

    Java Scala數(shù)據(jù)類(lèi)型與變量常量及類(lèi)和對(duì)象超詳細(xì)講解

    本文內(nèi)容主要分為3節(jié),依次講解:Scala的數(shù)據(jù)類(lèi)型有哪些? 變量常量如何使用? 類(lèi)和對(duì)象如何理解? 受限于博主的大腦容量,大概是無(wú)法做到事無(wú)巨細(xì)的,不過(guò)其實(shí)也沒(méi)必要那么"細(xì)",抓住主要脈絡(luò),加上大量的練習(xí),融會(huì)貫通只不過(guò)是時(shí)間的問(wèn)題
    2022-12-12
  • 在java中使用SPI創(chuàng)建可擴(kuò)展的應(yīng)用程序操作

    在java中使用SPI創(chuàng)建可擴(kuò)展的應(yīng)用程序操作

    這篇文章主要介紹了在java中使用SPI創(chuàng)建可擴(kuò)展的應(yīng)用程序操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-09-09
  • Java中Stream流對(duì)多個(gè)字段進(jìn)行排序的方法

    Java中Stream流對(duì)多個(gè)字段進(jìn)行排序的方法

    我們?cè)谔幚頂?shù)據(jù)的時(shí)候經(jīng)常會(huì)需要進(jìn)行排序后再返回給前端調(diào)用,比如按照時(shí)間升序排序,前端展示數(shù)據(jù)就是按時(shí)間先后進(jìn)行排序,下面這篇文章主要給大家介紹了關(guān)于Java中Stream流對(duì)多個(gè)字段進(jìn)行排序的相關(guān)資料,需要的朋友可以參考下
    2023-10-10
  • SpringBatch跳過(guò)異常和限制方式

    SpringBatch跳過(guò)異常和限制方式

    這篇文章主要介紹了SpringBatch跳過(guò)異常和限制方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Java Red5服務(wù)器實(shí)現(xiàn)流媒體視頻播放

    Java Red5服務(wù)器實(shí)現(xiàn)流媒體視頻播放

    這篇文章主要介紹了Java Red5服務(wù)器實(shí)現(xiàn)流媒體視頻播放,對(duì)視頻播放感興趣的同學(xué),可以參考下
    2021-04-04
  • 用SpringMVC編寫(xiě)一個(gè)HelloWorld的詳細(xì)過(guò)程

    用SpringMVC編寫(xiě)一個(gè)HelloWorld的詳細(xì)過(guò)程

    SpringMVC是Spring的一個(gè)后續(xù)產(chǎn)品,是Spring的一個(gè)子項(xiàng)目<BR>SpringMVC?是?Spring?為表述層開(kāi)發(fā)提供的一整套完備的解決方案,本文我們將用SpringMVC編寫(xiě)一個(gè)HelloWorld,文中有詳細(xì)的編寫(xiě)過(guò)程,需要的朋友可以參考下
    2023-08-08

最新評(píng)論

康定县| 平泉县| 梁河县| 泰和县| 迭部县| 马山县| 格尔木市| 封开县| 建德市| 洛浦县| 太和县| 新泰市| 巴林右旗| 合肥市| 襄汾县| 清涧县| 梁山县| 呼和浩特市| 兴仁县| 龙岩市| 石景山区| 松溪县| 尤溪县| 凤凰县| 新巴尔虎左旗| 通江县| 巴中市| 连山| 高雄县| 清新县| 榆社县| 平利县| 海安县| 洛宁县| 大姚县| 大洼县| 镇远县| 房产| 鄄城县| 日土县| 汉源县|