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