http://jira.mfd.msk:8088/browse/CLS-730 separate hazelcast and kafka instances removed, application.properties updated

This commit is contained in:
Ivan Nikolaev-Axenov 2024-09-12 18:46:16 +03:00
parent bcd1aa2b40
commit 0ae83b1c52
13 changed files with 118 additions and 317 deletions

View file

@ -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());
}
}

View file

@ -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<String, Object> pfSdf(ImportXMLServiceSettings settings) {
KafkaProducerSettings kafkaSettings = settings.getSdfHazelcastAndKafka().getKafkaProducer();
@Bean
public ProducerFactory<String, Object> 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<String, Object> 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<String, Object> kafkaTemplateSdf(@Qualifier("pfSdf") ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Bean("kafkaTemplateLks")
@ConditionalOnProperty(value = "import-xml-service.process-lks-files", havingValue = "true")
public KafkaTemplate<String, Object> kafkaTemplateLks(@Qualifier("pfLks") ProducerFactory<String, Object> pf) {
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean("kafkaSenderSdf")
@ConditionalOnProperty(value = "import-xml-service.process-sdf-files", havingValue = "true")
public Supplier<KafkaSender> kafkaSenderSdf(@Qualifier("kafkaTemplateSdf") KafkaTemplate<String, Object> 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<KafkaSender> kafkaSenderLks(@Qualifier("kafkaTemplateLks") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("lksImdgProvider") ImdgProvider imdgProvider) {
@Bean
public Supplier<KafkaSender> kafkaSender(KafkaTemplate<String, Object> 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<String, Object> 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<String, Object> createConsumerLks(ImportXMLServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getLksHazelcastAndKafka().getKafkaConsumer());
@Bean
public Consumer<String, Object> createConsumer(ImportXMLServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
}
}

View file

@ -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;

View file

@ -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() {

View file

@ -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;
}
}

View file

@ -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);
}

View file

@ -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<KafkaSender> kafkaSdf;
private final Supplier<KafkaSender> kafkaLks;
private final Supplier<KafkaSender> kafka;
private final Map<ETable, Consumer<Long>> messengers;
private final ImportXMLServiceSettings settings;
public XmlImportKafkaMessenger(@Qualifier("kafkaSenderSdf") Supplier<KafkaSender> kafkaSdf,
@Qualifier("kafkaSenderLks") Supplier<KafkaSender> kafkaLks,
public XmlImportKafkaMessenger(Supplier<KafkaSender> 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);
}
}

View file

@ -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;

View file

@ -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
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

View file

@ -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) {

View file

@ -130,9 +130,9 @@ public class KafkaTestConfig {
return kafkaTemplate;
}
@Bean("kafkaSenderSdf")
public Supplier<KafkaSender> kafkaSenderSdf(@Qualifier("kafkaTestTemplate") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("sdfImdgProvider") ImdgProvider imdgProvider,
@Bean("kafkaSender")
public Supplier<KafkaSender> kafkaSender(@Qualifier("kafkaTestTemplate") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("imdgProvider") ImdgProvider imdgProvider,
Producer<String, Object> mockProducer) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return () -> KafkaSender
@ -140,27 +140,7 @@ public class KafkaTestConfig {
.setKafkaTemplate(kafkaTemplate)
.producer(mockProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
@Bean("kafkaSenderLks")
public Supplier<KafkaSender> kafkaSenderLks(@Qualifier("kafkaTestTemplate") KafkaTemplate<String, Object> kafkaTemplate,
@Qualifier("lksImdgProvider") ImdgProvider imdgProvider,
Producer<String, Object> mockProducer) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return () -> KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.producer(mockProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.saveRequestInfo(false)
.build();
}

View file

@ -59,7 +59,7 @@ class ImportToDBTest {
private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("dd.MM.yy");
@Autowired
@Qualifier("sdfImdgProvider")
@Qualifier("imdgProvider")
private HazelcastService hazelcastService;
@Autowired

View file

@ -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
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