From 15cad890c9b9dec70a84ebb3a45578374db3eacc Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 28 Sep 2022 18:49:27 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-36 --- .../dbf/importer/config/KafkaConfig.java | 33 +++++++++---------- .../dbf/importer/logic/stages/ImportToDB.java | 7 ++-- 2 files changed, 20 insertions(+), 20 deletions(-) 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 e3d6ae042..86358d7c1 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 @@ -12,27 +12,26 @@ import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgProvider; +import java.util.function.Supplier; + @Configuration public class KafkaConfig { @Autowired @Bean - public Producer createProducer(ImportDBFServiceSettings settings) { - return KafkaProducerFactory.producer(settings.getKafkaProducer()); - } - - @Autowired - @Bean - public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { - ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); - return KafkaSender - .setup() - .producer(kafkaProducer) - .idGenerator(imdgIdGenerator::nextId) - .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); - return imdg::insert; - }) - .build(); + public Supplier kafkaSender(ImportDBFServiceSettings settings, ImdgProvider imdgProvider) { + return () -> { + Producer kafkaProducer = KafkaProducerFactory.producer(settings.getKafkaProducer()); + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .producer(kafkaProducer) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); + return imdg::insert; + }) + .build(); + }; } } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java index 03a98e92b..709fd1a9e 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/ImportToDB.java @@ -18,6 +18,7 @@ import java.io.IOException; import java.io.InputStream; import java.nio.charset.Charset; import java.util.Map; +import java.util.function.Supplier; /** * Заливка проверенных данных в базу @@ -27,11 +28,11 @@ public class ImportToDB extends Stage { private final AProperties properties; private final HazelcastService hazelcastService; private final Map mappingEnumTableObjectTable; - private final KafkaSender kafka; + private final Supplier kafka; public ImportToDB(@Qualifier("dbfImporterProperties") AProperties properties, HazelcastService hazelcastService, @Qualifier("mapOfTable") Map mappingEnumTableObjectTable, - KafkaSender kafka) { + Supplier kafka) { this.properties = properties; this.hazelcastService = hazelcastService; this.mappingEnumTableObjectTable = mappingEnumTableObjectTable; @@ -70,6 +71,6 @@ public class ImportToDB extends Stage { private void sendStatementRequest(Long fileId) { StatementRequest statementRequest = new StatementRequest(); statementRequest.setSdf01GroupId(fileId); - kafka.sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); } }