Интерпретация исходного кода PowerJob 1: интерпретация связи между сервером и исполнителем

задняя часть

1. Краткое введение в PowerJob

PowerJob(原OhMyScheduler)是全新一代分布式任务调度与计算框架。

Вышеприведенное введение взято из официальной документации PowerJob Насколько я понимаю, PowerJob — это промежуточная служба, которая управляет задачами, рассчитанными по времени, отложенными задачами и т. д. в нескольких других приложениях. На официальном веб-сайте сообщается, что процессоры Map/Reduce также можно использовать для распределенной обработки задач, чего я пока не знаю.

2. Связь PowerJob

PowerJob предоставляет независимо развернутую серверную часть для унифицированного управления задачами в разных клиентах (называемых Worker в PowerJob). Это включает в себя связь между серверной и рабочей сторонами. В распределенной системе взаимодействие между сервером и клиентом основывается на концепции обнаружения службы.

В общей модели обнаружения сервисов есть 3 роли:

  • Поставщик услуг: Сторона, предоставляющая услуги, часто предоставляет API для использования другими службами.
  • Потребитель услуг: Сторона, использующая услуги поставщика услуг, то есть сторона, использующая API поставщика услуг.
  • Центр регистрации: поставщик услуг регистрируется в центре регистрации, и регистрационная информация обычно включает имя службы и адрес службы. Затем потребитель службы обращается к реестру и использует имя службы для получения адреса службы, чтобы потребитель службы мог связаться с поставщиком службы.

Использование службы обнаружения может обеспечить высокую доступность распределенных систем, а в PowerJob нет службы реестра, так как же она реализует вышеуказанные функции?

2.1 Воркер получает адрес сервера

В общей распределенной системе поставщик услуг сначала регистрируется в центре регистрации, чтобы потребитель услуг мог получать информацию поставщика услуг через центр регистрации и пользоваться услугами поставщика услуг.

Это похоже на PowerJob, и аналогичные операции выполняются при инициализации Worker.

public void init() throws Exception {
    //初始化Akka
    ...
    // 服务发现
    currentServer = ServerDiscoveryService.discovery();
    if (StringUtils.isEmpty(currentServer) && !config.isEnableTestMode()) {
        throw new RuntimeException("can't find any available server, this worker has been quarantined.");
    }
    log.info("[OhMyWorker] discovery server succeed, current server is {}.", currentServer);
    ...
}

Выше приведен исходный код инициализации Worker. Вы можете видеть, что при инициализации Worker он также выполняет операцию обнаружения службы. Продолжайте просматривать исходный код ниже.

public static String discovery() {

        if (IP2ADDRESS.isEmpty()) {
            OhMyWorker.getConfig().getServerAddress().forEach(x -> IP2ADDRESS.put(x.split(":")[0], x));
        }

        String result = null;

        // 先对当前机器发起请求
        //1、判断之前是否已确定一个可用的Server,如果之前已经确定了一个Server,就再验证这个Server是否还是可用的
        String currentServer = OhMyWorker.getCurrentServer();
        if (!StringUtils.isEmpty(currentServer)) {
            String ip = currentServer.split(":")[0];
            // 直接请求当前Server的HTTP服务,可以少一次网络开销,减轻Server负担
            String firstServerAddress = IP2ADDRESS.get(ip);
            if (firstServerAddress != null) {
                result = acquire(firstServerAddress);
            }
        }

        //2、如果之前没有可用的Server,依次判断Server数组中的Server,查找出可用的Server地址
        for (String httpServerAddress : OhMyWorker.getConfig().getServerAddress()) {
            if (StringUtils.isEmpty(result)) {
                result = acquire(httpServerAddress);
            }else {
                break;
            }
        }

        //3、如果没有找到可用的Server,说明当前Worker与外界失联,进行错误处理
        if (StringUtils.isEmpty(result)) {
            log.warn("[OmsServerDiscovery] can't find any available server, this worker has been quarantined.");

            // 在 Server 高可用的前提下,连续失败多次,说明该节点与外界失联,Server已经将秒级任务转移到其他Worker,需要杀死本地的任务
            //错误处理
            return null;
        }else {
            // 重置失败次数
            FAILED_COUNT = 0;
            log.debug("[OmsServerDiscovery] current server is {}.", result);
            return result;
        }
    }

OhMyWorker.getConfig().getServerAddress()Возвращается адрес сервера, настроенный при инициализации рабочего процесса.

