clearing-service sdf08 обработка автосоздания счетов

This commit is contained in:
AKurakin 2023-06-09 16:50:12 +03:00
parent 5cffd47463
commit f9fa84ad64
2 changed files with 19 additions and 6 deletions

View file

@ -123,7 +123,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_57);
service.execute(sdfGroup, statementRequest);
Result res = service.execute(sdfGroup, statementRequest);
finishSendCommand(res, service, statementRequest);
}
private void processSdf04(StatementRequest statementRequest) {
@ -132,7 +133,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_04);
if (service != null) {
service.execute(sdfGroup, statementRequest);
Result res = service.execute(sdfGroup, statementRequest);
finishSendCommand(res, service, statementRequest);
} else {
log.warn("Executor for SDF_04 not set");
}
@ -144,7 +146,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_08);
if (service != null) {
service.execute(sdfGroup, statementRequest);
Result res = service.execute(sdfGroup, statementRequest);
finishSendCommand(res, service, statementRequest);
} else {
log.warn("Executor for SDF_08 not set");
}
@ -155,7 +158,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_13);
service.execute(sdfGroup, statementRequest);
Result res = service.execute(sdfGroup, statementRequest);
finishSendCommand(res, service, statementRequest);
}
private void processSdf01(StatementRequest statementRequest) {
@ -174,6 +178,10 @@ public class StatementService extends QueueConsumer implements InitializingBean
}
AbstractExecutor service = executorsMap.get(SdfTable.SDF_01);
Result res = service.execute(sdfGroup, statementRequest);
finishSendCommand(res, service, statementRequest);
}
void finishSendCommand(Result res, AbstractExecutor service, StatementRequest statementRequest) {
if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
} else if (service.isNeedToSendCommand()) {

View file

@ -1,5 +1,7 @@
package ru.spcex.clearing.dbf.importer.logic.stages;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
@ -16,6 +18,7 @@ import java.util.function.Supplier;
@Component
public class DbfImportKafkaMessenger implements InitializingBean {
final Logger log = LoggerFactory.getLogger(getClass());
private final Supplier<KafkaSender> kafka;
private final Map<ETable, Consumer<Long>> messengers;
@ -30,7 +33,7 @@ public class DbfImportKafkaMessenger implements InitializingBean {
messengers.put(ETable.DF_09, groupId -> messageBalance(groupId, SdfTable.SDF_09));
messengers.put(ETable.DF_16, groupId -> messageBalance(groupId, SdfTable.SDF_16));
messengers.put(ETable.DF_57, groupId -> messageBalance(groupId, SdfTable.SDF_57));
messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04));
messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04));
}
/**
@ -48,7 +51,9 @@ public class DbfImportKafkaMessenger implements InitializingBean {
StatementRequest statementRequest = new StatementRequest();
statementRequest.setGroupId(groupId);
statementRequest.setTable(table);
kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest);
Long msgId = kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest);
log.debug("Send StatementRequest({}, {}) message id={} to kafka \"{}\"",
groupId, table, msgId, Consts.STATEMENT_PROCESS);
}
private void messageDf04(Long groupId) {