From 2e46a484ffe263c76e40fa4d916b9b36459548f4 Mon Sep 17 00:00:00 2001 From: ialbert Date: Wed, 28 Sep 2022 18:39:02 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-36 --- .../balance/service/Sdf02Service.java | 2 + clearing-parent/dbf-importer/pom.xml | 4 ++ .../dbf/importer/config/KafkaConfig.java | 38 +++++++++++++++++++ .../settings/ImportDBFServiceSettings.java | 10 +++++ .../dbf/importer/logic/stages/ImportToDB.java | 20 +++++++++- .../src/main/resources/application.properties | 10 ++++- .../domain/cud/balance/StatementRequest.java | 1 - 7 files changed, 81 insertions(+), 4 deletions(-) create mode 100644 clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java index 1b873fde7..451a98c1c 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf02Service.java @@ -14,6 +14,8 @@ import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +//fixme remove +@Deprecated @Service public class Sdf02Service extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); diff --git a/clearing-parent/dbf-importer/pom.xml b/clearing-parent/dbf-importer/pom.xml index 9f8772d3e..d3d50fe59 100644 --- a/clearing-parent/dbf-importer/pom.xml +++ b/clearing-parent/dbf-importer/pom.xml @@ -62,6 +62,10 @@ ru.spcex.clearing classes + + ru.spcex.platform + platform-messaging + 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 new file mode 100644 index 000000000..e3d6ae042 --- /dev/null +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java @@ -0,0 +1,38 @@ +package ru.spcex.clearing.dbf.importer.config; + +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; +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; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@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(); + } +} 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 788715794..d697ec445 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.KafkaProducerSettings; import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @Component @@ -10,6 +11,7 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; @ConfigurationProperties("import-dbf-service") public class ImportDBFServiceSettings { private HazelcastClientParams hazelcast; + private KafkaProducerSettings kafkaProducer; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -18,4 +20,12 @@ public class ImportDBFServiceSettings { public void setHazelcast(HazelcastClientParams hazelcast) { this.hazelcast = hazelcast; } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } } 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 fce0faa4b..03a98e92b 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 @@ -8,6 +8,9 @@ import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; import ru.spcex.clearing.dbf.importer.logic.data.enums.StageResult; import ru.spcex.clearing.dbf.importer.logic.data.tables.AbstractTable; import ru.spcex.clearing.dbf.importer.properties.AProperties; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; import java.io.ByteArrayInputStream; @@ -24,12 +27,15 @@ public class ImportToDB extends Stage { private final AProperties properties; private final HazelcastService hazelcastService; private final Map mappingEnumTableObjectTable; + private final KafkaSender kafka; public ImportToDB(@Qualifier("dbfImporterProperties") AProperties properties, - HazelcastService hazelcastService, @Qualifier("mapOfTable") Map mappingEnumTableObjectTable) { + HazelcastService hazelcastService, @Qualifier("mapOfTable") Map mappingEnumTableObjectTable, + KafkaSender kafka) { this.properties = properties; this.hazelcastService = hazelcastService; this.mappingEnumTableObjectTable = mappingEnumTableObjectTable; + this.kafka = kafka; } @Override @@ -43,11 +49,15 @@ public class ImportToDB extends Stage { AbstractTable table = mappingEnumTableObjectTable.get(currTable); table.setHazelcastService(hazelcastService); table.setFilename(resultContainer.getDbfFile().getName()); - table.setFileId(hazelcastService.getImdgIdGenerator().nextId()); + Long fileId = hazelcastService.getImdgIdGenerator().nextId(); + table.setFileId(fileId); for (int i = 0; i < dbfReader.getRecordCount(); i++) { Object[] entity = dbfReader.nextRecord(); table.injectEntity(table.getEntity(entity)); } + if (ETable.DF_01.equals(currTable)) { + sendStatementRequest(fileId); + } } catch (IOException exception) { log.warn(exception.getMessage()); return StageResult.ERROR; @@ -56,4 +66,10 @@ public class ImportToDB extends Stage { return StageResult.OK; } + + private void sendStatementRequest(Long fileId) { + StatementRequest statementRequest = new StatementRequest(); + statementRequest.setSdf01GroupId(fileId); + kafka.sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + } } diff --git a/clearing-parent/dbf-importer/src/main/resources/application.properties b/clearing-parent/dbf-importer/src/main/resources/application.properties index df400cde1..79ce0025c 100644 --- a/clearing-parent/dbf-importer/src/main/resources/application.properties +++ b/clearing-parent/dbf-importer/src/main/resources/application.properties @@ -10,4 +10,12 @@ dbf.out-dir=/opt/clearing/file/importer/loaded/ dbf.threads-count=10 import-dbf-service.hazelcast.cluster-members=127.0.0.1:5701 import-dbf-service.hazelcast.login=dev -import-dbf-service.hazelcast.password=dev-pass \ No newline at end of file +import-dbf-service.hazelcast.password=dev-pass + + +import-dbf-service.kafka-producer.bootstrap-servers=localhost:9092 +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 diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java index 06193017e..e2477137a 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/balance/StatementRequest.java @@ -12,7 +12,6 @@ public class StatementRequest { //from accountBalance creation @JsonProperty List accountCreationResults = new ArrayList<>(); - //from sDf02 creation // public StatementRequestType getType() { // if (errorCode != null || errorText != null || status != null) {