Сделал экспорт файлов файла .lim для установки лимитов по валютным средствам и отправку сообщения об ошибке формирования.
This commit is contained in:
parent
8315e9814b
commit
fce78c196b
4 changed files with 119 additions and 1 deletions
|
|
@ -7,10 +7,14 @@ import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.lim.exporter.config.SFTPConfig;
|
import ru.spcex.clearing.lim.exporter.config.SFTPConfig;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
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.Priority;
|
||||||
import ru.spcex.platform.enumeration.Section;
|
import ru.spcex.platform.enumeration.Section;
|
||||||
import ru.spcex.platform.enumeration.ServiceStatus;
|
import ru.spcex.platform.enumeration.ServiceStatus;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -65,6 +69,7 @@ public abstract class AbstractExporterService {
|
||||||
gateway.sendToSftp(limFile);
|
gateway.sendToSftp(limFile);
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
log.error("Failed export {} file", fileName);
|
log.error("Failed export {} file", fileName);
|
||||||
|
sendErrorNotification(fileName);
|
||||||
throw new RuntimeException(e);
|
throw new RuntimeException(e);
|
||||||
} finally {
|
} finally {
|
||||||
try {
|
try {
|
||||||
|
|
@ -87,6 +92,15 @@ public abstract class AbstractExporterService {
|
||||||
kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest);
|
kafkaSender.sendRequestToQueue(LIM_EXPORTED, limExportedRequest);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void sendErrorNotification(String fileName) {
|
||||||
|
NotificationNewRequest newRequest = new NotificationNewRequest();
|
||||||
|
newRequest.setObjectType(ObjectType.rgst.getKey());
|
||||||
|
newRequest.setComment("Failed export file = " + fileName + " for section = " + section().getKey());
|
||||||
|
newRequest.setPriority(Priority.HIGH.getKey());
|
||||||
|
log.info("Sending notification to kafka: {}", newRequest);
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, newRequest);
|
||||||
|
}
|
||||||
|
|
||||||
protected String prepareFileName(String target) {
|
protected String prepareFileName(String target) {
|
||||||
StringBuilder result = new StringBuilder();
|
StringBuilder result = new StringBuilder();
|
||||||
String dt = dtFormatter.format(LocalDateTime.now());
|
String dt = dtFormatter.format(LocalDateTime.now());
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,97 @@
|
||||||
|
package ru.spcex.clearing.lim.exporter.services;
|
||||||
|
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
|
import ru.spcex.clearing.lim.exporter.config.LimFormat;
|
||||||
|
import ru.spcex.clearing.lim.exporter.config.SFTPConfig;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.platform.enumeration.RegistryStatus;
|
||||||
|
import ru.spcex.platform.enumeration.Section;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||||
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||||
|
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
|
||||||
|
|
||||||
|
import java.math.BigDecimal;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collection;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.function.Function;
|
||||||
|
|
||||||
|
import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class CurrencyExporterService extends AbstractExporterService {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
|
public CurrencyExporterService(SFTPConfig.LimGateway gateway,
|
||||||
|
KafkaSender kafkaSender,
|
||||||
|
ImdgProvider imdgProvider) {
|
||||||
|
super(gateway, kafkaSender, imdgProvider);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String getTargetFileName() {
|
||||||
|
return prepareFileName("money");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Section section() {
|
||||||
|
return Section.CURR;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Collection<String> getLimFileRows() {
|
||||||
|
log.debug("Started loading and formation of currency(money) file lines");
|
||||||
|
List<String> limFileRows = new ArrayList<>();
|
||||||
|
Collection<Registry> registriesA = registryImdg.getCollectionObjectsBySQL(RegistryCodeSqlBuilder.getInstance(AM_F).build());
|
||||||
|
log.debug("Selected {} registriesA.", registriesA.size());
|
||||||
|
|
||||||
|
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
||||||
|
Function<Registry, ImdgPredicate> prdctDM_T = registry -> pb.and(
|
||||||
|
pb.sql(RegistryCodeSqlBuilder.getInstance(DM_T).build()),
|
||||||
|
pb.equals("registryStatus", RegistryStatus.PROC.getKey()),
|
||||||
|
pb.equals("tradingClearingRegistry", registry.getTradingClearingRegistry()),
|
||||||
|
pb.equals("securityId", registry.getSecurityId())
|
||||||
|
);
|
||||||
|
for (Registry registry : registriesA) {
|
||||||
|
if (checkNotBlocked(registry)) {
|
||||||
|
BigDecimal sumPositive = BigDecimal.ZERO;
|
||||||
|
if (registry.getTradingClearingRegistry() != null && registry.getSecurityId() != null) {
|
||||||
|
Collection<Registry> regsDMT = registryImdg.getCollectionObjectsByPredicate(prdctDM_T.apply(registry));
|
||||||
|
log.trace("Found {} DM_T registries for registry AM_F {}", regsDMT.size(), registry);
|
||||||
|
sumPositive = sumPositive.add(sum(regsDMT));
|
||||||
|
}
|
||||||
|
limFileRows.add(getRow(registry, sumPositive));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
log.debug("Successfully completed the formation of rows: {} for export currency(money)", limFileRows.size());
|
||||||
|
return limFileRows;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getRow(Registry registryA, BigDecimal sumPositive) {
|
||||||
|
StringBuilder row = new StringBuilder();
|
||||||
|
|
||||||
|
row.append("MONEY: FIRM_ID = ");
|
||||||
|
row.append(registryA.getTradingCode());
|
||||||
|
|
||||||
|
row.append("; TAG = SPVB");
|
||||||
|
|
||||||
|
row.append("; CURR_CODE = ");
|
||||||
|
row.append(registryA.getSecuritySymbol());
|
||||||
|
|
||||||
|
row.append("; CLIENT_CODE = ");
|
||||||
|
row.append(registryA.getTradingClearingRegistry());
|
||||||
|
|
||||||
|
row.append("; OPEN_BALANCE = ");
|
||||||
|
BigDecimal balance = registryA.getBalance() != null ? registryA.getBalance().subtract(sumPositive) : BigDecimal.ZERO;
|
||||||
|
row.append(LimFormat.toStringD2(balance));
|
||||||
|
|
||||||
|
row.append("; OPEN_LIMIT = 0.00");
|
||||||
|
|
||||||
|
row.append("; LIMIT_KIND = 0;");
|
||||||
|
return row.toString();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -15,14 +15,17 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final MoneyExporterService moneyExporterService;
|
private final MoneyExporterService moneyExporterService;
|
||||||
private final SecurityExporterService securityExporterService;
|
private final SecurityExporterService securityExporterService;
|
||||||
|
private final CurrencyExporterService currencyExporterService;
|
||||||
|
|
||||||
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
|
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
|
||||||
Producer<String, Object> kafkaProducer,
|
Producer<String, Object> kafkaProducer,
|
||||||
MoneyExporterService moneyExporterService,
|
MoneyExporterService moneyExporterService,
|
||||||
SecurityExporterService securityExporterService) {
|
SecurityExporterService securityExporterService,
|
||||||
|
CurrencyExporterService currencyExporterService) {
|
||||||
super(kafkaQueue, kafkaProducer);
|
super(kafkaQueue, kafkaProducer);
|
||||||
this.moneyExporterService = moneyExporterService;
|
this.moneyExporterService = moneyExporterService;
|
||||||
this.securityExporterService = securityExporterService;
|
this.securityExporterService = securityExporterService;
|
||||||
|
this.currencyExporterService = currencyExporterService;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -33,6 +36,9 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi
|
||||||
callback(LauncherCommandRequest.class)
|
callback(LauncherCommandRequest.class)
|
||||||
.setFunction(action -> securityExporterService.process())
|
.setFunction(action -> securityExporterService.process())
|
||||||
.forDestination(Task.unloadingSession_LIMS.topic(), callbacks::put); // LIMS
|
.forDestination(Task.unloadingSession_LIMS.topic(), callbacks::put); // LIMS
|
||||||
|
callback(LauncherCommandRequest.class)
|
||||||
|
.setFunction(action -> currencyExporterService.process())
|
||||||
|
.forDestination(Task.unloadingSession_LIMC.topic(), callbacks::put); // LIMC
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ public enum Task implements IEnumKey {
|
||||||
createDealsReport_GREF("GREF"), // Формирование финальной отчетности по сделкам
|
createDealsReport_GREF("GREF"), // Формирование финальной отчетности по сделкам
|
||||||
unloadingSession_LIMM("LIMM"),//Выгрузка в торговую систему остатков по деньгам
|
unloadingSession_LIMM("LIMM"),//Выгрузка в торговую систему остатков по деньгам
|
||||||
unloadingSession_LIMS("LIMS"),//Выгрузка в торговую систему остатков по бумагам
|
unloadingSession_LIMS("LIMS"),//Выгрузка в торговую систему остатков по бумагам
|
||||||
|
unloadingSession_LIMC("LIMC"),//Выгрузка в торговую систему остатков по валютам
|
||||||
liquidationSession_LIQU("LIQU"),//Ликвидационная сессия по обязательтсвам участника
|
liquidationSession_LIQU("LIQU"),//Ликвидационная сессия по обязательтсвам участника
|
||||||
startSession_STRM("STRM"),//Начало торговой сессии секции МКР
|
startSession_STRM("STRM"),//Начало торговой сессии секции МКР
|
||||||
terminationSession_ETRM("ETRM"),//Завершение торговой сессии секции МКР
|
terminationSession_ETRM("ETRM"),//Завершение торговой сессии секции МКР
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue