This commit is contained in:
AKurakin 2024-08-06 19:08:46 +03:00
parent 8076bc22da
commit 21d1f180be
7 changed files with 81 additions and 26 deletions

View file

@ -8,6 +8,7 @@ import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
import ru.spcex.clearing.service.execution.ExecutionCurrencyComponent; import ru.spcex.clearing.service.execution.ExecutionCurrencyComponent;
import ru.spcex.clearing.service.execution.ExecutionDepositComponent; import ru.spcex.clearing.service.execution.ExecutionDepositComponent;
import ru.spcex.clearing.service.execution.ExecutionFondComponent; import ru.spcex.clearing.service.execution.ExecutionFondComponent;
@ -105,12 +106,14 @@ public class ClearingService implements DisposableBean {
} }
} }
public void executeSTrade() { public void executeSTrade(STradesImportedRequest eventBody) {
executor.execute(() -> { executor.execute(() -> {
try { try {
executionDepositComponent.processNewTS(); Long fromId = eventBody == null ? null : eventBody.getMinId();
executionFondComponent.processNewTS(); log.trace("executeSTrade(fromId={})", fromId);
executionCurrencyComponent.processNewTS(); executionDepositComponent.processNewTS(fromId);
executionFondComponent.processNewTS(fromId);
executionCurrencyComponent.processNewTS(fromId);
} catch (Throwable e) { } catch (Throwable e) {
log.error("{}", ExceptionUtils.getStackTrace(e)); log.error("{}", ExceptionUtils.getStackTrace(e));
} }

View file

@ -173,7 +173,7 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
.forDestination(Consts.ASSET_OPERATION_APPROVAL, callbacks::put); .forDestination(Consts.ASSET_OPERATION_APPROVAL, callbacks::put);
callback(STradesImportedRequest.class) callback(STradesImportedRequest.class)
.setConsumer(event -> clearingService.executeSTrade()) .setConsumer(event -> clearingService.executeSTrade(event.getRequestPayload()))
.forDestination(S_TRADES_IMPORTED, callbacks::put); .forDestination(S_TRADES_IMPORTED, callbacks::put);
callback(CreateRegistryRequest.class) callback(CreateRegistryRequest.class)
.setConsumer(event -> registryService.updateRegistryIfNeeded(event.getRequestPayload())) .setConsumer(event -> registryService.updateRegistryIfNeeded(event.getRequestPayload()))

View file

@ -33,6 +33,7 @@ import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side; import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider; 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.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumId;
@ -92,15 +93,26 @@ public class ExecutionCurrencyComponent {
// log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay); // log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay);
// } // }
public void processNewTS() { public void processNewTS(Long fromId) {
LocalDate today = LocalDate.now(); LocalDate today = LocalDate.now();
log.debug("Start check new S_TRADE at {}", today); log.debug("Start check new S_TRADE at {}, fromId={}", today, fromId);
//выбираем STrades на сегодня с правильным section //выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( Collection<STrades> sTrades;
"tradeDate", today, if (fromId == null) {
"section", Section.CURR.getKey() sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
)); "tradeDate", LocalDate.now(),
log.info("Found {} s_trade for today", sTrades.size()); "section", Section.CURR.getKey()
));
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
pb.equals("tradeDate", LocalDate.now()),
pb.equals("section", Section.CURR.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
if (sTrades.isEmpty()) { if (sTrades.isEmpty()) {
logError(ClearingError.NewDealsNotFound); logError(ClearingError.NewDealsNotFound);

View file

@ -35,6 +35,7 @@ import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side; import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider; 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.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumId;
@ -96,14 +97,25 @@ public class ExecutionDepositComponent {
log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay); log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay);
} }
public void processNewTS() { public void processNewTS(Long fromId) {
log.debug("Start check new S_TRADE after {}", tradingDay); log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId);
//выбираем STrades на сегодня с правильным section //выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( Collection<STrades> sTrades;
"tradeDate", LocalDate.now(), if (fromId == null) {
"section", Section.MKR.getKey() sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
)); "tradeDate", LocalDate.now(),
log.info("Found {} s_trade for today", sTrades.size()); "section", Section.MKR.getKey()
));
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
pb.equals("tradeDate", LocalDate.now()),
pb.equals("section", Section.MKR.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
if (sTrades.isEmpty()) { if (sTrades.isEmpty()) {
logError(ClearingError.NewDealsNotFound); logError(ClearingError.NewDealsNotFound);

View file

@ -116,14 +116,25 @@ public class ExecutionFondComponent {
log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay); log.info("Reset trading day for search STrade: tradeNum={}, tradeDat={}", tradeNum, tradingDay);
} }
public void processNewTS() { public void processNewTS(Long fromId) {
log.debug("Start check new S_TRADE after {}", tradingDay); log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId);
//выбираем STrades на сегодня с правильным section //выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of( Collection<STrades> sTrades;
"tradeDate", LocalDate.now(), if (fromId == null) {
"section", Section.FOND.getKey() sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
)); "tradeDate", LocalDate.now(),
log.info("Found {} s_trade for today", sTrades.size()); "section", Section.FOND.getKey()
));
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
pb.equals("tradeDate", LocalDate.now()),
pb.equals("section", Section.FOND.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
if (sTrades.isEmpty()) { if (sTrades.isEmpty()) {
logError(ClearingError.NewDealsNotFound); logError(ClearingError.NewDealsNotFound);

View file

@ -77,6 +77,7 @@ public class TradeImporterService {
AtomicLong rows = new AtomicLong(); AtomicLong rows = new AtomicLong();
AtomicLong created = new AtomicLong(); AtomicLong created = new AtomicLong();
AtomicLong updated = new AtomicLong(); AtomicLong updated = new AtomicLong();
AtomicLong minCreatedId = new AtomicLong(Long.MAX_VALUE);
final long logProgressTime = 60_000L; // интервал вывода в лог каждую минуту final long logProgressTime = 60_000L; // интервал вывода в лог каждую минуту
AtomicLong logTime = new AtomicLong(System.currentTimeMillis() + logProgressTime); AtomicLong logTime = new AtomicLong(System.currentTimeMillis() + logProgressTime);
String query = String.format("SELECT * FROM %s.Trades", schema); String query = String.format("SELECT * FROM %s.Trades", schema);
@ -104,6 +105,9 @@ public class TradeImporterService {
tradesDb.setUpdated(currentInstant); tradesDb.setUpdated(currentInstant);
sTradesImdg.insert(tradesDb); sTradesImdg.insert(tradesDb);
created.incrementAndGet(); created.incrementAndGet();
if (minCreatedId.get() > tradesDb.getId()) {
minCreatedId.set(tradesDb.getId());
}
} }
} else { } else {
log.warn(messageResolver.resolve(new EnumMessage(sTradesNotValid, tradesDb.getTradeNum()))); log.warn(messageResolver.resolve(new EnumMessage(sTradesNotValid, tradesDb.getTradeNum())));
@ -119,6 +123,9 @@ public class TradeImporterService {
log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED); log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED);
} }
STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest();
if (!byCommand) {
sTradesImportedRequest.setMinId(minCreatedId.get());
}
kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest); kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest);
} else { } else {
log.debug("No STrades were created or updated from DB. Kafka command will not be send"); log.debug("No STrades were created or updated from DB. Kafka command will not be send");

View file

@ -5,6 +5,8 @@ import com.fasterxml.jackson.annotation.JsonProperty;
public class STradesImportedRequest { public class STradesImportedRequest {
@JsonProperty @JsonProperty
private Long tradeNum; private Long tradeNum;
@JsonProperty
private Long minId;
public Long getTradeNum() { public Long getTradeNum() {
return tradeNum; return tradeNum;
@ -13,4 +15,12 @@ public class STradesImportedRequest {
public void setTradeNum(Long tradeNum) { public void setTradeNum(Long tradeNum) {
this.tradeNum = tradeNum; this.tradeNum = tradeNum;
} }
public Long getMinId() {
return minId;
}
public void setMinId(Long minId) {
this.minId = minId;
}
} }