diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java index 8285850b8..8f1eeadae 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java @@ -13,6 +13,7 @@ import org.springframework.stereotype.Controller; import org.springframework.web.bind.annotation.*; import ru.clearing.classes.statics.data.scheduler.Launcher; import ru.clearing.classes.statics.data.user.User; +import ru.clearing.platform.dictionary.AbstractDictionary; import ru.spcex.clearing.backendapi.controller.queue.AbstractQueueController; import ru.spcex.clearing.backendapi.controller.request.cud.schedule.LauncherNew; import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse; @@ -23,7 +24,6 @@ import ru.spcex.clearing.backendapi.security.KeycloakUtils; import ru.spcex.clearing.backendapi.service.IOperator; import ru.spcex.clearing.backendapi.service.IStateLoader; import ru.spcex.clearing.imdg.IMDGDistributedNames; -import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.EnumMessage; @@ -38,12 +38,14 @@ import java.util.concurrent.ExecutionException; public class LauncherController extends AbstractQueueController { private final IStateLoader stateLoader; private final Imdg userImdg; + private final Imdg taskDictionary; @Autowired public LauncherController(IStateLoader stateLoader, IOperator operator, ImdgProvider imdgProvider) { super(operator); this.stateLoader = stateLoader; this.userImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_User, User.class); + this.taskDictionary = imdgProvider.getImdg(IMDGDistributedNames.Map_TaskDictionary, AbstractDictionary.class); } @ApiOperation(value = "get all Launchers.") @@ -63,6 +65,10 @@ public class LauncherController extends AbstractQueueController { @ResponseBody public CudResponse add(@ApiParam(value = "Код задания из taskDictionary", required = true, example = "ABLK") @PathVariable("task-code") String dictionaryName) throws ExecutionException, InterruptedException { + AbstractDictionary taskEnum = taskDictionary.getSingleObjectByFieldValues(Map.of("code", dictionaryName)); + if (taskEnum == null) { + throw new NotFound404Exception("task dictionary element with code '" + dictionaryName + "'"); + } LauncherNew launcherCommand = new LauncherNew(); Authentication authentication = SecurityContextHolder.getContext().getAuthentication(); String username = KeycloakUtils.getUserNameFromAuthentication(authentication); @@ -72,7 +78,9 @@ public class LauncherController extends AbstractQueueController { } launcherCommand.setTask(dictionaryName); launcherCommand.setUserId(user.getId()); - return processRequest(Consts.LAUNCHER_NEW, launcherCommand); + //топики ограничиваются наличием в taskDictionary + //подписываются на разные топики в разных модулях, см. ru.spcex.platform.enumeration.Task#topic + return processRequest("launcher-" + dictionaryName, launcherCommand); } @Autowired diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java index 887fda87d..a2c9487a8 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java @@ -32,7 +32,7 @@ public class ClearingService { } public void paymentUpdateBySdf04(Long sdf04GroupId) { - log.info("adding sdf03/11 from STLD payments task to queue"); + log.info("updating payment.transactionStatus by sdf04 task added to queue"); executor.execute(() -> paymentUpdater.updatePayments(sdf04GroupId)); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java similarity index 65% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index 08db5dde1..bd4bfef11 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -5,12 +5,14 @@ import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Service; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.enumeration.Task; @Service -public class Sdf04Receiver extends QueueConsumer implements InitializingBean { +public class EventsReceiver extends QueueConsumer implements InitializingBean { private final ClearingService clearingService; - public Sdf04Receiver(Consumer kafkaQueue, ClearingService clearingService) { + public EventsReceiver(Consumer kafkaQueue, ClearingService clearingService) { super(kafkaQueue); this.clearingService = clearingService; } @@ -23,6 +25,9 @@ public class Sdf04Receiver extends QueueConsumer implements InitializingBean { clearingService.paymentUpdateBySdf04(requestPayload.getGroupId()); }) .forDestination(Consts.SDF04_PROCESS, callbacks::put); + callback(LauncherCommandRequest.class) + .setConsumer(event -> clearingService.sdfCreate()) + .forDestination(Task.createOrder.topic(), callbacks::put); init(); } } diff --git a/clearing-parent/dbf-importer/pom.xml b/clearing-parent/dbf-importer/pom.xml index d3d50fe59..0460db432 100644 --- a/clearing-parent/dbf-importer/pom.xml +++ b/clearing-parent/dbf-importer/pom.xml @@ -66,6 +66,10 @@ ru.spcex.platform platform-messaging + + ru.spcex.platform + platform-enum + diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java index bf533336b..fc448445b 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java @@ -1,11 +1,15 @@ package ru.spcex.clearing.dbf.importer.config; +import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; @@ -35,4 +39,11 @@ public class KafkaConfig { .build(); }; } + + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(ImportDBFServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/settings/ImportDBFServiceSettings.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/settings/ImportDBFServiceSettings.java index 8c4bd1fe2..af11e25ee 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/settings/ImportDBFServiceSettings.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/settings/ImportDBFServiceSettings.java @@ -3,6 +3,7 @@ package ru.spcex.clearing.dbf.importer.config.settings; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.PropertySource; import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @@ -12,6 +13,7 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; public class ImportDBFServiceSettings { private HazelcastClientParams hazelcast; private KafkaProducerSettings kafkaProducer; + private KafkaConsumerSettings kafkaConsumer; private Common common; private Store store; private Cron cron; @@ -32,6 +34,14 @@ public class ImportDBFServiceSettings { this.kafkaProducer = kafkaProducer; } + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + public Common getCommon() { return common; } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index 384056064..a4413c8a9 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -1,5 +1,7 @@ package ru.spcex.clearing.dbf.importer.services; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; @@ -16,6 +18,7 @@ import java.util.Map; @Service("dbfImporterService") @EnableScheduling public class DBFImporterService { + private final Logger log = LoggerFactory.getLogger(getClass()); private final FileChecker fileChecker; private final ThreadPoolTaskExecutor executorService; private final Processor processor; @@ -30,7 +33,12 @@ public class DBFImporterService { @Scheduled(cron = "${import-dbf-service.scheduler.check-src-dir-cron}") public void run() { - Map> newFiles = fileChecker.checkNewFiles(); + run(null); + } + + public void run(ETable specificTable) { + log.info("performing import {}", specificTable != null ? specificTable.name() : ""); + Map> newFiles = fileChecker.checkNewFiles(specificTable); for (Map.Entry> newFilesEntry : newFiles.entrySet()) { ETable currTable = newFilesEntry.getKey(); List fileList = newFilesEntry.getValue(); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java index 785523c3e..2ee390a94 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/FileChecker.java @@ -16,6 +16,10 @@ public class FileChecker { } public Map> checkNewFiles() { + return checkNewFiles(null); + } + + public Map> checkNewFiles(ETable specificTable) { Map> newFiles = new EnumMap<>(ETable.class); String srcDir = settings.getStore().getSrcDir(); @@ -27,11 +31,10 @@ public class FileChecker { ETable currTable = ETable.getTableForFilename(dbfFile.getName()); if (currTable == null) continue; - + if (specificTable != null && !specificTable.equals(currTable)) continue; List currList = newFiles.computeIfAbsent(currTable, list -> new LinkedList<>()); currList.add(dbfFile); } - return newFiles; } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java new file mode 100644 index 000000000..eec3f4f88 --- /dev/null +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java @@ -0,0 +1,33 @@ +package ru.spcex.clearing.dbf.importer.services; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.enumeration.Task; + +@Service +public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final DBFImporterService importer; + + @Autowired + public LauncherCommandReceiver(Consumer kafkaQueue, + DBFImporterService importer) { + super(kafkaQueue); + this.importer = importer; + } + + @Override + public void afterPropertiesSet() { + callback(LauncherCommandRequest.class) + .setConsumer(action -> importer.run(ETable.DF_04)) + .forDestination(Task.createOrderConfirm.topic(), callbacks::put); + init(); + } +} diff --git a/clearing-parent/dbf-importer/src/main/resources/application.properties b/clearing-parent/dbf-importer/src/main/resources/application.properties index 4c0438d76..c1c2a3fe0 100644 --- a/clearing-parent/dbf-importer/src/main/resources/application.properties +++ b/clearing-parent/dbf-importer/src/main/resources/application.properties @@ -22,4 +22,12 @@ import-dbf-service.kafka-producer.acks=all import-dbf-service.kafka-producer.retries=0 import-dbf-service.kafka-producer.batch-size=16384 import-dbf-service.kafka-producer.linger-ms=1 -import-dbf-service.kafka-producer.buffer-memory=33554432 \ No newline at end of file +import-dbf-service.kafka-producer.buffer-memory=33554432 + +import-dbf-service.kafka-consumer.bootstrap-servers=localhost:9092 +import-dbf-service.kafka-consumer.group-id=dev-group-clearing-service +import-dbf-service.kafka-consumer.enable-auto-commit=false +import-dbf-service.kafka-consumer.session-timeout-ms=30000 +import-dbf-service.kafka-consumer.auto-offset-reset=latest +import-dbf-service.kafka-consumer.linger-ms=1 +import-dbf-service.kafka-consumer.buffer-memory=33554432 \ No newline at end of file diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java new file mode 100644 index 000000000..762af41b4 --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/Task.java @@ -0,0 +1,22 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum Task implements IEnumKey { + createOrder("CORD"), createOrderConfirm("CORC"); + + private final String key; + + Task(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } + + public String topic() { + return "launcher-" + getKey(); + } +}