swt-exporter http://jira.mfd.msk:8088/browse/CLS-317 почти готов, но формат файла не понял
This commit is contained in:
parent
cf2b937cff
commit
adcb2e7a91
17 changed files with 449 additions and 287 deletions
|
|
@ -1,32 +0,0 @@
|
|||
package ru.spcex.clearing.swt.exporter.config.settings;
|
||||
|
||||
public class Common {
|
||||
|
||||
private String encoding;
|
||||
private int insertBatchSize;
|
||||
private int threadsCount;
|
||||
|
||||
public String getEncoding() {
|
||||
return encoding;
|
||||
}
|
||||
|
||||
public void setEncoding(String encoding) {
|
||||
this.encoding = encoding;
|
||||
}
|
||||
|
||||
public int getInsertBatchSize() {
|
||||
return insertBatchSize;
|
||||
}
|
||||
|
||||
public void setInsertBatchSize(int insertBatchSize) {
|
||||
this.insertBatchSize = insertBatchSize;
|
||||
}
|
||||
|
||||
public int getThreadsCount() {
|
||||
return threadsCount;
|
||||
}
|
||||
|
||||
public void setThreadsCount(int threadsCount) {
|
||||
this.threadsCount = threadsCount;
|
||||
}
|
||||
}
|
||||
|
|
@ -14,8 +14,10 @@ public class ExportSwtServiceSettings {
|
|||
private HazelcastClientParams hazelcast;
|
||||
private KafkaConsumerSettings kafkaConsumer;
|
||||
private KafkaProducerSettings kafkaProducer;
|
||||
private Store sFTPStore;
|
||||
private Common common;
|
||||
|
||||
private String docOut;
|
||||
//todo ??? private Long interval
|
||||
|
||||
|
||||
public HazelcastClientParams getHazelcast() {
|
||||
return hazelcast;
|
||||
|
|
@ -33,14 +35,6 @@ public class ExportSwtServiceSettings {
|
|||
this.kafkaConsumer = kafkaConsumer;
|
||||
}
|
||||
|
||||
public Store getStore() {
|
||||
return sFTPStore;
|
||||
}
|
||||
|
||||
public void setStore(Store Store) {
|
||||
this.sFTPStore = Store;
|
||||
}
|
||||
|
||||
public KafkaProducerSettings getKafkaProducer() {
|
||||
return kafkaProducer;
|
||||
}
|
||||
|
|
@ -49,11 +43,11 @@ public class ExportSwtServiceSettings {
|
|||
this.kafkaProducer = kafkaProducer;
|
||||
}
|
||||
|
||||
public Common getCommon() {
|
||||
return common;
|
||||
public String getDocOut() {
|
||||
return docOut;
|
||||
}
|
||||
|
||||
public void setCommon(Common common) {
|
||||
this.common = common;
|
||||
public void setDocOut(String docOut) {
|
||||
this.docOut = docOut;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,50 +0,0 @@
|
|||
package ru.spcex.clearing.swt.exporter.config.settings;
|
||||
|
||||
public class Store {
|
||||
|
||||
private String outDir;
|
||||
private String user;
|
||||
private String password;
|
||||
private String serverIp;
|
||||
private int serverPort;
|
||||
|
||||
public String getUser() {
|
||||
return user;
|
||||
}
|
||||
|
||||
public void setUser(String user) {
|
||||
this.user = user;
|
||||
}
|
||||
|
||||
public String getPassword() {
|
||||
return password;
|
||||
}
|
||||
|
||||
public void setPassword(String password) {
|
||||
this.password = password;
|
||||
}
|
||||
|
||||
public String getServerIp() {
|
||||
return serverIp;
|
||||
}
|
||||
|
||||
public void setServerIp(String serverIp) {
|
||||
this.serverIp = serverIp;
|
||||
}
|
||||
|
||||
public int getServerPort() {
|
||||
return serverPort;
|
||||
}
|
||||
|
||||
public void setServerPort(int serverPort) {
|
||||
this.serverPort = serverPort;
|
||||
}
|
||||
|
||||
public String getOutDir() {
|
||||
return outDir;
|
||||
}
|
||||
|
||||
public void setOutDir(String outDir) {
|
||||
this.outDir = outDir;
|
||||
}
|
||||
}
|
||||
|
|
@ -2,34 +2,28 @@ package ru.spcex.clearing.swt.exporter.services;
|
|||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import ru.clearing.classes.statics.data.registry.Registry;
|
||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.JournalEventExportedRequest;
|
||||
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.enumeration.ResultStatuses;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.io.*;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.*;
|
||||
|
||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.LIM_EXPORTED;
|
||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.JOURNAL_SERVICE;
|
||||
|
||||
public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
protected final Logger log = LoggerFactory.getLogger(getClass());
|
||||
protected final SwtTable type;
|
||||
protected final Imdg<T> sdfImdg;
|
||||
private final DateTimeFormatter dtFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
|
||||
private final DateTimeFormatter dtFileNameFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
|
||||
private final DateTimeFormatter dtInFileHeaderFormatter = DateTimeFormatter.ofPattern("yyyyMMdd'/'HHmm");
|
||||
private final KafkaSender kafkaSender;
|
||||
protected final FileStorage fileStorage;
|
||||
|
||||
|
|
@ -48,59 +42,113 @@ public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
|||
return type;
|
||||
}
|
||||
|
||||
protected abstract String typeForFileName();
|
||||
|
||||
protected abstract String sectionForFileName();
|
||||
|
||||
protected abstract String typeForHeader();
|
||||
|
||||
public void process() {
|
||||
String fileName = getTargetFileName();
|
||||
LocalDateTime exportAt = LocalDateTime.now();
|
||||
|
||||
String fileName = formatFileName(typeForFileName(), sectionForFileName(), exportAt);
|
||||
log.debug("Start export {} Lim file", fileName);
|
||||
|
||||
try {
|
||||
byte[] data;
|
||||
{
|
||||
ByteArrayOutputStream outBuffer=new ByteArrayOutputStream();
|
||||
Collection<T> records = selectItems();log.debug("Prepared {} record from {} to file {}",
|
||||
records.size(), sdfImdg.getMapName(), fileName);
|
||||
makeSWTData(records, outBuffer);
|
||||
data= outBuffer.toByteArray();
|
||||
ByteArrayOutputStream outBuffer = new ByteArrayOutputStream();
|
||||
Collection<T> records = selectItems();
|
||||
log.debug("Prepared {} record from {} to file {}",
|
||||
records.size(), sdfImdg.getMapName(), fileName);
|
||||
makeSWTData(typeForHeader(), exportAt, records, outBuffer);
|
||||
data = outBuffer.toByteArray();
|
||||
}
|
||||
try {
|
||||
fileStorage.saveFile(fileName, data);
|
||||
} catch (IOException e) {
|
||||
} catch (Exception e) { // IOException, ...
|
||||
log.error("Failed export {} file", fileName);
|
||||
sendSwtExportedNotification(exportAt, null, ResultStatuses.notSuccess);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
log.debug("Successfully exported {} file", fileName);
|
||||
|
||||
sendSwtExportedNotification(fileName);
|
||||
sendSwtExportedNotification(exportAt, null, ResultStatuses.success);
|
||||
}
|
||||
|
||||
void sendSwtExportedNotification(String fileName) {
|
||||
LimExportedRequest limExportedRequest = new LimExportedRequest();
|
||||
limExportedRequest.setLimFileName(fileName);
|
||||
log.debug("Send message to kafka \"{}\": {}", LIM_EXPORTED, LogFormatter.toStringWrapper(limExportedRequest));
|
||||
kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest);//todo rewrite!!!
|
||||
protected abstract String getDocumentNameForJournal();
|
||||
|
||||
void sendSwtExportedNotification(LocalDateTime registrationAt, Long registrationNumber, ResultStatuses resultStatus) {
|
||||
JournalEventExportedRequest exportedRequest = new JournalEventExportedRequest();
|
||||
exportedRequest.setRegistratoinDate(registrationAt.toLocalDate());
|
||||
exportedRequest.setRegistrationTime(registrationAt.toLocalTime());
|
||||
exportedRequest.setRegistrationNumber(registrationNumber);
|
||||
exportedRequest.setDocumentName(getDocumentNameForJournal());
|
||||
exportedRequest.setResultStatus(resultStatus.getKey());
|
||||
log.debug("Send message to kafka \"{}\": {}", JOURNAL_SERVICE, LogFormatter.toStringWrapper(exportedRequest));
|
||||
kafkaSender.sendRequestToQueue(JOURNAL_SERVICE, exportedRequest);
|
||||
}
|
||||
|
||||
protected String prepareFileName(Long counter, String target) {
|
||||
String dt = dtFormatter.format(LocalDateTime.now());
|
||||
String type="09";
|
||||
String section="U";
|
||||
String counterS = counter==null?"":"_"+counter;
|
||||
String partyCode="";
|
||||
String result="KS_RDC_DF-%s_%s_PRC%s%s%s.swt".formatted(type, section,dt,counterS,partyCode);
|
||||
return result;
|
||||
|
||||
/**
|
||||
* @param type DF-09
|
||||
* @param section bond/fund/""
|
||||
* @param atTime LocalDateTime.now(), если null - текущее время
|
||||
* @return пример "KS_RDC_DF-12_bond_220907151804503.txt"
|
||||
*/
|
||||
protected String formatFileName(String type, String section, LocalDateTime atTime) {
|
||||
Objects.requireNonNull(type);
|
||||
if (section == null) section = "";
|
||||
if (section.length() > 0) section += "_";
|
||||
if (atTime == null) atTime = LocalDateTime.now();
|
||||
String dt = dtFileNameFormatter.format(atTime);
|
||||
return String.format("KS_RDC_%s_%s%s.txt", type, section, dt);
|
||||
}
|
||||
|
||||
// Выборка
|
||||
protected Collection<T> selectItems() {
|
||||
//todo select criteria?
|
||||
return sdfImdg.getAllValues();
|
||||
}
|
||||
|
||||
// Конвертация (см. meta.xml)
|
||||
protected abstract Map<String, Object> convertRecord(T record);
|
||||
protected abstract String[] swtHeader();
|
||||
// Конвертация (поля см. meta.xml)
|
||||
protected abstract LinkedHashMap<String, Object> convertRecord(T record);
|
||||
|
||||
protected void makeSWTData(Collection<T> records, OutputStream out) {
|
||||
//todo ... header + validation + check + convert types
|
||||
// protected abstract String[] swtHeader();
|
||||
|
||||
protected void makeSWTData(String type, LocalDateTime time, Collection<T> records, OutputStream outStream) {
|
||||
PrintWriter out = new PrintWriter(outStream);
|
||||
// todo SWT txt не понял формат. Надо уточнить формат файла. Должен соответствовать мете.
|
||||
out.println("To:CSO");
|
||||
out.println("From:SPCE");
|
||||
if (type != null)
|
||||
out.println("Type:" + type);// Type:009
|
||||
if (time != null) {
|
||||
String timeS = dtInFileHeaderFormatter.format(time);
|
||||
out.println("Date/Time:" + timeS);// Date/Time:20230227/0932
|
||||
}
|
||||
/*
|
||||
To:CSO
|
||||
From:SPCE
|
||||
Type:009
|
||||
Date/Time:20230227/0932
|
||||
:20:0ef63e17-83ac-4e22-a76b-fd8a4df10de3
|
||||
:21:SDC230227084839
|
||||
:18A:46
|
||||
*/
|
||||
StringBuilder line = new StringBuilder();
|
||||
for (T row : records) {
|
||||
LinkedHashMap<String, Object> rowData = convertRecord(row);
|
||||
line.setLength(0);
|
||||
for (Map.Entry<String, Object> r : rowData.entrySet()) {
|
||||
String value = convertItem(r.getValue());
|
||||
line.append(value).append(':');
|
||||
}
|
||||
if (line.length() > 0) // remove :
|
||||
line.setLength(line.length() - 1);
|
||||
out.println(line);
|
||||
}
|
||||
//out.println("2"); // todo что значит 2?
|
||||
}
|
||||
|
||||
protected String convertItem(Object o) {
|
||||
//todo date/time/etc.
|
||||
return String.valueOf(o);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.apache.commons.io.FileUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
|
|
@ -11,15 +12,15 @@ import java.io.IOException;
|
|||
|
||||
@Service
|
||||
public class FileStorage {
|
||||
protected final Logger log= LoggerFactory.getLogger(getClass());
|
||||
protected final Logger log = LoggerFactory.getLogger(getClass());
|
||||
protected File outPath;
|
||||
|
||||
@Autowired
|
||||
public FileStorage(ExportSwtServiceSettings config) {
|
||||
if (config.getStore().getOutDir()==null || config.getStore().getOutDir().isBlank()) {
|
||||
if (config.getDocOut() == null || config.getDocOut().isBlank()) {
|
||||
throw new IllegalArgumentException("Out directory settings is empty.");
|
||||
}
|
||||
this.outPath = new File(config.getStore().getOutDir()); //todo ...
|
||||
this.outPath = new File(config.getDocOut());
|
||||
if (!outPath.isDirectory()) {
|
||||
log.info("Path not exist. mkdir \"{}\"", outPath.getAbsolutePath());
|
||||
if (!outPath.mkdir()) {
|
||||
|
|
@ -30,6 +31,7 @@ public class FileStorage {
|
|||
}
|
||||
|
||||
public void saveFile(String fileName, byte[] data) throws IOException {
|
||||
//todo ...
|
||||
File toFile = new File(outPath, fileName);
|
||||
FileUtils.writeByteArrayToFile(toFile, data);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,8 +10,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||
|
||||
import java.util.List;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,62 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf09;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
|
||||
@Service
|
||||
public class DF09Exporter extends AbstractExporterService<SDf09> {
|
||||
|
||||
public DF09Exporter(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider,
|
||||
SwtTable.SDF_09,
|
||||
IMDGDistributedNames.Map_SDf09, SDf09.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForFileName() {
|
||||
return "DF-09";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String sectionForFileName() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForHeader() {
|
||||
return "009";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String getDocumentNameForJournal() {
|
||||
return "Уведомление об исполнении операции загрузки ценных бумаг или уведомление об ошибке";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected LinkedHashMap<String, Object> convertRecord(SDf09 record) {
|
||||
LinkedHashMap<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("id", record.getId());
|
||||
|
||||
row.put("outDocument", record.getOutDocument());
|
||||
row.put("inDocument", record.getInDocument());
|
||||
row.put("depoCode", record.getDepoCode());
|
||||
row.put("quantity", record.getQuantity());
|
||||
row.put("securityCode", record.getSecurityCode());
|
||||
row.put("clientName", record.getClientName());
|
||||
row.put("result", record.getResult());
|
||||
row.put("generationTime", record.getGenerationTime());
|
||||
row.put("generationId", record.getGenerationId());
|
||||
return row;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,63 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf11;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
|
||||
|
||||
@Service
|
||||
public class DF11Exporter extends AbstractExporterService<SDf11> {
|
||||
|
||||
public DF11Exporter(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider,
|
||||
SwtTable.SDF_11,
|
||||
IMDGDistributedNames.Map_SDf11, SDf11.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForFileName() {
|
||||
return "DF-11";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String sectionForFileName() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForHeader() {
|
||||
return "011";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String getDocumentNameForJournal() {
|
||||
return "Ответ на Запрос на Зачисление или списание ценных бумаг";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected LinkedHashMap<String, Object> convertRecord(SDf11 record) {
|
||||
LinkedHashMap<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("id", record.getId());
|
||||
|
||||
row.put("outDocument", record.getOutDocument());
|
||||
row.put("inDocument", record.getInDocument());
|
||||
row.put("depoCode", record.getDepoCode());
|
||||
row.put("quantity", record.getQuantity());
|
||||
row.put("securityCode", record.getSecurityCode());
|
||||
row.put("clientName", record.getClientName());
|
||||
row.put("result", record.getResult());
|
||||
row.put("generationTime", record.getGenerationTime());
|
||||
row.put("generationId", record.getGenerationId());
|
||||
return row;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,66 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf12;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
|
||||
|
||||
@Service
|
||||
public class DF12Exporter extends AbstractExporterService<SDf12> {
|
||||
|
||||
public DF12Exporter(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider,
|
||||
SwtTable.SDF_12,
|
||||
IMDGDistributedNames.Map_SDf12, SDf12.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForFileName() {
|
||||
return "DF-12";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String sectionForFileName() {
|
||||
return "bond"; // todo bond / fund
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForHeader() {
|
||||
return "012";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String getDocumentNameForJournal() {
|
||||
//todo Выбор:
|
||||
// Распоряжение на проведение операций по итогам клиринга (Фондовая секция)
|
||||
//или
|
||||
// Распоряжение на проведение операций по итогам клиринга (ОФЗ, ОБР)
|
||||
return "Распоряжение на проведение операций по итогам клиринга (Фондовая секция / ОФЗ, ОБР)";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected LinkedHashMap<String, Object> convertRecord(SDf12 record) {
|
||||
LinkedHashMap<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("id", record.getId());
|
||||
|
||||
row.put("outDocument", record.getOutDocument());
|
||||
row.put("quantity", record.getQuantity());
|
||||
row.put("securityCode", record.getSecurityCode());
|
||||
row.put("depoCodeSender", record.getDepoCodeSender());
|
||||
row.put("depoCodeAdressee", record.getDepoCodeAdressee());
|
||||
row.put("result", record.getResult());
|
||||
row.put("generationTime", record.getGenerationTime());
|
||||
row.put("generationId", record.getGenerationId());
|
||||
return row;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,62 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf14;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
|
||||
|
||||
@Service
|
||||
public class DF14Exporter extends AbstractExporterService<SDf14> {
|
||||
|
||||
public DF14Exporter(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider,
|
||||
SwtTable.SDF_14,
|
||||
IMDGDistributedNames.Map_SDf14, SDf14.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForFileName() {
|
||||
return "DF-14";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String sectionForFileName() {
|
||||
return "bond"; // todo bond / fund
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String typeForHeader() {
|
||||
return "014";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected String getDocumentNameForJournal() {
|
||||
//todo Выбор:
|
||||
// Уведомление о завершении расчетов (ОФЗ, ОБР)
|
||||
//или
|
||||
// Уведомление о завершении расчетов (ОФЗ, ОБР)
|
||||
return "Уведомление о завершении расчетов (ОФЗ, ОБР)";
|
||||
}
|
||||
|
||||
@Override
|
||||
protected LinkedHashMap<String, Object> convertRecord(SDf14 record) {
|
||||
LinkedHashMap<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("id", record.getId());
|
||||
|
||||
row.put("outDocument", record.getOutDocument());
|
||||
row.put("result", record.getResult());
|
||||
row.put("generationTime", record.getGenerationTime());
|
||||
row.put("generationId", record.getGenerationId());
|
||||
return row;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1,36 +0,0 @@
|
|||
package ru.spcex.clearing.swt.exporter.services.exportimpl;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.registry.Registry;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf12;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.swt.exporter.services.AbstractExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.FileStorage;
|
||||
import ru.spcex.platform.enumeration.SwtTable;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.math.BigDecimal;
|
||||
import java.time.LocalDate;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Service
|
||||
public class MoneyExporterService extends AbstractExporterService {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
||||
public MoneyExporterService(FileStorage fileStorage,
|
||||
KafkaSender kafkaSender,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(fileStorage, kafkaSender, imdgProvider,
|
||||
SwtTable.SDF_12,
|
||||
IMDGDistributedNames.Map_SDf12, SDf12.class);
|
||||
}
|
||||
|
||||
//todo impl...
|
||||
}
|
||||
|
|
@ -7,7 +7,7 @@ export-swt-service.hazelcast.password=dev-pass
|
|||
export-swt-service.common.encoding=cp866
|
||||
export-swt-service.common.threads-count=10
|
||||
|
||||
export-swt-service.out-dir=DocOut
|
||||
export-swt-service.out-dir=/opt/spcex/clearing/filedata/swt/SettlementHouse_DocOut
|
||||
|
||||
export-swt-service.kafka-consumer.bootstrap-servers=localhost:9092
|
||||
export-swt-service.kafka-consumer.group-id=dev-group-balance-service
|
||||
|
|
|
|||
|
|
@ -2,9 +2,11 @@ package ru.spcex.clearing.swt.exporter.services;
|
|||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import ru.spcex.clearing.swt.exporter.AbstractServiceTest;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.DF09Exporter;
|
||||
import ru.spcex.platform.enumeration.ResultStatuses;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
class AbstractExporterServiceTest extends AbstractServiceTest {
|
||||
@PostConstruct
|
||||
|
|
@ -14,10 +16,10 @@ class AbstractExporterServiceTest extends AbstractServiceTest {
|
|||
|
||||
@Test
|
||||
void sendSwtExportedNotification() {
|
||||
AbstractExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider);
|
||||
AbstractExporterService moneyExporterService = new DF09Exporter(null, kafkaSender, imdgProvider);
|
||||
|
||||
String fileName = "KS_RDC_DF-14_fund_202305241832.swt";
|
||||
moneyExporterService.sendSwtxportedNotification(fileName);
|
||||
String fileName = "KS_RDC_DF-14_bond_221005134616035.txt";
|
||||
moneyExporterService.sendSwtExportedNotification(LocalDateTime.now(), 1L, ResultStatuses.notSuccess);
|
||||
//TestUtils.waitingSendAndCheckRecord(null, mockProducer);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,102 +0,0 @@
|
|||
package ru.spcex.clearing.swt.exporter.services;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.util.StringUtils;
|
||||
import ru.clearing.classes.statics.data.registry.Registry;
|
||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||
import ru.spcex.clearing.swt.exporter.AbstractServiceTest;
|
||||
import ru.spcex.clearing.swt.exporter.services.exportimpl.MoneyExporterService;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
import java.math.BigDecimal;
|
||||
import java.util.Collection;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static ru.spcex.clearing.test.TestUtils.clearAllInImdg;
|
||||
|
||||
class MoneyExporterServiceTest extends AbstractServiceTest {
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
super.init();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link MoneyExporterService#getSwtFileRows()}<br>
|
||||
* Тест проверяет создание строк документа lim.<br>
|
||||
*/
|
||||
@Test // todo rewrite
|
||||
void getSwtFileRows() {
|
||||
clearAllInImdg(tradingClearingRegistryImdg);
|
||||
Registry registryA = getRegistryA(tcrA, securityIdFirst);
|
||||
Registry registryD = getRegistryD(tcrA, securityIdFirst);
|
||||
registryImdg.insert(registryA);
|
||||
registryImdg.insert(registryD);
|
||||
registryA = getRegistryA(tcrD, securityIdSecond);
|
||||
registryD = getRegistryD(tcrD, securityIdSecond);
|
||||
registryImdg.insert(registryA);
|
||||
registryImdg.insert(registryD);
|
||||
|
||||
MoneyExporterService moneyExporterService = new MoneyExporterService(null, kafkaSender, imdgProvider);
|
||||
Collection<String> limFileRows = moneyExporterService.getLimFileRows();
|
||||
assertEquals(0, limFileRows.size());
|
||||
|
||||
TradingClearingRegistry tradingClearingRegistry = new TradingClearingRegistry();
|
||||
tradingClearingRegistry.setId(securityIdFirst);
|
||||
tradingClearingRegistry.setStatus("ACTV");
|
||||
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
||||
tradingClearingRegistry.setId(securityIdSecond);
|
||||
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
||||
|
||||
limFileRows = moneyExporterService.getLimFileRows();
|
||||
assertEquals(2, limFileRows.size());
|
||||
assertTrue(limFileRows.contains(moneyExporterService.getRow(registryA, registryD)));
|
||||
}
|
||||
|
||||
private Registry getRegistryA(String tradingClearingRegistry, Long securityId) {
|
||||
Registry registry = new Registry();
|
||||
registry.setTradingCode("1A12323");
|
||||
registry.setBalance(new BigDecimal("10.00"));
|
||||
registry.setTradingClearingRegistry(tradingClearingRegistry);
|
||||
registry.setRegistryDesignation("A");
|
||||
registry.setRegistryInstrumentType("M");
|
||||
registry.setRegistryUnit("F");
|
||||
registry.setClearingCode(clearingCode(registry));
|
||||
registry.setClearingDate(currentDate);
|
||||
registry.setSecuritySymbol("RUB");
|
||||
registry.setSecurityId(securityId);
|
||||
registry.setTradingClearingRegistryId(securityId);
|
||||
return registry;
|
||||
}
|
||||
|
||||
private Registry getRegistryD(String tradingClearingRegistry, Long securityId) {
|
||||
Registry registry = new Registry();
|
||||
registry.setTradingCode("1A12323");
|
||||
registry.setBalance(new BigDecimal("5.00"));
|
||||
registry.setTradingClearingRegistry(tradingClearingRegistry);
|
||||
registry.setRegistryDesignation("D");
|
||||
registry.setRegistryInstrumentType("M");
|
||||
registry.setRegistryUnit("T");
|
||||
registry.setClearingCode(clearingCode(registry));
|
||||
registry.setClearingDate(currentDate);
|
||||
registry.setSecuritySymbol("RUB");
|
||||
registry.setSecurityId(securityId);
|
||||
registry.setTradingClearingRegistryId(securityId);
|
||||
return registry;
|
||||
}
|
||||
|
||||
|
||||
// see clearing-service ReistryUtil:
|
||||
|
||||
public static String clearingCode(Registry ofRegistry) {
|
||||
return clearingCode(ofRegistry.getRegistryDesignation(), ofRegistry.getRegistryInstrumentType(), ofRegistry.getRegistryCapacity(), ofRegistry.getRegistryUnit());
|
||||
}
|
||||
|
||||
public static String clearingCode(String registryDesignation, String registryInstrumentType, String registryCapacity, String registryUnit) {
|
||||
if (StringUtils.isEmpty(registryDesignation)) registryDesignation = "-";
|
||||
if (StringUtils.isEmpty(registryInstrumentType)) registryInstrumentType = "-";
|
||||
if (StringUtils.isEmpty(registryCapacity)) registryCapacity = "-";
|
||||
if (StringUtils.isEmpty(registryUnit)) registryUnit = "-";
|
||||
return registryDesignation + registryInstrumentType + registryCapacity + registryUnit;
|
||||
}
|
||||
}
|
||||
|
|
@ -3,6 +3,7 @@ package ru.spcex.platform.enumeration;
|
|||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum SwtTable implements IEnumKey {
|
||||
SDF_09("SDF_09"), SDF_11("SDF_11"),
|
||||
SDF_12("SDF_12"), SDF_14("SDF_14");
|
||||
|
||||
SwtTable(String key) {
|
||||
|
|
|
|||
|
|
@ -128,6 +128,7 @@ public interface Consts {
|
|||
String EXPORT_COMPLETED = "export_completed";
|
||||
String S_TRADES_IMPORTED = "s_trades-imported";
|
||||
String LIM_EXPORTED = "lim_exported";
|
||||
String JOURNAL_SERVICE = "journal-service-exported";
|
||||
String ACCOUNT_TERMINATION = "account-termination";
|
||||
String BALANCE_ACCOUNT_NEW = "balance-account-new";
|
||||
String BALANCE_ACCOUNT_UPDATE = "balance-account-update";
|
||||
|
|
|
|||
|
|
@ -0,0 +1,83 @@
|
|||
package ru.spcex.clearing.platform.messaging.domain.cud.utilities;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
|
||||
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
|
||||
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer;
|
||||
import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer;
|
||||
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer;
|
||||
import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer;
|
||||
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalTime;
|
||||
|
||||
public class JournalEventExportedRequest {
|
||||
@JsonSerialize(using = LocalDateSerializer.class)
|
||||
@JsonDeserialize(using = LocalDateDeserializer.class)
|
||||
@JsonProperty
|
||||
private LocalDate registratoinDate;
|
||||
@JsonSerialize(using = LocalTimeSerializer.class)
|
||||
@JsonDeserialize(using = LocalTimeDeserializer.class)
|
||||
@JsonProperty
|
||||
private LocalTime registrationTime;
|
||||
@JsonProperty
|
||||
private Long registrationNumber;
|
||||
@JsonProperty
|
||||
private String documentName;
|
||||
@JsonProperty
|
||||
private String dossierNumber;
|
||||
/**
|
||||
* ACK при успешной загрузке
|
||||
* NACK при ошибке загрузки
|
||||
*/
|
||||
@JsonProperty
|
||||
private String resultStatus;
|
||||
|
||||
public LocalDate getRegistratoinDate() {
|
||||
return registratoinDate;
|
||||
}
|
||||
|
||||
public void setRegistratoinDate(LocalDate registratoinDate) {
|
||||
this.registratoinDate = registratoinDate;
|
||||
}
|
||||
|
||||
public LocalTime getRegistrationTime() {
|
||||
return registrationTime;
|
||||
}
|
||||
|
||||
public void setRegistrationTime(LocalTime registrationTime) {
|
||||
this.registrationTime = registrationTime;
|
||||
}
|
||||
|
||||
public Long getRegistrationNumber() {
|
||||
return registrationNumber;
|
||||
}
|
||||
|
||||
public void setRegistrationNumber(Long registrationNumber) {
|
||||
this.registrationNumber = registrationNumber;
|
||||
}
|
||||
|
||||
public String getDocumentName() {
|
||||
return documentName;
|
||||
}
|
||||
|
||||
public void setDocumentName(String documentName) {
|
||||
this.documentName = documentName;
|
||||
}
|
||||
|
||||
public String getDossierNumber() {
|
||||
return dossierNumber;
|
||||
}
|
||||
|
||||
public void setDossierNumber(String dossierNumber) {
|
||||
this.dossierNumber = dossierNumber;
|
||||
}
|
||||
|
||||
public String getResultStatus() {
|
||||
return resultStatus;
|
||||
}
|
||||
|
||||
public void setResultStatus(String resultStatus) {
|
||||
this.resultStatus = resultStatus;
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue