From 394bcec2b38c8a25fdfc3937504b70653dfa2a7e Mon Sep 17 00:00:00 2001 From: AKurakin Date: Mon, 16 Jan 2023 12:19:47 +0300 Subject: [PATCH] =?UTF-8?q?scheduler-service=20=D0=BE=D1=82=D0=BF=D1=80?= =?UTF-8?q?=D0=B0=D0=B2=D0=BA=D0=B0=20=D0=BA=D0=BE=D0=BC=D0=B0=D0=BD=D0=B4?= =?UTF-8?q?=20=D0=B2=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D0=B8=20=D0=BF?= =?UTF-8?q?=D0=BE=20=D1=80=D0=B0=D1=81=D0=BF=D0=B8=D1=81=D0=B0=D0=BD=D0=B8?= =?UTF-8?q?=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduler/service/LauncherSender.java | 68 +++++++++++++++++++ .../scheduler/service/TaskManager.java | 11 ++- 2 files changed, 77 insertions(+), 2 deletions(-) create mode 100644 clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java new file mode 100644 index 000000000..bbf98c9d5 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java @@ -0,0 +1,68 @@ +package ru.spcex.clearing.scheduler.service; + + +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.platform.enumeration.Task; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; + +@Service +public class LauncherSender { + Logger log = LoggerFactory.getLogger(getClass()); + private final Producer kafka; + private final ImdgId idGenerator; + + @Autowired + public LauncherSender(Producer kafka, ImdgProvider imdgProvider) { + this.kafka = kafka; + this.idGenerator = imdgProvider.getImdgIdGenerator(); + } + + protected BaseRequest makeCmdRequest(Task byTask, Long userId) { + Object toRequest; + if (Task.createReport_GREP.equals(byTask)) { // см. ru.spcex.clearing.backendapi.controller.request.cud.schedule.LauncherNew#toRequest + ReportWithPeriodRequest reportCommand = new ReportWithPeriodRequest(); + reportCommand.setReportId(byTask.getKey()); + toRequest = reportCommand; + } else { + LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest(); + taskRunnerCommandRequest.setTaskName(byTask.getKey()); + taskRunnerCommandRequest.setUserId(userId); + toRequest = taskRunnerCommandRequest; + } + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.nextId()); + request.setActionType(ActionType.NEW); + request.setRequestPayload(toRequest);// iAction.toRequest() + return request; + } + + public void sendCommandToQueue(Task toTaskQueue, Long userId) { + BaseRequest request = makeCmdRequest(toTaskQueue, userId); + String destination = toTaskQueue.topic(); // "launcher-" + getKey() + log.debug("Send command {} to {}", toTaskQueue, destination); + Future send = kafka.send(new ProducerRecord<>(destination, request)); + try { + send.get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException("Command " + toTaskQueue + " not send", e); + } catch (ExecutionException e) { + throw new RuntimeException("Command " + toTaskQueue + " not send", e); + } + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java index ee5b53975..b0a8708d5 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java @@ -48,12 +48,15 @@ public class TaskManager implements EntryAddedListener, private Imdg launcherMap; private Imdg plannerAllTodayMap; private ConcurrentHashMap scheduledJobs; + private LauncherSender launcherSender; @Autowired TaskManager(TaskScheduler taskScheduler, - ImdgProvider imdgProvider) { + ImdgProvider imdgProvider, + LauncherSender launcherSender) { this.taskScheduler = taskScheduler; this.imdgProvider = imdgProvider; + this.launcherSender = launcherSender; } private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) { @@ -216,8 +219,10 @@ public class TaskManager implements EntryAddedListener, // --- Реализация выполнения задач --- protected void doJob(PlannerAllToday task) { - if (getEnumByKey(Task.class, task.getTask()) != null) { + Task taskE = getEnumByKey(Task.class, task.getTask()); + if (taskE != null) { log.debug("LauncherCommandRequest received"); + // 1. cохранить команду Instant created = Instant.now(); Launcher launcher = new Launcher(); launcher.setTask(task.getTask()); @@ -225,6 +230,8 @@ public class TaskManager implements EntryAddedListener, launcher.setCreated(created); launcher.setUpdated(created); launcherMap.insert(launcher); + // 2. отправить сообщение + launcherSender.sendCommandToQueue(taskE, task.getParentId()); log.debug("successfully processed, new id {}", launcher.getId()); } }