Как видно из приведенного выше исходного кода, полное обнаружение службы в основном состоит из 3 шагов:

  • 1. Если ранее был найден доступный сервер, определите, доступен ли еще сервер, и верните сервер, если он доступен.
  • 2. Если ранее доступный сервер не был найден (может быть, это первое обнаружение службы), по очереди оценивайте сервер в массиве серверов (инициализированный список служб) и возвращайте первый доступный сервер.
  • 3. Доступный сервер не найден, что указывает на то, что текущий рабочий процесс отключен, и выполняется обработка ошибок.

2.1.1 Как определить, доступен ли сервер

использовать в приведенном выше кодеprivate static String acquire(String httpServerAddress)Чтобы определить, доступен ли сервер, давайте посмотрим, проверен ли он.

private static String acquire(String httpServerAddress) {
    String result = null;
    String url = String.format(DISCOVERY_URL, httpServerAddress, OhMyWorker.getAppId(), OhMyWorker.getCurrentServer());
    try {
        result = CommonUtils.executeWithRetry0(() -> HttpUtils.get(url));
    }catch (Exception ignore) {
    }
    if (!StringUtils.isEmpty(result)) {
        try {
            ResultDTO resultDTO = JsonUtils.parseObject(result, ResultDTO.class);
            if (resultDTO.isSuccess()) {
                return resultDTO.getData().toString();
            }
        }catch (Exception ignore) {
        }
    }
    return null;
}

Приведенный выше код очень прост, он инициирует HTTP-запрос к адресу сервера, который необходимо проверить, а параметр запроса имеет AppId, соответствующий рабочему. Судя по возвращаемым данным этого запроса, если данные, возвращаемые этим сервером, соответствуют ожиданиям, это означает, что этот сервер доступен.

resultDTO.getData()Значением является адрес сервера и номер порта Akka на сервере.

Далее давайте посмотрим, как серверная сторона обрабатывает этот запрос.

public String getServer(Long appId) {

        Set<String> downServerCache = Sets.newHashSet();

        for (int i = 0; i < RETRY_TIMES; i++) {

            // 无锁获取当前数据库中的Server
            //1、从数据库中获取这个应用的信息,信息中会有这个应用绑定的Server地址
            Optional<AppInfoDO> appInfoOpt = appInfoRepository.findById(appId);
            if (!appInfoOpt.isPresent()) {
                throw new OmsException(appId + " is not registered!");
            }
            String appName = appInfoOpt.get().getAppName();
            String originServer = appInfoOpt.get().getCurrentServer();
            //判断当前应用绑定的Server是否存活,如果存活就返回这个Server信息
            if (isActive(originServer, downServerCache)) {
                return originServer;
            }

            //2、如果Server不可用,就尝试将应用绑定的Server换成本机
            // 无可用Server,重新进行Server选举,需要加锁
            String lockName = String.format(SERVER_ELECT_LOCK, appId);
            boolean lockStatus = lockService.lock(lockName, 30000);
            if (!lockStatus) {
                try {
                    Thread.sleep(500);
                }catch (Exception ignore) {
                }
                continue;
            }
            try {

                // 可能上一台机器已经完成了Server选举,需要再次判断
                AppInfoDO appInfo = appInfoRepository.findById(appId).orElseThrow(() -> new RuntimeException("impossible, unless we just lost our database."));
                if (isActive(appInfo.getCurrentServer(), downServerCache)) {
                    return appInfo.getCurrentServer();
                }

                // 篡位,本机作为Server
                appInfo.setCurrentServer(OhMyServer.getActorSystemAddress());
                appInfo.setGmtModified(new Date());

                appInfoRepository.saveAndFlush(appInfo);
                log.info("[ServerSelectService] this server({}) become the new server for app(appId={}).", appInfo.getCurrentServer(), appId);
                return appInfo.getCurrentServer();
            }catch (Exception e) {
                log.warn("[ServerSelectService] write new server to db failed for app {}.", appName);
            }finally {
                lockService.unlock(lockName);
            }
        }
        throw new RuntimeException("server elect failed for app " + appId);
    }

Операция обработки в основном включает следующие этапы:

  • 1. Используйте параметр запроса appId, чтобы получить информацию о приложении, проверить, доступен ли сервер в информации о приложении, и напрямую вернуть информацию о сервере, если она доступна.
  • 2. Если сервер недоступен, выполните выбор сервера. По сути, под управлением распределенных блокировок каждый сервер ставит себя сервером привязки данного приложения. Наконец, вернитесь на выбранный сервер.

