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 позже, и я надеюсь, что смогу понять это позже .