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