fix statementService when we should send command for continue session
This commit is contained in:
parent
0379f82276
commit
c922351170
5 changed files with 18 additions and 25 deletions
|
|
@ -2,27 +2,22 @@ package ru.spcex.clearing.config;
|
|||
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import ru.clearing.classes.statics.data.misc.Market;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.platform.enumeration.MarketType;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Configuration
|
||||
public class MarketCodesBySessionConfig {
|
||||
|
||||
@Bean(name = "marketCodesForBn")
|
||||
public List<String> marketCodesForBn(ImdgProvider imdgProvider) {
|
||||
Imdg<Market> marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class);
|
||||
Collection<Market> markets = marketImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"section", Section.FOND.getKey(),
|
||||
"marketType", MarketType.PRMR.getKey()));
|
||||
return markets.stream().map(Market::getCode).distinct().collect(Collectors.toList());
|
||||
// Imdg<Market> marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class);
|
||||
// Collection<Market> markets = marketImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
// "section", Section.FOND.getKey(),
|
||||
// "marketType", MarketType.PRMR.getKey()));
|
||||
// return markets.stream().map(Market::getCode).distinct().collect(Collectors.toList());
|
||||
return List.of("BMFC", "SMMC", "AEPC", "WBIC", "KBVC", "WBFC", "BMVC", "ABIC",
|
||||
"AESC", "ABFC", "KBIC", "WBMC", "BMMC", "ABEC", "SMIC", "ABMC", "SMVC", "KBFC", "SKVC", "ABVC", "WEPC",
|
||||
"KBMC", "SMFC", "BKVC", "WBVC", "BMIC", "WESC");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,11 +8,14 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRe
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
|
||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
||||
|
||||
@Service
|
||||
public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||
private final ClearingService clearingService;
|
||||
|
|
@ -65,9 +68,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
.setConsumer(primaryAuctionBnSession::continueSession)
|
||||
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
|
||||
|
||||
callback(Object.class)
|
||||
callback(STradesImportedRequest.class)
|
||||
.setConsumer(event -> clearingService.executeSTrade())
|
||||
.forDestination(Task.getOfTrades.topic(), callbacks::put);
|
||||
.forDestination(S_TRADES_IMPORTED, callbacks::put);
|
||||
callback(CreateRegistryRequest.class)
|
||||
.setConsumer(event -> registryService.createRegistry(event.getRequestPayload()))
|
||||
.forDestination(Consts.REGISTRY_NEW, callbacks::put);
|
||||
|
|
|
|||
|
|
@ -199,13 +199,13 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
Optional<Long> uncompletedPair = Optional.empty();
|
||||
switch (sdfTable) {
|
||||
case SDF_01 ->
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey);
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey);
|
||||
case SDF_57 ->
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey);
|
||||
case SDF_04 ->
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey);
|
||||
case SDF_13 ->
|
||||
case SDF_04 ->
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() == null && entry.getValue().getSecond() != null).findFirst().map(Map.Entry::getKey);
|
||||
case SDF_13 ->
|
||||
uncompletedPair = pairOfSdfRequest.entrySet().stream().filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() == null).findFirst().map(Map.Entry::getKey);
|
||||
}
|
||||
return uncompletedPair;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.service.LoggingService;
|
||||
import ru.spcex.clearing.service.model.Result;
|
||||
|
|
@ -95,10 +94,6 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
|
|||
|
||||
@Override
|
||||
public void sendCommand(KafkaSender kafkaSender, Result result) {
|
||||
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
|
||||
continueSessionBn.setGenerationId(result.getGenerationId());
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
|
||||
ExportToFileRequest exportRequest = new ExportToFileRequest();
|
||||
exportRequest.setSdfGroupId(result.getGenerationId());
|
||||
exportRequest.setNameOfTable(exportTableName());
|
||||
|
|
|
|||
|
|
@ -61,7 +61,7 @@ public class AbstractSession {
|
|||
|
||||
protected boolean checkStage(TaskType t) {
|
||||
synchronized (this.currStage) {
|
||||
return this.currStage.get().equals(t);
|
||||
return t.equals(this.currStage.get());
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue