diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/DBFImporterConfig.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/DBFImporterConfig.java index cfb331986..50b976f3f 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/DBFImporterConfig.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/DBFImporterConfig.java @@ -57,6 +57,7 @@ public class DBFImporterConfig { public Map getMapOfTables() { Map map = new HashMap<>(); map.put(ETable.DF_01, new SDf01Table()); + map.put(ETable.DF_04, new SDf04Table()); map.put(ETable.DF_09, new SDf09Table()); map.put(ETable.DF_12, new SDf12Table()); map.put(ETable.DF_16, new SDf16Table()); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java index ff802e79f..c99d96dac 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/enums/ETable.java @@ -2,6 +2,7 @@ package ru.spcex.clearing.dbf.importer.logic.data.enums; public enum ETable { DF_01("DF-01"), + DF_04("DF-04"), DF_09("DF-09"), DF_12("DF-12"), DF_16("DF-16"); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/SDf04Table.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/SDf04Table.java new file mode 100644 index 000000000..47022da4d --- /dev/null +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/data/tables/SDf04Table.java @@ -0,0 +1,74 @@ +package ru.spcex.clearing.dbf.importer.logic.data.tables; + +import ru.clearing.classes.statics.data.sdf.SDf04; +import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; +import ru.spcex.clearing.imdg.IMDGDistributedNames; + +import java.time.Instant; + +public class SDf04Table extends AbstractTable { + + private static final String PREFIX = ETable.DF_04.name(); + private static final Class CLAZZ = SDf04.class; + private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf04; + + public SDf04Table() { + super(PREFIX, CLAZZ, NAME_OF_HZ_MAP); + } + + @Override + public SDf04 getEntity(Object[] entity) { + SDf04 result = new SDf04(); + result.setSeg_type((String) entity[0]); + result.setDoc_type((String) entity[1]); + result.setDocnm_ref((String) entity[2]); + result.setDocnmprev((String) entity[3]); + result.setPriority((String) entity[4]); + result.setSbankcode((String) entity[5]); + result.setC_acc_deb((String) entity[6]); + result.setSbanknam1((String) entity[7]); + result.setSbanknam2((String) entity[8]); + result.setSbanknam3((String) entity[9]); + result.setSbanknam4((String) entity[10]); + result.setSbanknam5((String) entity[11]); + result.setRbankcode((String) entity[12]); + result.setC_acc_cred((String) entity[13]); + result.setRbanknam1((String) entity[14]); + result.setRbanknam2((String) entity[15]); + result.setRbanknam3((String) entity[16]); + result.setRbanknam4((String) entity[17]); + result.setRbanknam5((String) entity[18]); + result.setPay_date((String) entity[19]); + result.setExt_date((String) entity[20]); + result.setPay_val((String) entity[21]); + result.setSum_deb((String) entity[22]); + result.setSclientn1((String) entity[23]); + result.setSclientn2((String) entity[24]); + result.setSclientn3((String) entity[25]); + result.setSclientn4((String) entity[26]); + result.setSc_code((String) entity[27]); + result.setAcc_deb((String) entity[28]); + result.setRclientn1((String) entity[29]); + result.setRclientn2((String) entity[30]); + result.setRclientn3((String) entity[31]); + result.setRclientn4((String) entity[32]); + result.setAcc_kr_1((String) entity[33]); + result.setAcc_kr_2((String) entity[34]); + result.setSp_code((String) entity[35]); + result.setSpecif_1((String) entity[36]); + result.setSpecif_2((String) entity[37]); + result.setSpecif_3((String) entity[38]); + result.setSpecif_4((String) entity[39]); + result.setSpecif_5((String) entity[40]); + result.setSpecif_6((String) entity[41]); + result.setSend_type((String) entity[42]); + result.setServdate((String) entity[43]); + result.setDoc_result((String) entity[44]); + result.setImp_result((String) entity[45]); + result.setFile_name(filename); + result.setGenerationTime(Instant.now()); + result.setGenerationId(fileId); + return result; + } + +} \ No newline at end of file diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java new file mode 100644 index 000000000..e72cf67f4 --- /dev/null +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -0,0 +1,54 @@ +package ru.spcex.clearing.dbf.importer.logic.stages; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; + +import java.util.HashMap; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +@Component +public class DbfImportKafkaMessenger implements InitializingBean { + private final Supplier kafka; + private final Map> messengers; + + public DbfImportKafkaMessenger(Supplier kafka) { + this.kafka = kafka; + this.messengers = new HashMap<>(); + } + + @Override + public void afterPropertiesSet() { + messengers.put(ETable.DF_01, this::messageDf01); + messengers.put(ETable.DF_04, this::messageDf04); + } + + /** + * отправляет в кафку сообщение, при необходимости + * (обрабатывается, например, в balance-service, clearing-service) + */ + public void notifySystemIfNeeded(ETable table, Long groupId) { + Consumer messenger = messengers.get(table); + if (messenger != null) { + messenger.accept(groupId); + } + } + + private void messageDf01(Long groupId) { + StatementRequest statementRequest = new StatementRequest(); + statementRequest.setGroupId(groupId); + kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + } + + private void messageDf04(Long groupId) { + Sdf04Request sdf04ImportNotification = new Sdf04Request(); + sdf04ImportNotification.setGroupId(groupId); + kafka.get().sendRequestToQueue(Consts.SDF04_PROCESS, sdf04ImportNotification); + } +} 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 d9a5a933e..59f0d63d0 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,9 +8,6 @@ import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer; 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.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; @@ -18,7 +15,6 @@ import java.io.IOException; import java.io.InputStream; import java.nio.charset.Charset; import java.util.Map; -import java.util.function.Supplier; /** * Заливка проверенных данных в базу @@ -28,16 +24,16 @@ public class ImportToDB extends Stage { private final ImportDBFServiceSettings settings; private final HazelcastService hazelcastService; private final Map mappingEnumTableObjectTable; - private final Supplier kafka; + private final DbfImportKafkaMessenger kafkaMessenger; public ImportToDB(ImportDBFServiceSettings settings, HazelcastService hazelcastService, @Qualifier("mapOfTable") Map mappingEnumTableObjectTable, - Supplier kafka) { + DbfImportKafkaMessenger kafkaMessenger) { this.settings = settings; this.hazelcastService = hazelcastService; this.mappingEnumTableObjectTable = mappingEnumTableObjectTable; - this.kafka = kafka; + this.kafkaMessenger = kafkaMessenger; } @Override @@ -57,9 +53,7 @@ public class ImportToDB extends Stage { Object[] entity = dbfReader.nextRecord(); table.injectEntity(table.getEntity(entity)); } - if (ETable.DF_01.equals(currTable)) { - sendStatementRequest(fileId); - } + kafkaMessenger.notifySystemIfNeeded(currTable, fileId); } catch (IOException exception) { log.warn(exception.getMessage()); return StageResult.ERROR; @@ -68,10 +62,4 @@ public class ImportToDB extends Stage { return StageResult.OK; } - - private void sendStatementRequest(Long fileId) { - StatementRequest statementRequest = new StatementRequest(); - statementRequest.setGroupId(fileId); - kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); - } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index fab1a1636..759cb6328 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -30,6 +30,7 @@ public interface Consts { //todo String STATEMENT_PROCESS = "statement-process"; + String SDF04_PROCESS = "sdf04-process"; String EXPORT_PROCESS = "export-process"; String ACCOUNT_NEW = "account-new"; String BALANCE_ACCOUNT_NEW = "balance-account-new"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf04Request.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf04Request.java new file mode 100644 index 000000000..b1a80557c --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/Sdf04Request.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.clearing; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class Sdf04Request { + @JsonProperty + private Long groupId; + + public Long getGroupId() { + return groupId; + } + + public void setGroupId(Long groupId) { + this.groupId = groupId; + } +} \ No newline at end of file