Вернитесь к заголовку этого раздела, если установлено, что сервер доступен.

private boolean isActive(String serverAddress, Set<String> downServerCache) {
        ...
        ActorSelection serverActor = OhMyServer.getFriendActor(serverAddress);
        try {
            CompletionStage<Object> askCS = Patterns.ask(serverActor, ping, Duration.ofMillis(PING_TIMEOUT_MS));
            AskResponse response = (AskResponse) askCS.toCompletableFuture().get(PING_TIMEOUT_MS, TimeUnit.MILLISECONDS);
            downServerCache.remove(serverAddress);
            return response.isSuccess();
        }catch (Exception e) {
            log.warn("[ServerSelectService] server({}) was down.", serverAddress);
        }
        downServerCache.add(serverAddress);
        return false;
    }

На самом деле это очень просто, просто отправьте сигнал Ping сервису Akka на этом сервере, и если вы сможете получить корректный ответ, значит, этот сервер доступен.

До сих пор мы разбирали информацию об адресе сервера, связанную соответствующим приложением с рабочим. Таким образом, Worker может сообщать серверу о некоторых своих условиях, например, об использовании пульса для поддержания соединения с сервером.

3. Сервер получает рабочий адрес

Выше показано, что Worker получает соответствующий адрес Сервера, а следующее показывает, как Сервер получает управляемую им информацию Worker.

После того, как рабочий процесс получит адрес сервера, он будет использовать Akka для периодической отправки на сервер информации о пульсе. Информация пульса включает в себя адрес локального компьютера (IP: порт), который является портом локальной службы Akka; appName и appId, соответствующие локальному рабочему процессу; использование ресурсов локальной системы и т. д.

Давайте посмотрим, что делает сервер после получения информации о пульсе.

/**
* 更新状态
* @param heartbeat Worker的心跳包
*/
public static void updateStatus(WorkerHeartbeat heartbeat) {
    Long appId = heartbeat.getAppId();
    String appName = heartbeat.getAppName();
    ClusterStatusHolder clusterStatusHolder = appId2ClusterStatus.computeIfAbsent(appId, ignore -> new ClusterStatusHolder(appName));
    clusterStatusHolder.updateStatus(heartbeat);
}

В основном для обновления состояния Worker, где классClusterStatusHolderЭто класс, который управляет одним и тем же рабочим кластером и внутренне поддерживает карту для хранения состояния каждой машины в кластере.

рядом сclusterStatusHolder.updateStatus(heartbeat)Посмотри внутри.

public void updateStatus(WorkerHeartbeat heartbeat) {

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

        Long oldTime = address2ActiveTime.getOrDefault(workerAddress, -1L);
        if (heartbeatTime < oldTime) {
            log.warn("[ClusterStatusHolder-{}] receive the expired heartbeat from {}, serverTime: {}, heartTime: {}", appName, heartbeat.getWorkerAddress(), System.currentTimeMillis(), heartbeat.getHeartbeatTime());
            return;
        }

        address2ActiveTime.put(workerAddress, heartbeatTime);
        address2Metrics.put(workerAddress, heartbeat.getSystemMetrics());

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

Он основан на периоде сердцебиения и использовании ресурсов рабочей машины.

Хотя этот код относительно прост, адрес рабочего процесса получается через сервер пульса, и сервер может использовать этот адрес для отправки информации о выполнении задачи рабочему процессу.

4. Резюме

Исходя из вышеприведенного анализа, Worker является поставщиком услуг в распределенной системе, Server — потребителем услуг, а потребляемый контент — это запланированные задачи Worker, управляемые сервером.

Мы также упомянули реестр в начале, что такое реестр в PowerJob? Это также сервер.Рабочий процесс регистрирует имя приложения и адрес приложения на сервере, отправляя информацию о сердцебиении.

Упомянутая выше концепция заключается в том, что Сервер связан с Работником, то есть управление задачами Работника выполняется Сервером, независимо от того, сколько экземпляров Сервера имеется в этой распределенной системе. Преимуществом этого является группировка и изоляция задач. В чем польза от этой группировки и изоляции задач? Сейчас я немного запутался. Я опубликую блог об анализе планирования задач PowerJob позже, и я надеюсь, что смогу понять это позже .