From 1442ab6c50b3170dac24596f25317f478361c0ca Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 20 Apr 2026 17:14:59 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-961 --- .../csv/importer/config/KafkaConfig.java | 39 +++++++++++++++++++ .../SOrdersScheduler.java | 20 +++++++++- .../tro/processor/TroProcessorTest.java | 2 +- 3 files changed, 58 insertions(+), 3 deletions(-) create mode 100644 clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/config/KafkaConfig.java rename clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/{component => service}/SOrdersScheduler.java (79%) diff --git a/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/config/KafkaConfig.java b/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/config/KafkaConfig.java new file mode 100644 index 000000000..6ba7d935e --- /dev/null +++ b/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/config/KafkaConfig.java @@ -0,0 +1,39 @@ +package ru.spcex.clearing.csv.importer.config; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.csv.importer.config.settings.CsvImporterSettings; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Configuration +public class KafkaConfig { + @Bean + public ProducerFactory pf(CsvImporterSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean("kafkaSenderWithoutRequestInfo") + public KafkaSender kafkaSenderWithoutRequestInfo(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .saveRequestInfo(false) + .build(); + } +} \ No newline at end of file diff --git a/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/component/SOrdersScheduler.java b/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/service/SOrdersScheduler.java similarity index 79% rename from clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/component/SOrdersScheduler.java rename to clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/service/SOrdersScheduler.java index 3e92282c4..908982efc 100644 --- a/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/component/SOrdersScheduler.java +++ b/clearing-parent/csv-importer/src/main/java/ru/spcex/clearing/csv/importer/service/SOrdersScheduler.java @@ -1,4 +1,4 @@ -package ru.spcex.clearing.csv.importer.component; +package ru.spcex.clearing.csv.importer.service; import java.nio.file.Files; import java.nio.file.Path; @@ -10,12 +10,16 @@ import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import ru.clearing.classes.statics.data.misc.SOrders; +import ru.spcex.clearing.csv.importer.component.FileService; import ru.spcex.clearing.csv.importer.component.tro.data.SOrder; import ru.spcex.clearing.csv.importer.component.tro.parser.TroParserType; import ru.spcex.clearing.csv.importer.component.tro.processor.TroProcessor; import ru.spcex.clearing.csv.importer.component.tro.processor.TroProcessorProvider; import ru.spcex.clearing.csv.importer.config.settings.CsvImporterSettings; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.order.OrderCurrencyStatusUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.log.ExceptionUtils; @@ -29,13 +33,15 @@ public class SOrdersScheduler { private final FileService fileService; private final TroProcessor processor; private final Imdg sOrderImdg; + private final KafkaSender kafkaSender; @Autowired - public SOrdersScheduler(CsvImporterSettings settings, FileService fileService, TroProcessorProvider provider, ImdgProvider imdgProvider) { + public SOrdersScheduler(CsvImporterSettings settings, FileService fileService, TroProcessorProvider provider, ImdgProvider imdgProvider, KafkaSender kafkaSender) { this.importDir = Path.of(settings.getSrcDir()); this.fileService = fileService; this.processor = provider.createProcessor(TroParserType.SOrder, SOrder.class); this.sOrderImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SOrders, SOrders.class); + this.kafkaSender = kafkaSender; } @Scheduled(fixedDelayString = "${csv-importer.poll-period:5000}") @@ -74,10 +80,20 @@ public class SOrdersScheduler { sOrderImdg.insert(order); } } + sendToDefaultManagement(sOrders); fileService.moveSuccess(path); } catch (Throwable e) { log.error("error for {}: {}", path, ExceptionUtils.getStackTrace(e)); fileService.moveError(path); } } + + private void sendToDefaultManagement(List sOrders) { + String destination = Consts.ORDER_CURRENCY_STATUS_UPDATE_IMPORTER; + OrderCurrencyStatusUpdateRequest req = new OrderCurrencyStatusUpdateRequest(); + sOrders.forEach(order -> req.addId(order.getTransId())); + log.info("sending request to topic {}...", destination); + kafkaSender.sendRequestToQueue(destination, req); + log.info("kafka request sent"); + } } diff --git a/clearing-parent/csv-importer/src/test/java/ru/spcex/clearing/csv/importer/component/tro/processor/TroProcessorTest.java b/clearing-parent/csv-importer/src/test/java/ru/spcex/clearing/csv/importer/component/tro/processor/TroProcessorTest.java index 259b28dc1..3865617b2 100644 --- a/clearing-parent/csv-importer/src/test/java/ru/spcex/clearing/csv/importer/component/tro/processor/TroProcessorTest.java +++ b/clearing-parent/csv-importer/src/test/java/ru/spcex/clearing/csv/importer/component/tro/processor/TroProcessorTest.java @@ -27,6 +27,6 @@ class TroProcessorTest { ); TroProcessor processor = provider.createProcessor(TroParserType.SOrder, SOrder.class); List sOrders = processor.readFile(path); - Assertions.assertTrue(sOrders.size() > 0, "read SOrder from tro file"); + Assertions.assertFalse(sOrders.isEmpty(), "read SOrder's from tro file"); } } \ No newline at end of file