This commit is contained in:
ialbert 2022-11-25 12:58:10 +03:00
parent 59e748b18a
commit e463d83084
7 changed files with 151 additions and 16 deletions

View file

@ -57,6 +57,7 @@ public class DBFImporterConfig {
public Map<ETable, AbstractTable> getMapOfTables() {
Map<ETable, AbstractTable> 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());

View file

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

View file

@ -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<SDf04> {
private static final String PREFIX = ETable.DF_04.name();
private static final Class<SDf04> 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;
}
}

View file

@ -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<KafkaSender> kafka;
private final Map<ETable, Consumer<Long>> messengers;
public DbfImportKafkaMessenger(Supplier<KafkaSender> 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<Long> 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);
}
}

View file

@ -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<ETable, AbstractTable> mappingEnumTableObjectTable;
private final Supplier<KafkaSender> kafka;
private final DbfImportKafkaMessenger kafkaMessenger;
public ImportToDB(ImportDBFServiceSettings settings,
HazelcastService hazelcastService,
@Qualifier("mapOfTable") Map<ETable, AbstractTable> mappingEnumTableObjectTable,
Supplier<KafkaSender> 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);
}
}

View file

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

View file

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