Sdf08Executor registry capacity: DepoAccount search fix

This commit is contained in:
ialbert 2023-07-06 19:03:15 +03:00
parent b469c162b2
commit 9f018d13c4
5 changed files with 76 additions and 39 deletions

View file

@ -273,7 +273,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
StatementRequest request = new StatementRequest(); StatementRequest request = new StatementRequest();
request.setGroupId(groupingSdf01Id); request.setGroupId(groupingSdf01Id);
request.setAccountCreationResults(results); request.setAccountCreationResults(results);
request.setContinueSdf01(true); request.setContinueSdf(true);
request.setTable(SdfTable.SDF_01); // по нему запрос получили request.setTable(SdfTable.SDF_01); // по нему запрос получили
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request)); log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));
kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request); kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request);

View file

@ -211,7 +211,7 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea
public void sendStatementRequestBack(Long groupingSdf01Id, List<AccountSdfToStatementRequestPart> results) { public void sendStatementRequestBack(Long groupingSdf01Id, List<AccountSdfToStatementRequestPart> results) {
StatementRequest request = new StatementRequest(); StatementRequest request = new StatementRequest();
request.setGroupId(groupingSdf01Id); request.setGroupId(groupingSdf01Id);
request.setContinueSdf01(true); //fixme???? request.setContinueSdf(true); //fixme????
request.setAccountCreationResults(results); request.setAccountCreationResults(results);
request.setTable(SdfTable.SDF_08); // по нему запрос получили request.setTable(SdfTable.SDF_08); // по нему запрос получили
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request)); log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));

View file

@ -66,20 +66,9 @@ public class StatementService extends QueueConsumer implements InitializingBean
public void afterPropertiesSet() throws Exception { public void afterPropertiesSet() throws Exception {
callback(StatementRequest.class) callback(StatementRequest.class)
.setConsumer(systemRequest -> { .setConsumer(systemRequest -> {
if (systemRequest.getRequestPayload() != null && systemRequest.getRequestPayload().isContinueSdf01()) { StatementRequest payload = systemRequest.getRequestPayload();
Optional<Long> sdf01And57Key = findCompletePair(); if (payload.isContinueSdf()) {
if (sdf01And57Key.isEmpty()) { processAccountAnswer(systemRequest);
log.error("FATAL: response from account-service received id={} sdf ids={}", systemRequest.getId(),
systemRequest.getRequestPayload()
.getAccountCreationResults()
.stream()
.map(res -> String.valueOf(res.getSdfId()))
.collect(Collectors.joining(",", "[", "]"))
);
return;
}
processSdf01And57(sdf01And57Key.get(), systemRequest.getRequestPayload());
log.info("Sdf01 and Sdf57 processed successfully after account-service command.");
} else { } else {
processPaired(systemRequest); processPaired(systemRequest);
} }
@ -88,10 +77,15 @@ public class StatementService extends QueueConsumer implements InitializingBean
init(); init();
} }
/**
* fromAccService передается когда пришел ответ от account-service
* в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01
* вместо старого sdf01 запроса выполняем новый пришедший от account-service
*/
private void processSdf01And57(Long key, StatementRequest fromAccService) { private void processSdf01And57(Long key, StatementRequest fromAccService) {
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key); Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
boolean allProcessed = processSdf01(fromAccService == null ? pair.getFirst() : fromAccService); Result sdf01Res = processSdf01(fromAccService == null ? pair.getFirst() : fromAccService);
if (!allProcessed) { if (sdf01Res.getAccountRequests().size() > 0) {
log.info("sdf01 execution wasn't complete, waiting for an answer from account-service"); log.info("sdf01 execution wasn't complete, waiting for an answer from account-service");
return; return;
} }
@ -106,6 +100,31 @@ public class StatementService extends QueueConsumer implements InitializingBean
} }
private void processAccountAnswer(BaseRequest<StatementRequest> systemRequest) {
StatementRequest payload = systemRequest.getRequestPayload();
switch (payload.getTable()) {
case SDF_01 -> {
//там нужно переиспользовать SDF57 запрос приходивший ранее и сохраненный в pairOfSdfRequest
Optional<Long> key = findCompleteKey(payload.getTable());
if (key.isEmpty()) {
log.error("FATAL: response from account-service received id={} table {} sdf ids={}", systemRequest.getId(),
payload.getTable(),
payload
.getAccountCreationResults()
.stream()
.map(res -> String.valueOf(res.getSdfId()))
.collect(Collectors.joining(",", "[", "]"))
);
return;
}
processSdf01And57(key.get(), payload);
}
case SDF_08 -> processSdf08(payload);
}
log.info("{} processed successfully after account-service command.", payload.getTable());
}
private void processPaired(BaseRequest<StatementRequest> systemRequest) { private void processPaired(BaseRequest<StatementRequest> systemRequest) {
StatementRequest statementRequest = systemRequest.getRequestPayload(); StatementRequest statementRequest = systemRequest.getRequestPayload();
SdfTable table = statementRequest.getTable(); SdfTable table = statementRequest.getTable();
@ -172,8 +191,17 @@ public class StatementService extends QueueConsumer implements InitializingBean
private void processSdf08(StatementRequest statementRequest) { private void processSdf08(StatementRequest statementRequest) {
Imdg<SDf08> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); Imdg<SDf08> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( Collection<? extends SpcexObjectBase> sdfGroup;
"generationId", statementRequest.getGroupId())); if (statementRequest.getAccountCreationResults().size() == 0) {
sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
} else {
sdfGroup = statementRequest.getAccountCreationResults()
.stream()
.filter(part -> part.getErrorCode() == null)
.map(part -> sdfImdg.getSingleObjectByID(part.getSdfId()))
.collect(Collectors.toList());
}
AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); AbstractExecutor service = executorsMap.get(SdfTable.SDF_08);
if (service != null) { if (service != null) {
Result res = service.execute(sdfGroup, statementRequest); Result res = service.execute(sdfGroup, statementRequest);
@ -196,7 +224,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
// finishSendCommand(res, service, statementRequest); // finishSendCommand(res, service, statementRequest);
} }
private boolean processSdf01(StatementRequest statementRequest) { private Result processSdf01(StatementRequest statementRequest) {
Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
Collection<? extends SpcexObjectBase> sdfGroup; Collection<? extends SpcexObjectBase> sdfGroup;
if (statementRequest.getAccountCreationResults().size() == 0) { if (statementRequest.getAccountCreationResults().size() == 0) {
@ -214,12 +242,12 @@ public class StatementService extends QueueConsumer implements InitializingBean
Result res = service.execute(sdfGroup, statementRequest); Result res = service.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() != 0) { if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
return false; return res;
} else if (service.isNeedToSendCommand()) { } else if (service.isNeedToSendCommand()) {
service.sendCommand(kafkaSender, res); service.sendCommand(kafkaSender, res);
return true; return res;
} }
return true; return res;
} }
private Optional<Long> saveRequest(StatementRequest statementRequest) { private Optional<Long> saveRequest(StatementRequest statementRequest) {
@ -282,14 +310,24 @@ public class StatementService extends QueueConsumer implements InitializingBean
return uncompletedPair; return uncompletedPair;
} }
private Optional<Long> findCompletePair() { private Optional<Long> findCompleteKey(SdfTable sdfTable) {
return pairOfSdfRequest.entrySet() Optional<Long> key = Optional.empty();
.stream() switch (sdfTable) {
.filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() != null) case SDF_01 -> key = pairOfSdfRequest.entrySet()
.filter(entry -> SdfTable.SDF_01.equals(entry.getValue().getFirst().getTable()) .stream()
&& SdfTable.SDF_57.equals(entry.getValue().getSecond().getTable())) .filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() != null)
.map(Map.Entry::getKey) .filter(entry -> SdfTable.SDF_01.equals(entry.getValue().getFirst().getTable())
.findFirst(); && SdfTable.SDF_57.equals(entry.getValue().getSecond().getTable()))
.map(Map.Entry::getKey)
.findFirst();
case SDF_08 -> key = pairOfSdfRequest.entrySet()
.stream()
.filter(entry -> entry.getValue().getFirst() != null)
.filter(entry -> SdfTable.SDF_08.equals(entry.getValue().getFirst().getTable()))
.map(Map.Entry::getKey)
.findFirst();
}
return key;
} }
private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests) { private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests) {

View file

@ -42,7 +42,6 @@ import java.util.Collection;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.function.Function; import java.util.function.Function;
import java.util.stream.Collectors;
@Service @Service
public class Sdf08Executor extends AbstractExecutor<SDf08> { public class Sdf08Executor extends AbstractExecutor<SDf08> {
@ -267,7 +266,7 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
rgs.setAccount(account.getAccount()); rgs.setAccount(account.getAccount());
rgs.setRegistryDesignation(RegistryDesignation.A.getKey()); rgs.setRegistryDesignation(RegistryDesignation.A.getKey());
rgs.setRegistryInstrumentType(RegistryInstrumentType.S.getKey()); rgs.setRegistryInstrumentType(RegistryInstrumentType.S.getKey());
DepoAccount depoAccount = depoAccountImdg.getSingleObjectByID(statement.getAccountId()); DepoAccount depoAccount = depoAccountImdg.getSingleObjectByFieldValues(Map.of("accountId", statement.getAccountId()));
if (depoAccount != null) { if (depoAccount != null) {
rgs.setRegistryCapacity(depoAccount.getDepoAccountType()); rgs.setRegistryCapacity(depoAccount.getDepoAccountType());
} }

View file

@ -16,18 +16,18 @@ public class StatementRequest {
@JsonProperty @JsonProperty
List<AccountSdfToStatementRequestPart> accountCreationResults = new ArrayList<>(); List<AccountSdfToStatementRequestPart> accountCreationResults = new ArrayList<>();
@JsonProperty @JsonProperty
private boolean continueSdf01 = false; private boolean continueSdf = false;
public Long getGroupId() { public Long getGroupId() {
return groupId; return groupId;
} }
public boolean isContinueSdf01() { public boolean isContinueSdf() {
return continueSdf01; return continueSdf;
} }
public void setContinueSdf01(boolean continueSdf01) { public void setContinueSdf(boolean continueSdf) {
this.continueSdf01 = continueSdf01; this.continueSdf = continueSdf;
} }
public void setGroupId(Long groupId) { public void setGroupId(Long groupId) {