Добавил логирование пар sdf в pairsdf.
This commit is contained in:
parent
67c0599284
commit
ef9f9b3ca2
15 changed files with 400 additions and 109 deletions
|
|
@ -1,6 +1,5 @@
|
||||||
package ru.spcex.clearing.dbf.exporter.logic.stages;
|
package ru.spcex.clearing.dbf.exporter.logic.stages;
|
||||||
|
|
||||||
import com.linuxense.javadbf.DBFField;
|
|
||||||
import com.linuxense.javadbf.DBFWriter;
|
import com.linuxense.javadbf.DBFWriter;
|
||||||
import org.springframework.beans.factory.InitializingBean;
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
|
|
@ -11,6 +10,9 @@ import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer;
|
||||||
import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult;
|
import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult;
|
||||||
import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table;
|
import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table;
|
||||||
import ru.spcex.clearing.dbf.exporter.services.converters.DFConverter;
|
import ru.spcex.clearing.dbf.exporter.services.converters.DFConverter;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
import ru.spcex.platform.enumeration.CurrencyCode;
|
import ru.spcex.platform.enumeration.CurrencyCode;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -19,8 +21,10 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import java.io.File;
|
import java.io.File;
|
||||||
import java.nio.charset.Charset;
|
import java.nio.charset.Charset;
|
||||||
import java.util.*;
|
import java.util.*;
|
||||||
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
import static ru.spcex.clearing.dbf.exporter.services.converters.DFConverter.ruCurrency;
|
import static ru.spcex.clearing.dbf.exporter.services.converters.DFConverter.ruCurrency;
|
||||||
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Выгрузка данных из мапы hazelcast и их запись в файлы
|
* Выгрузка данных из мапы hazelcast и их запись в файлы
|
||||||
|
|
@ -31,19 +35,20 @@ public class ExportFromHazelcast extends Stage implements InitializingBean {
|
||||||
private final ImdgProvider imdgProvider;
|
private final ImdgProvider imdgProvider;
|
||||||
private final SFTPConfig.DbfGateway gateway;
|
private final SFTPConfig.DbfGateway gateway;
|
||||||
private final List<DFConverter> converters;
|
private final List<DFConverter> converters;
|
||||||
|
private final Supplier<KafkaSender> kafkaSender;
|
||||||
|
|
||||||
|
|
||||||
private final Map<Table, DBFField[]> dbfFieldsForTable = new HashMap<>();
|
|
||||||
private Charset dbfCharset;
|
private Charset dbfCharset;
|
||||||
|
|
||||||
public ExportFromHazelcast(ExportDBFServiceSettings settings,
|
public ExportFromHazelcast(ExportDBFServiceSettings settings,
|
||||||
ImdgProvider imdgProvider,
|
ImdgProvider imdgProvider,
|
||||||
SFTPConfig.DbfGateway gateway,
|
SFTPConfig.DbfGateway gateway,
|
||||||
List<DFConverter> converters) {
|
List<DFConverter> converters,
|
||||||
|
Supplier<KafkaSender> kafkaSender) {
|
||||||
this.settings = settings;
|
this.settings = settings;
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.gateway = gateway;
|
this.gateway = gateway;
|
||||||
this.converters = converters;
|
this.converters = converters;
|
||||||
|
this.kafkaSender = kafkaSender;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -110,6 +115,9 @@ public class ExportFromHazelcast extends Stage implements InitializingBean {
|
||||||
} else {
|
} else {
|
||||||
log.debug("uuid {}, send to SFTP path \"{}\"", resultContainer.getUuid(), path);
|
log.debug("uuid {}, send to SFTP path \"{}\"", resultContainer.getUuid(), path);
|
||||||
gateway.sendToSftp(dbfFile, path);
|
gateway.sendToSftp(dbfFile, path);
|
||||||
|
if (Table.S_DF02.equals(resultContainer.getTableForExport())) {
|
||||||
|
sendPairSdfRequest(resultContainer);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -133,4 +141,13 @@ public class ExportFromHazelcast extends Stage implements InitializingBean {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void sendPairSdfRequest(ResultContainer resultContainer) {
|
||||||
|
PairSdfRequest request = new PairSdfRequest();
|
||||||
|
request.setGenerationId(resultContainer.getGroupId());
|
||||||
|
request.setFileNameSDf(resultContainer.getFileForExport().getName());
|
||||||
|
request.setTableSDf(resultContainer.getTableForExport().getFilePrefix());
|
||||||
|
Long msgId = kafkaSender.get().sendRequestToQueue(PAIR_SDF, request);
|
||||||
|
log.info("Send PairSdfRequest={} message id={} to kafka \"{}\"", LogFormatter.toString(request), msgId, PAIR_SDF);
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -11,6 +11,10 @@ public enum ETable {
|
||||||
DF_55("DF-55"),
|
DF_55("DF-55"),
|
||||||
DF_57("DF-57");
|
DF_57("DF-57");
|
||||||
|
|
||||||
|
public String getPrefix() {
|
||||||
|
return prefix;
|
||||||
|
}
|
||||||
|
|
||||||
private final String prefix;
|
private final String prefix;
|
||||||
|
|
||||||
ETable(String prefix) {
|
ETable(String prefix) {
|
||||||
|
|
|
||||||
|
|
@ -96,4 +96,12 @@ public abstract class AbstractTable<T extends SpcexObjectBase> {
|
||||||
}
|
}
|
||||||
throw new IllegalArgumentException("Can not convert to BigDecimal " + value.getClass().getName() + " " + value);
|
throw new IllegalArgumentException("Can not convert to BigDecimal " + value.getClass().getName() + " " + value);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public String getFilename() {
|
||||||
|
return filename;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long getFileId() {
|
||||||
|
return fileId;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,10 +7,12 @@ import org.springframework.stereotype.Component;
|
||||||
import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer;
|
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.ETable;
|
||||||
import ru.spcex.clearing.dbf.importer.logic.data.enums.StageResult;
|
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.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
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.domain.cud.clearing.Sdf04Request;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.platform.enumeration.ObjectType;
|
import ru.spcex.platform.enumeration.ObjectType;
|
||||||
|
|
@ -22,6 +24,8 @@ import java.util.Map;
|
||||||
import java.util.function.Consumer;
|
import java.util.function.Consumer;
|
||||||
import java.util.function.Supplier;
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
|
||||||
|
|
||||||
@Component
|
@Component
|
||||||
public class DbfImportKafkaMessenger implements InitializingBean {
|
public class DbfImportKafkaMessenger implements InitializingBean {
|
||||||
final Logger log = LoggerFactory.getLogger(getClass());
|
final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
@ -103,6 +107,15 @@ public class DbfImportKafkaMessenger implements InitializingBean {
|
||||||
groupId, table, msgId, destination);
|
groupId, table, msgId, destination);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void sendPairSdfRequest(AbstractTable table, String tableSdf) {
|
||||||
|
PairSdfRequest request = new PairSdfRequest();
|
||||||
|
request.setGenerationId(table.getFileId());
|
||||||
|
request.setFileNameSDf(table.getFilename());
|
||||||
|
request.setTableSDf(tableSdf);
|
||||||
|
Long msgId = kafka.get().sendRequestToQueue(PAIR_SDF, request);
|
||||||
|
log.info("Send PairSdfRequest={} message id={} to kafka \"{}\"", LogFormatter.toString(request), msgId, PAIR_SDF);
|
||||||
|
}
|
||||||
|
|
||||||
private void messageDf04(Long groupId) {
|
private void messageDf04(Long groupId) {
|
||||||
Sdf04Request sdf04ImportNotification = new Sdf04Request();
|
Sdf04Request sdf04ImportNotification = new Sdf04Request();
|
||||||
sdf04ImportNotification.setGroupId(groupId);
|
sdf04ImportNotification.setGroupId(groupId);
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,8 @@ import java.io.InputStream;
|
||||||
import java.nio.charset.Charset;
|
import java.nio.charset.Charset;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.dbf.importer.logic.data.enums.ETable.DF_01;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Заливка проверенных данных в базу
|
* Заливка проверенных данных в базу
|
||||||
*/
|
*/
|
||||||
|
|
@ -58,6 +60,9 @@ public class ImportToDB extends Stage {
|
||||||
}
|
}
|
||||||
table.injectEntity(table.getEntity(entity));
|
table.injectEntity(table.getEntity(entity));
|
||||||
}
|
}
|
||||||
|
if (DF_01.equals(currTable)){
|
||||||
|
kafkaMessenger.sendPairSdfRequest(table, currTable.getPrefix());
|
||||||
|
}
|
||||||
kafkaMessenger.notifySystemIfNeeded(currTable, fileId);
|
kafkaMessenger.notifySystemIfNeeded(currTable, fileId);
|
||||||
kafkaMessenger.notifyUserAboutSuccessLoad(resultContainer);
|
kafkaMessenger.notifyUserAboutSuccessLoad(resultContainer);
|
||||||
} catch (IOException exception) {
|
} catch (IOException exception) {
|
||||||
|
|
|
||||||
|
|
@ -54,6 +54,11 @@
|
||||||
<artifactId>spring-boot-starter-test</artifactId>
|
<artifactId>spring-boot-starter-test</artifactId>
|
||||||
<scope>test</scope>
|
<scope>test</scope>
|
||||||
</dependency>
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.clearing</groupId>
|
||||||
|
<artifactId>test-clearing</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
</dependencies>
|
</dependencies>
|
||||||
|
|
||||||
<build>
|
<build>
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,75 @@
|
||||||
|
package ru.spcex.clearing.utility.Enum;
|
||||||
|
|
||||||
|
import ru.clearing.classes.statics.data.sdf.*;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
|
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
public enum Table {
|
||||||
|
//from dbf-importer
|
||||||
|
S_DF01("DF-01", IMDGDistributedNames.Map_SDf01, SDf01.class),
|
||||||
|
S_DF04("DF-04", IMDGDistributedNames.Map_SDf04, SDf04.class),
|
||||||
|
S_DF06("DF-06", IMDGDistributedNames.Map_SDf06, SDf06.class),
|
||||||
|
S_DF52("DF-52", IMDGDistributedNames.Map_SDf52, SDf52.class),
|
||||||
|
S_DF55("DF-55", IMDGDistributedNames.Map_SDf55, SDf55.class),
|
||||||
|
//from dbf-exporter
|
||||||
|
S_DF02("DF-02", IMDGDistributedNames.Map_SDf02, SDf02.class),
|
||||||
|
S_DF03("DF-03", IMDGDistributedNames.Map_SDf03, SDf03.class),
|
||||||
|
S_DF07("DF-07", IMDGDistributedNames.Map_SDf07, SDf07.class),
|
||||||
|
S_DF53("DF-53", IMDGDistributedNames.Map_SDf53, SDf53.class),
|
||||||
|
S_DF54("DF-54", IMDGDistributedNames.Map_SDf54, SDf54.class);
|
||||||
|
|
||||||
|
public static final Map<Table, TableProcessor> outTables = Map.of(
|
||||||
|
S_DF02, new TableProcessor(S_DF01, value -> ((SDf01) value).getGenerationId(), value -> ((SDf02) value).getInSDfId()),
|
||||||
|
S_DF04, new TableProcessor(S_DF03, value -> ((SDf03) value).getGenerationId(), value -> null),
|
||||||
|
S_DF07, new TableProcessor(S_DF06, value -> ((SDf06) value).getGenerationId(), value -> ((SDf07) value).getInSDfId()),
|
||||||
|
S_DF53, new TableProcessor(S_DF52, value -> ((SDf52) value).getGenerationId(), value -> ((SDf53) value).getInSDfId()),
|
||||||
|
S_DF55, new TableProcessor(S_DF54, value -> ((SDf54) value).getGenerationId(), value -> null));
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Имя таблицы
|
||||||
|
*/
|
||||||
|
public final String tableSDf;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Имя таблицы hazelcast в которой будет поиск generationId
|
||||||
|
*/
|
||||||
|
private final String hazelcastMapName;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Класс объекта
|
||||||
|
*/
|
||||||
|
private final Class<? extends SpcexObjectBase> entityClass;
|
||||||
|
|
||||||
|
Table(String tableSDf, String hazelcastMapName, Class<? extends SpcexObjectBase> entityClass) {
|
||||||
|
this.tableSDf = tableSDf;
|
||||||
|
this.hazelcastMapName = hazelcastMapName;
|
||||||
|
this.entityClass = entityClass;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getHazelcastMapName() {
|
||||||
|
return hazelcastMapName;
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean canProcess(String tableSDf) {
|
||||||
|
return this.tableSDf.equals(tableSDf);
|
||||||
|
}
|
||||||
|
|
||||||
|
public Class<? extends SpcexObjectBase> getEntityClass() {
|
||||||
|
return entityClass;
|
||||||
|
}
|
||||||
|
|
||||||
|
public static class TableProcessor {
|
||||||
|
public final Table inTable;
|
||||||
|
public final Function<SpcexObjectBase, Long> inGenerationId;
|
||||||
|
public final Function<SpcexObjectBase, Long> inSdfId;
|
||||||
|
|
||||||
|
public TableProcessor(Table inTable, Function<SpcexObjectBase, Long> inGenerationId, Function<SpcexObjectBase, Long> inSdfId) {
|
||||||
|
this.inTable = inTable;
|
||||||
|
this.inGenerationId = inGenerationId;
|
||||||
|
this.inSdfId = inSdfId;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,113 @@
|
||||||
|
package ru.spcex.clearing.utility.service;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.clearing.classes.statics.data.register.PairSdf;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
import ru.spcex.clearing.utility.Enum.Table;
|
||||||
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import java.time.Instant;
|
||||||
|
import java.util.Arrays;
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.Optional;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
|
||||||
|
import static ru.spcex.clearing.utility.Enum.Table.outTables;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class PairSdfService extends QueueConsumer implements InitializingBean {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final Imdg<PairSdf> pairSdfImdg;
|
||||||
|
private final ImdgProvider imdgProvider;
|
||||||
|
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public PairSdfService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
|
||||||
|
ImdgProvider imdgProvider) {
|
||||||
|
super(kafkaQueue, kafkaProducer);
|
||||||
|
this.pairSdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PairSdf, PairSdf.class);
|
||||||
|
this.imdgProvider = imdgProvider;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
callback(PairSdfRequest.class)
|
||||||
|
.setConsumer(this::processRequest)
|
||||||
|
.forDestination(PAIR_SDF, callbacks::put);
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
|
||||||
|
public void processRequest(BaseRequest<PairSdfRequest> pairSdfNewRequestBaseRequest) {
|
||||||
|
log.info("Starting BaseRequest<PairSdfRequest>={} processing...", LogFormatter.toString(pairSdfNewRequestBaseRequest));
|
||||||
|
PairSdfRequest request = pairSdfNewRequestBaseRequest.getRequestPayload();
|
||||||
|
Optional<Table> table = Arrays.stream(Table.values()).filter(t -> t.canProcess(request.getTableSDf())).findFirst();
|
||||||
|
if (table.isPresent()) {
|
||||||
|
if (outTables.containsKey(table.get())) {
|
||||||
|
processOut(request, table.get());
|
||||||
|
} else {
|
||||||
|
processIn(request);
|
||||||
|
}
|
||||||
|
log.info("Request with BaseRequest<PairSdfRequest>.id={} successfully processed!", pairSdfNewRequestBaseRequest.getId());
|
||||||
|
} else
|
||||||
|
log.info("Request with BaseRequest<PairSdfRequest>.id={} not processed, because not supported table = {}!", pairSdfNewRequestBaseRequest.getId(), request.getTableSDf());
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
private void processIn(PairSdfRequest request) {
|
||||||
|
log.info("Process PairSdfRequest={} (by IN)!", LogFormatter.toString(request));
|
||||||
|
Instant currDt = Instant.now();
|
||||||
|
PairSdf pairSdf = new PairSdf();
|
||||||
|
pairSdf.setCreated(currDt);
|
||||||
|
pairSdf.setUpdated(currDt);
|
||||||
|
pairSdf.setInSDfId(String.valueOf(request.getGenerationId()));
|
||||||
|
pairSdf.setInSDf(request.getFileNameSDf());
|
||||||
|
pairSdfImdg.insert(pairSdf);
|
||||||
|
log.info("Successfully inserted PairSdf={} in map!", LogFormatter.toString(pairSdf));
|
||||||
|
}
|
||||||
|
|
||||||
|
private void processOut(PairSdfRequest request, Table outTable) {
|
||||||
|
log.info("Process PairSdfRequest={} (by OUT)!", LogFormatter.toString(request));
|
||||||
|
Table.TableProcessor processor = outTables.get(outTable);
|
||||||
|
PairSdf pairSdf = null;
|
||||||
|
// out generationId -> outSdf inSdfId -> inSdf by id -> inSdf generationId -> pairSdf by in_s_df_id
|
||||||
|
if (request.getGenerationId() != null) {
|
||||||
|
Imdg<? extends SpcexObjectBase> outMap = imdgProvider.getImdg(outTable.getHazelcastMapName(), outTable.getEntityClass());
|
||||||
|
SpcexObjectBase outSdf = outMap.getFirstObjectByFieldValues(Map.of("generationId", request.getGenerationId()));
|
||||||
|
if (outSdf != null && processor.inSdfId.apply(outSdf) != null) {
|
||||||
|
Imdg<? extends SpcexObjectBase> inMap = imdgProvider.getImdg(processor.inTable.getHazelcastMapName(), processor.inTable.getEntityClass());
|
||||||
|
SpcexObjectBase inSdf = inMap.getSingleObjectByID(processor.inSdfId.apply(outSdf));
|
||||||
|
if (inSdf != null && processor.inGenerationId.apply(inSdf) != null) {
|
||||||
|
pairSdf = pairSdfImdg.getFirstObjectByFieldValues(Map.of("inSDfId", processor.inGenerationId.apply(inSdf)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Instant currDt = Instant.now();
|
||||||
|
if (pairSdf != null) {
|
||||||
|
pairSdf.setUpdated(currDt);
|
||||||
|
pairSdf.setOutSDfId(String.valueOf(request.getGenerationId()));
|
||||||
|
pairSdf.setOutSDf(request.getFileNameSDf());
|
||||||
|
pairSdfImdg.update(pairSdf);
|
||||||
|
log.info("Successfully updated PairSdf {} in map!", LogFormatter.toString(pairSdf));
|
||||||
|
} else {
|
||||||
|
pairSdf = new PairSdf();
|
||||||
|
pairSdf.setCreated(currDt);
|
||||||
|
pairSdf.setUpdated(currDt);
|
||||||
|
pairSdf.setOutSDfId(String.valueOf(request.getGenerationId()));
|
||||||
|
pairSdf.setOutSDf(request.getFileNameSDf());
|
||||||
|
pairSdfImdg.insert(pairSdf);
|
||||||
|
log.info("Not found PairSdf(IN), inserted PairSdf={} in map!", LogFormatter.toString(pairSdf));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -1,23 +0,0 @@
|
||||||
package ru.spcex.clearing.utility.config;
|
|
||||||
|
|
||||||
import com.hazelcast.client.HazelcastClient;
|
|
||||||
import com.hazelcast.config.Config;
|
|
||||||
import com.hazelcast.config.JoinConfig;
|
|
||||||
import com.hazelcast.config.MulticastConfig;
|
|
||||||
import com.hazelcast.config.NetworkConfig;
|
|
||||||
import com.hazelcast.core.Hazelcast;
|
|
||||||
import com.hazelcast.core.HazelcastInstance;
|
|
||||||
import org.springframework.context.annotation.Bean;
|
|
||||||
import org.springframework.context.annotation.Configuration;
|
|
||||||
|
|
||||||
import java.util.Random;
|
|
||||||
|
|
||||||
@Configuration
|
|
||||||
public class HazelcastInstanceTestConfiguration {
|
|
||||||
|
|
||||||
@Bean(name = "hazelcastInstance")
|
|
||||||
public HazelcastInstance hazelcastInstance() {
|
|
||||||
return HazelcastClient.newHazelcastClient();
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
@ -1,73 +0,0 @@
|
||||||
package ru.spcex.clearing.utility.config;
|
|
||||||
|
|
||||||
import com.hazelcast.config.*;
|
|
||||||
import com.hazelcast.core.Hazelcast;
|
|
||||||
import com.hazelcast.core.HazelcastInstance;
|
|
||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
|
||||||
import org.springframework.context.annotation.Bean;
|
|
||||||
import org.springframework.context.annotation.Configuration;
|
|
||||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
|
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.util.HazelcastHelper;
|
|
||||||
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Random;
|
|
||||||
|
|
||||||
@Configuration
|
|
||||||
public class HazelcastServiceTestConfiguration {
|
|
||||||
private HazelcastInstance hazelcastInstance;
|
|
||||||
|
|
||||||
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
|
|
||||||
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
|
|
||||||
if (maxPoolSz > 2) {
|
|
||||||
pool.setKeepAliveSeconds(60);
|
|
||||||
pool.setAllowCoreThreadTimeOut(true);
|
|
||||||
}
|
|
||||||
pool.setCorePoolSize(maxPoolSz);
|
|
||||||
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
|
|
||||||
return pool;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Bean(name = "hazelcastServiceTest")
|
|
||||||
public HazelcastService hazelcastService(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, 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.newHazelcastInstance(cfg);
|
|
||||||
|
|
||||||
HazelcastHelper.imdgSystem_setStorageState(true, hazelcastInstance);
|
|
||||||
return new HazelcastService(taskExecutorHazelcastClientInitializer,
|
|
||||||
taskExecutorIdGeneratorAwaiter, params);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Bean(name = "taskExecutorHazelcastClientInitializer")
|
|
||||||
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
|
|
||||||
return createThreadPoolTaskExecutor(1, true);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Bean(name = "taskExecutorIdGeneratorAwaiter")
|
|
||||||
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
|
|
||||||
return createThreadPoolTaskExecutor(1, false);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Bean(name = "hazelcastClientParams")
|
|
||||||
public HazelcastClientParams getHazelcastClientParams() {
|
|
||||||
HazelcastClientParams params = new HazelcastClientParams();
|
|
||||||
params.setLogin("dev");
|
|
||||||
params.setPassword("dev-pass");
|
|
||||||
params.setClusterMembers("127.0.0.1");
|
|
||||||
params.setInstanceName("hzTestClient" + new Random().nextInt());
|
|
||||||
params.setNearCacheConfig(new NearCacheConfig());
|
|
||||||
|
|
||||||
return params;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -0,0 +1,16 @@
|
||||||
|
package ru.spcex.clearing.utility.config;
|
||||||
|
|
||||||
|
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
||||||
|
import org.springframework.boot.test.context.TestConfiguration;
|
||||||
|
import org.springframework.context.annotation.ComponentScan;
|
||||||
|
import org.springframework.context.annotation.FilterType;
|
||||||
|
import ru.spcex.clearing.utility.config.settings.UtilityServiceSettings;
|
||||||
|
|
||||||
|
@TestConfiguration
|
||||||
|
@ComponentScan(basePackages = {"ru.spcex.clearing.utility.config",
|
||||||
|
"ru.spcex.clearing.utility.service"},
|
||||||
|
excludeFilters = {@ComponentScan.Filter(type = FilterType.ASSIGNABLE_TYPE, value = KafkaConfig.class),
|
||||||
|
@ComponentScan.Filter(type = FilterType.ASSIGNABLE_TYPE, value = UtilityServiceImdgConfig.class)})
|
||||||
|
@EnableConfigurationProperties(value = UtilityServiceSettings.class)
|
||||||
|
public class TestConfig {
|
||||||
|
}
|
||||||
|
|
@ -15,8 +15,6 @@ import org.apache.kafka.common.TopicPartition;
|
||||||
import org.junit.jupiter.api.Assertions;
|
import org.junit.jupiter.api.Assertions;
|
||||||
import org.junit.jupiter.api.extension.ExtendWith;
|
import org.junit.jupiter.api.extension.ExtendWith;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
|
||||||
import org.springframework.test.context.ContextConfiguration;
|
|
||||||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||||
import ru.clearing.classes.statics.data.misc.KeyRate;
|
import ru.clearing.classes.statics.data.misc.KeyRate;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
|
@ -26,8 +24,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest;
|
||||||
import ru.spcex.clearing.utility.config.HazelcastInstanceTestConfiguration;
|
|
||||||
import ru.spcex.clearing.utility.config.HazelcastServiceTestConfiguration;
|
|
||||||
import ru.spcex.clearing.utility.utils.MatcherFactory;
|
import ru.spcex.clearing.utility.utils.MatcherFactory;
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
||||||
|
|
||||||
|
|
@ -37,9 +33,6 @@ import java.util.Collections;
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
|
|
||||||
@ExtendWith(SpringExtension.class)
|
@ExtendWith(SpringExtension.class)
|
||||||
@ContextConfiguration(classes = {
|
|
||||||
HazelcastServiceTestConfiguration.class,
|
|
||||||
HazelcastInstanceTestConfiguration.class})
|
|
||||||
public class KeyRateServiceDisabled {
|
public class KeyRateServiceDisabled {
|
||||||
static final ObjectMapper objectMapper = new ObjectMapper();
|
static final ObjectMapper objectMapper = new ObjectMapper();
|
||||||
private final int TIMEOUT = 1000;
|
private final int TIMEOUT = 1000;
|
||||||
|
|
@ -47,10 +40,8 @@ public class KeyRateServiceDisabled {
|
||||||
private MockConsumer<String, Object> kafkaMockQueue;
|
private MockConsumer<String, Object> kafkaMockQueue;
|
||||||
private MockProducer<String, Object> mockProducer = new MockProducer<>();
|
private MockProducer<String, Object> mockProducer = new MockProducer<>();
|
||||||
@Autowired
|
@Autowired
|
||||||
@Qualifier("hazelcastServiceTest")
|
|
||||||
private HazelcastService hazelcastService;
|
private HazelcastService hazelcastService;
|
||||||
@Autowired
|
@Autowired
|
||||||
@Qualifier("hazelcastInstance")
|
|
||||||
private HazelcastInstance hz;
|
private HazelcastInstance hz;
|
||||||
|
|
||||||
private void setIsUsed(boolean isUsed) {
|
private void setIsUsed(boolean isUsed) {
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,103 @@
|
||||||
|
package ru.spcex.clearing.utility.service;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
import org.junit.jupiter.api.Disabled;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.junit.jupiter.api.extension.ExtendWith;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.test.context.ContextConfiguration;
|
||||||
|
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||||
|
import ru.clearing.classes.statics.data.register.PairSdf;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf01;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf02;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.PairSdfRequest;
|
||||||
|
import ru.spcex.clearing.test.config.ImdgTestConfig;
|
||||||
|
import ru.spcex.clearing.test.config.KafkaTestConfig;
|
||||||
|
import ru.spcex.clearing.utility.config.TestConfig;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import javax.annotation.PostConstruct;
|
||||||
|
import java.io.File;
|
||||||
|
import java.nio.file.Path;
|
||||||
|
import java.nio.file.Paths;
|
||||||
|
|
||||||
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
|
||||||
|
import static ru.spcex.clearing.test.TestUtils.*;
|
||||||
|
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
|
||||||
|
import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId;
|
||||||
|
import static ru.spcex.clearing.utility.Enum.Table.S_DF01;
|
||||||
|
import static ru.spcex.clearing.utility.Enum.Table.S_DF02;
|
||||||
|
|
||||||
|
@ExtendWith(SpringExtension.class)
|
||||||
|
@ContextConfiguration(classes = {TestConfig.class,
|
||||||
|
ImdgTestConfig.class,
|
||||||
|
KafkaTestConfig.class})
|
||||||
|
class PairSdfServiceTest {
|
||||||
|
private static final Long inGenerationId = currentID.getAndIncrement();
|
||||||
|
private static final Long outGenerationId = currentID.getAndIncrement();
|
||||||
|
private static final Long inSdfId = currentID.getAndIncrement();
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
private PairSdfService pairSdfService;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Qualifier("hazelcastServiceTest")
|
||||||
|
protected ImdgProvider imdgProvider;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Qualifier("mockProducer")
|
||||||
|
protected Producer<String, Object> mockProducer;
|
||||||
|
|
||||||
|
private Imdg<PairSdf> pairSdfImdg;
|
||||||
|
private Imdg<SDf01> sDf01Imdg;
|
||||||
|
private Imdg<SDf02> sDf02Imdg;
|
||||||
|
|
||||||
|
static {
|
||||||
|
Path path = Paths.get("src", "main", "resources");
|
||||||
|
String currentPath = path.toAbsolutePath().toString();
|
||||||
|
System.setProperty("spring.config.location", currentPath + File.separator);
|
||||||
|
}
|
||||||
|
|
||||||
|
@PostConstruct
|
||||||
|
public void init(){
|
||||||
|
waitAvailableImdgProviderAndAddAdminWithDefaultId();
|
||||||
|
this.sDf01Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
||||||
|
SDf01 sDf01 = new SDf01();
|
||||||
|
sDf01.setId(inSdfId);
|
||||||
|
sDf01.setGenerationId(inGenerationId);
|
||||||
|
sDf01Imdg.insert(sDf01);
|
||||||
|
this.sDf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
||||||
|
SDf02 sDf02 = new SDf02();
|
||||||
|
sDf02.setGenerationId(outGenerationId);
|
||||||
|
sDf02.setInSDfId(inSdfId);
|
||||||
|
sDf02Imdg.insert(sDf02);
|
||||||
|
this.pairSdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PairSdf, PairSdf.class);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Disabled//Для отладки
|
||||||
|
@Test
|
||||||
|
void processRequest() throws InterruptedException {
|
||||||
|
Long ID = currentID.getAndIncrement();
|
||||||
|
|
||||||
|
PairSdfRequest pairSdfRequest = new PairSdfRequest();
|
||||||
|
pairSdfRequest.setFileNameSDf("S_DF01_file_name");
|
||||||
|
pairSdfRequest.setGenerationId(inGenerationId);
|
||||||
|
pairSdfRequest.setTableSDf(S_DF01.tableSDf);
|
||||||
|
|
||||||
|
String jsonString = getJsonStringForNew(pairSdfRequest, ID);
|
||||||
|
addRecordToKafka((MockConsumer) pairSdfService.getConsumer(), PAIR_SDF, 0, 0, jsonString);
|
||||||
|
waitingSendAndCheckRecord(ID, mockProducer);
|
||||||
|
|
||||||
|
pairSdfRequest.setFileNameSDf("S_DF02_file_name");
|
||||||
|
pairSdfRequest.setGenerationId(outGenerationId);
|
||||||
|
pairSdfRequest.setTableSDf(S_DF02.tableSDf);
|
||||||
|
|
||||||
|
jsonString = getJsonStringForNew(pairSdfRequest, ID);
|
||||||
|
addRecordToKafka((MockConsumer) pairSdfService.getConsumer(), PAIR_SDF, 0, 1, jsonString);
|
||||||
|
waitingSendAndCheckRecord(ID, mockProducer);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -115,6 +115,7 @@ public interface Consts {
|
||||||
String ACCOUNT_NOTIFICATION_FEEDBACK = "account-notification-feedback";
|
String ACCOUNT_NOTIFICATION_FEEDBACK = "account-notification-feedback";
|
||||||
String NOTIFICATION_NEW = "notification-new";
|
String NOTIFICATION_NEW = "notification-new";
|
||||||
String NOTIFICATION_UPDATE = "notification-update";
|
String NOTIFICATION_UPDATE = "notification-update";
|
||||||
|
String PAIR_SDF = "pair-sdf";
|
||||||
|
|
||||||
|
|
||||||
String DESTINATION_SDF08_NEW = "s-df-08-new";
|
String DESTINATION_SDF08_NEW = "s-df-08-new";
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.utilities;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
|
public class PairSdfRequest {
|
||||||
|
@JsonProperty
|
||||||
|
public Long generationId;
|
||||||
|
@JsonProperty
|
||||||
|
public String fileNameSDf;
|
||||||
|
@JsonProperty
|
||||||
|
public String tableSDf;
|
||||||
|
|
||||||
|
public Long getGenerationId() {
|
||||||
|
return generationId;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setGenerationId(Long generationId) {
|
||||||
|
this.generationId = generationId;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getFileNameSDf() {
|
||||||
|
return fileNameSDf;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setFileNameSDf(String fileNameSDf) {
|
||||||
|
this.fileNameSDf = fileNameSDf;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getTableSDf() {
|
||||||
|
return tableSDf;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setTableSDf(String tableSDf) {
|
||||||
|
this.tableSDf = tableSDf;
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue