diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/ImporterImdgConfig.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/ImporterImdgConfig.java index 32243b8f3..0c54ab12d 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/ImporterImdgConfig.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/ImporterImdgConfig.java @@ -2,7 +2,6 @@ package ru.spcex.clearing.xml.importer.config; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; @@ -33,28 +32,14 @@ public class ImporterImdgConfig { } @Autowired - @Bean(name = "sdfImdgProvider") - @ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true") - public HazelcastService sdfImdgProvider( + @Bean + public HazelcastService imdgProvider( @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, ImportXMLServiceSettings settings ) { return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, - settings.getSdfHazelcastAndKafka().getHazelcast()); - } - - @Autowired - @Bean(name = "lksImdgProvider") - @ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true") - public HazelcastService lksImdgProvider( - @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, - @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, - ImportXMLServiceSettings settings - ) { - return new HazelcastService(taskExecutorHazelcastClientInitializer, - taskExecutorIdGeneratorAwaiter, - settings.getLksHazelcastAndKafka().getHazelcast()); + settings.getHazelcast()); } } diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/KafkaConfig.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/KafkaConfig.java index 0a106ef19..77c196d6c 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/KafkaConfig.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/KafkaConfig.java @@ -3,9 +3,7 @@ package ru.spcex.clearing.xml.importer.config; import java.util.function.Supplier; import org.apache.kafka.clients.consumer.Consumer; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.ConfigurableBeanFactory; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; @@ -21,51 +19,21 @@ import ru.spcex.platform.imdg.api.ImdgProvider; @Configuration public class KafkaConfig { - @Bean(name = "pfSdf") - @ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true") - public ProducerFactory pfSdf(ImportXMLServiceSettings settings) { - KafkaProducerSettings kafkaSettings = settings.getSdfHazelcastAndKafka().getKafkaProducer(); + @Bean + public ProducerFactory pf(ImportXMLServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); return KafkaProducerFactory.producerFactory(kafkaSettings); } - @Bean(name = "pfLks") - @ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true") - public ProducerFactory pfLks(ImportXMLServiceSettings settings) { - KafkaProducerSettings kafkaSettings = settings.getLksHazelcastAndKafka().getKafkaProducer(); - return KafkaProducerFactory.producerFactory(kafkaSettings); - } - - @Bean("kafkaTemplateSdf") - @ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true") - public KafkaTemplate kafkaTemplateSdf(@Qualifier("pfSdf") ProducerFactory pf) { - return new KafkaTemplate<>(pf); - } - - @Bean("kafkaTemplateLks") - @ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true") - public KafkaTemplate kafkaTemplateLks(@Qualifier("pfLks") ProducerFactory pf) { + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { return new KafkaTemplate<>(pf); } @Autowired - @Bean("kafkaSenderSdf") - @ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true") - public Supplier kafkaSenderSdf(@Qualifier("kafkaTemplateSdf") KafkaTemplate kafkaTemplate, - @Qualifier("sdfImdgProvider") ImdgProvider imdgProvider) { - ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); - return () -> KafkaSender - .setup() - .setKafkaTemplate(kafkaTemplate) - .idGenerator(imdgIdGenerator::nextId) - .saveRequestInfo(false) - .build(); - } - - @Autowired - @Bean("kafkaSenderLks") - @ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true") - public Supplier kafkaSenderLks(@Qualifier("kafkaTemplateLks") KafkaTemplate kafkaTemplate, - @Qualifier("lksImdgProvider") ImdgProvider imdgProvider) { + @Bean + public Supplier kafkaSender(KafkaTemplate kafkaTemplate, + ImdgProvider imdgProvider) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return () -> KafkaSender .setup() @@ -77,17 +45,8 @@ public class KafkaConfig { @Autowired @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) - @Bean("createConsumerSdf") - @ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true") - public Consumer createConsumerSdf(ImportXMLServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getSdfHazelcastAndKafka().getKafkaConsumer()); - } - - @Autowired - @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) - @Bean("createConsumerLks") - @ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true") - public Consumer createConsumerLks(ImportXMLServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getLksHazelcastAndKafka().getKafkaConsumer()); + @Bean + public Consumer createConsumer(ImportXMLServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } } diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/XMLImporterConfig.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/XMLImporterConfig.java index 092313f9c..0a7f74107 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/XMLImporterConfig.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/XMLImporterConfig.java @@ -13,7 +13,6 @@ import javax.xml.stream.XMLInputFactory; import javax.xml.stream.XMLOutputFactory; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.ApplicationContext; @@ -43,7 +42,7 @@ public class XMLImporterConfig { public XMLImporterConfig(ImportXMLServiceSettings settings, ApplicationContext context, - @Qualifier("sdfImdgProvider") ImdgProvider hazelcastService) { + ImdgProvider hazelcastService) { this.settings = settings; this.context = context; this.hazelcastService = hazelcastService; diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/ImportXMLServiceSettings.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/ImportXMLServiceSettings.java index e32cf7da4..121f2eb40 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/ImportXMLServiceSettings.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/ImportXMLServiceSettings.java @@ -4,15 +4,20 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.NestedConfigurationProperty; 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; @Component @PropertySource("file:${spring.config.location}/application.properties") @ConfigurationProperties("import-xml-service") public class ImportXMLServiceSettings { @NestedConfigurationProperty - private KafkaHazelcastSettings sdfHazelcastAndKafka; + private HazelcastClientParams hazelcast; @NestedConfigurationProperty - private KafkaHazelcastSettings lksHazelcastAndKafka; + private KafkaProducerSettings kafkaProducer; + @NestedConfigurationProperty + private KafkaConsumerSettings kafkaConsumer; @NestedConfigurationProperty private StoreSettings storeSdf; @@ -27,20 +32,28 @@ public class ImportXMLServiceSettings { private boolean processSdfFiles; private boolean processLksFiles; - public KafkaHazelcastSettings getSdfHazelcastAndKafka() { - return sdfHazelcastAndKafka; + public HazelcastClientParams getHazelcast() { + return hazelcast; } - public void setSdfHazelcastAndKafka(KafkaHazelcastSettings sdfHazelcastAndKafka) { - this.sdfHazelcastAndKafka = sdfHazelcastAndKafka; + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; } - public KafkaHazelcastSettings getLksHazelcastAndKafka() { - return lksHazelcastAndKafka; + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; } - public void setLksHazelcastAndKafka(KafkaHazelcastSettings lksHazelcastAndKafka) { - this.lksHazelcastAndKafka = lksHazelcastAndKafka; + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; } public StoreSettings getStoreSdf() { diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/KafkaHazelcastSettings.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/KafkaHazelcastSettings.java deleted file mode 100644 index 0b689551f..000000000 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/config/settings/KafkaHazelcastSettings.java +++ /dev/null @@ -1,41 +0,0 @@ -package ru.spcex.clearing.xml.importer.config.settings; - -import org.springframework.boot.context.properties.NestedConfigurationProperty; -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; - -public class KafkaHazelcastSettings { - @NestedConfigurationProperty - private HazelcastClientParams hazelcast; - - @NestedConfigurationProperty - private KafkaProducerSettings kafkaProducer; - - @NestedConfigurationProperty - private KafkaConsumerSettings kafkaConsumer; - - public HazelcastClientParams getHazelcast() { - return hazelcast; - } - - public void setHazelcast(HazelcastClientParams hazelcast) { - this.hazelcast = hazelcast; - } - - public KafkaProducerSettings getKafkaProducer() { - return kafkaProducer; - } - - public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { - this.kafkaProducer = kafkaProducer; - } - - public KafkaConsumerSettings getKafkaConsumer() { - return kafkaConsumer; - } - - public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { - this.kafkaConsumer = kafkaConsumer; - } -} diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDB.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDB.java index d3f18527b..545a61ef6 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDB.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDB.java @@ -2,13 +2,11 @@ package ru.spcex.clearing.xml.importer.logic.steps; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Component; import ru.clearing.classes.statics.data.sdf.SDf57; import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewRequest; import ru.spcex.clearing.xml.importer.logic.data.ResultContainer; import ru.spcex.clearing.xml.importer.logic.data.enums.ETable; -import ru.spcex.clearing.xml.importer.logic.data.enums.FileType; import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult; import ru.spcex.clearing.xml.importer.logic.data.tags.lks.AccountListCur; import ru.spcex.clearing.xml.importer.logic.data.tags.lks.AccountListRub; @@ -24,15 +22,12 @@ import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; @Component public class ImportToDB { private final Logger log = LoggerFactory.getLogger(getClass()); - private final HazelcastService hazelcastServiceSdf; - private final HazelcastService hazelcastServiceLks; + private final HazelcastService hazelcastService; private final XmlImportKafkaMessenger kafkaMessenger; - public ImportToDB(@Qualifier("sdfImdgProvider") HazelcastService hazelcastServiceSdf, - @Qualifier("lksImdgProvider") HazelcastService hazelcastServiceLks, + public ImportToDB(HazelcastService hazelcastService, XmlImportKafkaMessenger kafkaMessenger) { - this.hazelcastServiceSdf = hazelcastServiceSdf; - this.hazelcastServiceLks = hazelcastServiceLks; + this.hazelcastService = hazelcastService; this.kafkaMessenger = kafkaMessenger; } @@ -50,9 +45,9 @@ public class ImportToDB { DocumentTag documentTag = (DocumentTag) resultContainer.getXmlFile(); String fileName = resultContainer.getXmlFilePath().getName(); - Long fileId = hazelcastServiceSdf.getImdgIdGenerator().nextId(); + Long fileId = hazelcastService.getImdgIdGenerator().nextId(); for (ObjectTag o : documentTag.getObjects()) { - o.setHazelcastService(hazelcastServiceSdf); + o.setHazelcastService(hazelcastService); o.setFileName(fileName); o.setGenerationId(fileId); @@ -60,7 +55,7 @@ public class ImportToDB { if (!o.checkOnExisting(entityTable) && entityTable instanceof SDf57) { kafkaMessenger.sendUserNotification(ObjectType.rgst, "Номер транзакции " + ((SDf57) entityTable).getDbfId() + " в полученном df57 уже был обработан ранее", - Priority.HIGH, FileType.SDF); + Priority.HIGH); } else { o.injectEntity(entityTable); } diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/XmlImportKafkaMessenger.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/XmlImportKafkaMessenger.java index fd58beb1c..bc54f1f82 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/XmlImportKafkaMessenger.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/logic/steps/XmlImportKafkaMessenger.java @@ -9,7 +9,6 @@ import java.util.function.Supplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Component; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewRequest; @@ -22,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.xml.importer.config.settings.ImportXMLServiceSettings; import ru.spcex.clearing.xml.importer.logic.data.ResultContainer; import ru.spcex.clearing.xml.importer.logic.data.enums.ETable; -import ru.spcex.clearing.xml.importer.logic.data.enums.FileType; import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult; import ru.spcex.platform.enumeration.ObjectType; import ru.spcex.platform.enumeration.Priority; @@ -31,34 +29,31 @@ import ru.spcex.platform.enumeration.SdfTable; @Component public class XmlImportKafkaMessenger implements InitializingBean { final Logger log = LoggerFactory.getLogger(getClass()); - private final Supplier kafkaSdf; - private final Supplier kafkaLks; + private final Supplier kafka; private final Map> messengers; private final ImportXMLServiceSettings settings; - public XmlImportKafkaMessenger(@Qualifier("kafkaSenderSdf") Supplier kafkaSdf, - @Qualifier("kafkaSenderLks") Supplier kafkaLks, + public XmlImportKafkaMessenger(Supplier kafka, ImportXMLServiceSettings settings) { - this.kafkaSdf = kafkaSdf; - this.kafkaLks = kafkaLks; + this.kafka = kafka; this.settings = settings; this.messengers = new HashMap<>(); } @Override public void afterPropertiesSet() { - messengers.put(ETable.DF_01, groupId -> messageStatement(groupId, SdfTable.SDF_01, FileType.SDF)); - messengers.put(ETable.DF_04, groupId -> messageStatement(groupId, SdfTable.SDF_04, FileType.SDF)); - messengers.put(ETable.DF_06, groupId -> messageStatement(groupId, SdfTable.SDF_06, Consts.STATEMENT_PROCESS_SDF06, FileType.SDF)); - messengers.put(ETable.DF_52, groupId -> messageStatement(groupId, SdfTable.SDF_52, Consts.ACCOUNT_PROCESS_SDF52, FileType.SDF)); - messengers.put(ETable.DF_55, groupId -> messageStatement(groupId, SdfTable.SDF_55, FileType.SDF)); - messengers.put(ETable.DF_57, groupId -> messageStatement(groupId, SdfTable.SDF_57, FileType.SDF)); + messengers.put(ETable.DF_01, groupId -> messageStatement(groupId, SdfTable.SDF_01)); + messengers.put(ETable.DF_04, groupId -> messageStatement(groupId, SdfTable.SDF_04)); + messengers.put(ETable.DF_06, groupId -> messageStatement(groupId, SdfTable.SDF_06, Consts.STATEMENT_PROCESS_SDF06)); + messengers.put(ETable.DF_52, groupId -> messageStatement(groupId, SdfTable.SDF_52, Consts.ACCOUNT_PROCESS_SDF52)); + messengers.put(ETable.DF_55, groupId -> messageStatement(groupId, SdfTable.SDF_55)); + messengers.put(ETable.DF_57, groupId -> messageStatement(groupId, SdfTable.SDF_57)); } public boolean sendToBankAccount(BankAccountNewRequest bankAccountNewRequest) { try { - kafkaSdf.get().sendRequestToQueue(Consts.DESTINATION_BANK_ACCOUNT_NEW, bankAccountNewRequest); + kafka.get().sendRequestToQueue(Consts.DESTINATION_BANK_ACCOUNT_NEW, bankAccountNewRequest); } catch (Exception e) { log.error("Error occurred on sending message to kafka {}", e.getMessage(), e); return false; @@ -85,7 +80,7 @@ public class XmlImportKafkaMessenger implements InitializingBean { String fileName = resultContainer.getXmlFilePath().getName(); log.debug("Notify user about error {} \"{}\"", resultContainer.getXmlTable(), fileName); - sendUserNotification(ObjectType.rgst, String.format("Файл \"%s\" не сохранен", fileName), Priority.HIGH, resultContainer.getXmlTable().getFileType()); + sendUserNotification(ObjectType.rgst, String.format("Файл \"%s\" не сохранен", fileName), Priority.HIGH); } public void notifyUserAboutSuccessLoad(ResultContainer resultContainer) { @@ -96,11 +91,11 @@ public class XmlImportKafkaMessenger implements InitializingBean { String fileName = resultContainer.getXmlFilePath().getName(); log.debug("Notify user about success {} \"{}\"", resultContainer.getXmlTable(), fileName); - sendUserNotification(ObjectType.rgst, String.format("Загружен \"%s\" - успешно", fileName), Priority.LOW, resultContainer.getXmlTable().getFileType()); + sendUserNotification(ObjectType.rgst, String.format("Загружен \"%s\" - успешно", fileName), Priority.LOW); } } - public void sendUserNotification(ObjectType objectType, String comment, Priority priority, FileType fileType) { + public void sendUserNotification(ObjectType objectType, String comment, Priority priority) { final String destination = Consts.NOTIFICATION_NEW; NotificationNewRequest request = new NotificationNewRequest(); request.setObjectType(objectType.getKey()); @@ -108,39 +103,22 @@ public class XmlImportKafkaMessenger implements InitializingBean { request.setComment(comment); log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); - switch (fileType) { - case SDF -> { - Long rKey = kafkaSdf.get().sendRequestToQueue(destination, request); - log.trace("For user send message to SDF kafka, request id={}", rKey); - } - case LKS -> { - Long rKey = kafkaLks.get().sendRequestToQueue(destination, request); - log.trace("For user send message to LKS kafka, request id={}", rKey); - } - } + Long rKey = kafka.get().sendRequestToQueue(destination, request); + log.trace("For user send message to kafka, request id={}", rKey); } - private void messageStatement(Long groupId, SdfTable table, FileType fileType) { - messageStatement(groupId, table, Consts.STATEMENT_PROCESS, fileType); + private void messageStatement(Long groupId, SdfTable table) { + messageStatement(groupId, table, Consts.STATEMENT_PROCESS); } - private void messageStatement(Long groupId, SdfTable table, String destination, FileType fileType) { + private void messageStatement(Long groupId, SdfTable table, String destination) { StatementRequest statementRequest = new StatementRequest(); statementRequest.setGroupId(groupId); statementRequest.setTable(table); - switch (fileType) { - case SDF -> { - Long msgId = kafkaSdf.get().sendRequestToQueue(destination, statementRequest); - log.debug("Send StatementRequest({}, {}) message id={} to SDF kafka \"{}\"", - groupId, table, msgId, destination); - } - case LKS -> { - Long msgId = kafkaLks.get().sendRequestToQueue(destination, statementRequest); - log.debug("Send StatementRequest({}, {}) message id={} to LKS kafka \"{}\"", - groupId, table, msgId, destination); - } - } + Long msgId = kafka.get().sendRequestToQueue(destination, statementRequest); + log.debug("Send StatementRequest({}, {}) message id={} to kafka \"{}\"", + groupId, table, msgId, destination); } public void sendPairSdfRequest(String fileName, Long fileId, String tableSdf) { @@ -149,7 +127,7 @@ public class XmlImportKafkaMessenger implements InitializingBean { request.setGenerationId(fileId); request.setTableSDf(tableSdf); - Long msgId = kafkaSdf.get().sendRequestToQueue(PAIR_SDF, request); + Long msgId = kafka.get().sendRequestToQueue(PAIR_SDF, request); log.info("Send PairSdfRequest={} message id={} to SDF kafka \"{}\"", LogFormatter.toString(request), msgId, PAIR_SDF); } @@ -157,6 +135,6 @@ public class XmlImportKafkaMessenger implements InitializingBean { Sdf04Request sdf04ImportNotification = new Sdf04Request(); sdf04ImportNotification.setGroupId(groupId); - kafkaSdf.get().sendRequestToQueue(Consts.SDF04_PROCESS, sdf04ImportNotification); + kafka.get().sendRequestToQueue(Consts.SDF04_PROCESS, sdf04ImportNotification); } } diff --git a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/services/FileChecker.java b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/services/FileChecker.java index f8662b972..c53823314 100644 --- a/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/services/FileChecker.java +++ b/clearing-parent/xml-importer/src/main/java/ru/spcex/clearing/xml/importer/services/FileChecker.java @@ -12,7 +12,6 @@ import java.util.LinkedList; import java.util.List; import java.util.Locale; import java.util.Map; -import java.util.stream.Collectors; import java.util.stream.Stream; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/clearing-parent/xml-importer/src/main/resources/application.properties b/clearing-parent/xml-importer/src/main/resources/application.properties index a7caae4a6..c67a3ba19 100644 --- a/clearing-parent/xml-importer/src/main/resources/application.properties +++ b/clearing-parent/xml-importer/src/main/resources/application.properties @@ -11,19 +11,13 @@ import-xml-service.common.encoding-source=cp866 import-xml-service.common.insert-batch-size=100 import-xml-service.common.threads-count=10 -# SDF settings -# Hazelcast cluster -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.cluster-members=10.200.200.181:5701 -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.login=dev -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.password=dev-pass - -# Store directories +# Store directories for SDf import-xml-service.store-sdf.delete-src-files=false import-xml-service.store-sdf.src-dir=/opt/clearing/file/xml-importer/ import-xml-service.store-sdf.out-dir=/opt/clearing/file/xml-importer/loaded/ -import-xml-service.store-sdf.out-dir-error=/opt/clearing/file/xml-importer/error/= +import-xml-service.store-sdf.out-dir-error=/opt/clearing/file/xml-importer/error/ -# Sftp directories and credentials +# Sftp directories and credentials for SDf import-xml-service.store-sdf.sftp-in.sftp-src-pay-val-dir.rub=clearing_xml-importer_sftp/rub/ import-xml-service.store-sdf.sftp-in.sftp-src-pay-val-dir.eur=clearing_xml-importer_sftp/eur/ import-xml-service.store-sdf.sftp-in.sftp-src-dir=clearing_xml-importer_sftp @@ -32,35 +26,13 @@ import-xml-service.store-sdf.sftp-in.password=********* import-xml-service.store-sdf.sftp-in.server-ip=127.0.0.1 import-xml-service.store-sdf.sftp-in.server-port=22 -# Kafka settings -# Kafka producer -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.bootstrap-servers=localhost:9092 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.acks=all -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.retries=0 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.batch-size=16384 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.linger-ms=1 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.buffer-memory=33554432 - -# Kafka consumer -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.bootstrap-servers=localhost:9092 -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.group-id=dev-group-clearing-service -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.enable-auto-commit=false -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.session-timeout-ms=30000 -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.auto-offset-reset=latest - -# LKS settings -# Hazelcast cluster -import-xml-service.lks-hazelcast-and-kafka.hazelcast.cluster-members=10.200.200.181:5701 -import-xml-service.lks-hazelcast-and-kafka.hazelcast.login=dev -import-xml-service.lks-hazelcast-and-kafka.hazelcast.password=dev-pass - -# Store directories +# Store directories for LKS import-xml-service.store-lks.delete-src-files=false import-xml-service.store-lks.src-dir=/opt/clearing/file/xml-importer/ import-xml-service.store-lks.out-dir=/opt/clearing/file/xml-importer/loaded/ -import-xml-service.store-lks.out-dir-error=/opt/clearing/file/xml-importer/error/= +import-xml-service.store-lks.out-dir-error=/opt/clearing/file/xml-importer/error/ -# Sftp directories and credentials +# Sftp directories and credentials for LKS import-xml-service.store-lks.sftp-in.sftp-src-pay-val-dir.rub=clearing_xml-importer_sftp/rub/ import-xml-service.store-lks.sftp-in.sftp-src-pay-val-dir.eur=clearing_xml-importer_sftp/eur/ import-xml-service.store-lks.sftp-in.sftp-src-dir=clearing_xml-importer_sftp @@ -69,18 +41,23 @@ import-xml-service.store-lks.sftp-in.password=********* import-xml-service.store-lks.sftp-in.server-ip=127.0.0.1 import-xml-service.store-lks.sftp-in.server-port=22 +# Hazelcast cluster +import-xml-service.hazelcast.cluster-members=10.200.200.181:5701 +import-xml-service.hazelcast.login=dev +import-xml-service.hazelcast.password=dev-pass + # Kafka settings # Kafka producer -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.bootstrap-servers=localhost:9092 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.acks=all -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.retries=0 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.batch-size=16384 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.linger-ms=1 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.buffer-memory=33554432 +import-xml-service.kafka-producer.bootstrap-servers=localhost:9092 +import-xml-service.kafka-producer.acks=all +import-xml-service.kafka-producer.retries=0 +import-xml-service.kafka-producer.batch-size=16384 +import-xml-service.kafka-producer.linger-ms=1 +import-xml-service.kafka-producer.buffer-memory=33554432 # Kafka consumer -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.bootstrap-servers=localhost:9092 -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.group-id=dev-group-clearing-service -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.enable-auto-commit=false -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.session-timeout-ms=30000 -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.auto-offset-reset=latest \ No newline at end of file +import-xml-service.kafka-consumer.bootstrap-servers=localhost:9092 +import-xml-service.kafka-consumer.group-id=dev-group-clearing-service +import-xml-service.kafka-consumer.enable-auto-commit=false +import-xml-service.kafka-consumer.session-timeout-ms=30000 +import-xml-service.kafka-consumer.auto-offset-reset=latest \ No newline at end of file diff --git a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/ImporterImdgTestConfig.java b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/ImporterImdgTestConfig.java index eb4f76190..7327c8829 100644 --- a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/ImporterImdgTestConfig.java +++ b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/ImporterImdgTestConfig.java @@ -64,28 +64,8 @@ public class ImporterImdgTestConfig { } @Autowired - @Bean(name = "sdfImdgProvider") - public HazelcastService imdgTestProviderSdf( - @Qualifier("taskExecutorHazelcastTestClientInitializerXmlImporter") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, - @Qualifier("taskExecutorTestIdGeneratorAwaiterXmlImporter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, - @Qualifier("hazelcastClientParamsXmlImporter") HazelcastClientParams params) { - Config cfg = new Config(); - cfg.setInstanceName("localhost"); - - NetworkConfig networkConfig = new NetworkConfig(); - JoinConfig joinConfig = new JoinConfig(); - joinConfig.setMulticastConfig(new MulticastConfig().setEnabled(false)); - joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1"))); - networkConfig.setJoin(joinConfig); - cfg.setNetworkConfig(networkConfig); - hazelcastInstance = Hazelcast.getOrCreateHazelcastInstance(cfg); - HazelcastHelper.imdgSystem_setStorageState(true, hazelcastInstance); - return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); - } - - @Autowired - @Bean(name = "lksImdgProvider") - public HazelcastService imdgTestProviderLks( + @Bean(name = "imdgProvider") + public HazelcastService imdgTestProvider( @Qualifier("taskExecutorHazelcastTestClientInitializerXmlImporter") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorTestIdGeneratorAwaiterXmlImporter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, @Qualifier("hazelcastClientParamsXmlImporter") HazelcastClientParams params) { diff --git a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/KafkaTestConfig.java b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/KafkaTestConfig.java index 01767a7ec..2770b648c 100644 --- a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/KafkaTestConfig.java +++ b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/config/KafkaTestConfig.java @@ -130,9 +130,9 @@ public class KafkaTestConfig { return kafkaTemplate; } - @Bean("kafkaSenderSdf") - public Supplier kafkaSenderSdf(@Qualifier("kafkaTestTemplate") KafkaTemplate kafkaTemplate, - @Qualifier("sdfImdgProvider") ImdgProvider imdgProvider, + @Bean("kafkaSender") + public Supplier kafkaSender(@Qualifier("kafkaTestTemplate") KafkaTemplate kafkaTemplate, + @Qualifier("imdgProvider") ImdgProvider imdgProvider, Producer mockProducer) { ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); return () -> KafkaSender @@ -140,27 +140,7 @@ public class KafkaTestConfig { .setKafkaTemplate(kafkaTemplate) .producer(mockProducer) .idGenerator(imdgIdGenerator::nextId) - .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); - return imdg::insert; - }) - .build(); - } - - @Bean("kafkaSenderLks") - public Supplier kafkaSenderLks(@Qualifier("kafkaTestTemplate") KafkaTemplate kafkaTemplate, - @Qualifier("lksImdgProvider") ImdgProvider imdgProvider, - Producer mockProducer) { - ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); - return () -> KafkaSender - .setup() - .setKafkaTemplate(kafkaTemplate) - .producer(mockProducer) - .idGenerator(imdgIdGenerator::nextId) - .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); - return imdg::insert; - }) + .saveRequestInfo(false) .build(); } diff --git a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDBTest.java b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDBTest.java index 7ef3616ab..34a191bbd 100644 --- a/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDBTest.java +++ b/clearing-parent/xml-importer/src/test/java/ru/spcex/clearing/xml/importer/logic/steps/ImportToDBTest.java @@ -59,7 +59,7 @@ class ImportToDBTest { private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("dd.MM.yy"); @Autowired - @Qualifier("sdfImdgProvider") + @Qualifier("imdgProvider") private HazelcastService hazelcastService; @Autowired diff --git a/clearing-parent/xml-importer/src/test/resources/application.properties b/clearing-parent/xml-importer/src/test/resources/application.properties index 386256555..c67a3ba19 100644 --- a/clearing-parent/xml-importer/src/test/resources/application.properties +++ b/clearing-parent/xml-importer/src/test/resources/application.properties @@ -11,19 +11,13 @@ import-xml-service.common.encoding-source=cp866 import-xml-service.common.insert-batch-size=100 import-xml-service.common.threads-count=10 -# SDF settings -# Hazelcast cluster -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.cluster-members=localhost:5701 -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.login=dev -import-xml-service.sdf-hazelcast-and-kafka.hazelcast.password=dev-pass - -# Store directories +# Store directories for SDf import-xml-service.store-sdf.delete-src-files=false import-xml-service.store-sdf.src-dir=/opt/clearing/file/xml-importer/ import-xml-service.store-sdf.out-dir=/opt/clearing/file/xml-importer/loaded/ -import-xml-service.store-sdf.out-dir-error=/opt/clearing/file/xml-importer/error/= +import-xml-service.store-sdf.out-dir-error=/opt/clearing/file/xml-importer/error/ -# Sftp directories and credentials +# Sftp directories and credentials for SDf import-xml-service.store-sdf.sftp-in.sftp-src-pay-val-dir.rub=clearing_xml-importer_sftp/rub/ import-xml-service.store-sdf.sftp-in.sftp-src-pay-val-dir.eur=clearing_xml-importer_sftp/eur/ import-xml-service.store-sdf.sftp-in.sftp-src-dir=clearing_xml-importer_sftp @@ -32,35 +26,13 @@ import-xml-service.store-sdf.sftp-in.password=********* import-xml-service.store-sdf.sftp-in.server-ip=127.0.0.1 import-xml-service.store-sdf.sftp-in.server-port=22 -# Kafka settings -# Kafka producer -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.bootstrap-servers=localhost:9092 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.acks=all -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.retries=0 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.batch-size=16384 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.linger-ms=1 -import-xml-service.sdf-hazelcast-and-kafka.kafka-producer.buffer-memory=33554432 - -# Kafka consumer -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.bootstrap-servers=localhost:9092 -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.group-id=dev-group-clearing-service -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.enable-auto-commit=false -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.session-timeout-ms=30000 -import-xml-service.sdf-hazelcast-and-kafka.kafka-consumer.auto-offset-reset=latest - -# LKS settings -# Hazelcast cluster -import-xml-service.lks-hazelcast-and-kafka.hazelcast.cluster-members=localhost:5701 -import-xml-service.lks-hazelcast-and-kafka.hazelcast.login=dev -import-xml-service.lks-hazelcast-and-kafka.hazelcast.password=dev-pass - -# Store directories +# Store directories for LKS import-xml-service.store-lks.delete-src-files=false import-xml-service.store-lks.src-dir=/opt/clearing/file/xml-importer/ import-xml-service.store-lks.out-dir=/opt/clearing/file/xml-importer/loaded/ -import-xml-service.store-lks.out-dir-error=/opt/clearing/file/xml-importer/error/= +import-xml-service.store-lks.out-dir-error=/opt/clearing/file/xml-importer/error/ -# Sftp directories and credentials +# Sftp directories and credentials for LKS import-xml-service.store-lks.sftp-in.sftp-src-pay-val-dir.rub=clearing_xml-importer_sftp/rub/ import-xml-service.store-lks.sftp-in.sftp-src-pay-val-dir.eur=clearing_xml-importer_sftp/eur/ import-xml-service.store-lks.sftp-in.sftp-src-dir=clearing_xml-importer_sftp @@ -69,18 +41,23 @@ import-xml-service.store-lks.sftp-in.password=********* import-xml-service.store-lks.sftp-in.server-ip=127.0.0.1 import-xml-service.store-lks.sftp-in.server-port=22 +# Hazelcast cluster +import-xml-service.hazelcast.cluster-members=10.200.200.181:5701 +import-xml-service.hazelcast.login=dev +import-xml-service.hazelcast.password=dev-pass + # Kafka settings # Kafka producer -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.bootstrap-servers=localhost:9092 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.acks=all -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.retries=0 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.batch-size=16384 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.linger-ms=1 -import-xml-service.lks-hazelcast-and-kafka.kafka-producer.buffer-memory=33554432 +import-xml-service.kafka-producer.bootstrap-servers=localhost:9092 +import-xml-service.kafka-producer.acks=all +import-xml-service.kafka-producer.retries=0 +import-xml-service.kafka-producer.batch-size=16384 +import-xml-service.kafka-producer.linger-ms=1 +import-xml-service.kafka-producer.buffer-memory=33554432 # Kafka consumer -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.bootstrap-servers=localhost:9092 -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.group-id=dev-group-clearing-service -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.enable-auto-commit=false -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.session-timeout-ms=30000 -import-xml-service.lks-hazelcast-and-kafka.kafka-consumer.auto-offset-reset=latest \ No newline at end of file +import-xml-service.kafka-consumer.bootstrap-servers=localhost:9092 +import-xml-service.kafka-consumer.group-id=dev-group-clearing-service +import-xml-service.kafka-consumer.enable-auto-commit=false +import-xml-service.kafka-consumer.session-timeout-ms=30000 +import-xml-service.kafka-consumer.auto-offset-reset=latest \ No newline at end of file