account-service dbf-importer http://jira.mfd.msk:8088/browse/CLS-489 изменение статусов клиринговых счетов по S_DF52
This commit is contained in:
parent
8d864e8aca
commit
bd65726668
5 changed files with 163 additions and 8 deletions
|
|
@ -10,6 +10,10 @@ import org.springframework.beans.factory.annotation.Qualifier;
|
|||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.account.Account;
|
||||
import ru.clearing.classes.statics.data.account.ClearingAccount;
|
||||
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
||||
import ru.clearing.classes.statics.data.company.Company;
|
||||
import ru.clearing.classes.statics.data.company.relation.Relation;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf52;
|
||||
import ru.spcex.clearing.account.errors.AccountError;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
|
|
@ -26,22 +30,20 @@ import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
|||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.util.services.RequestHelper;
|
||||
import ru.spcex.clearing.validation.common.ValidationHelper;
|
||||
import ru.spcex.platform.enumeration.AccountType;
|
||||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
import ru.spcex.platform.enumeration.ServiceStatus;
|
||||
import ru.spcex.platform.enumeration.WorkflowStatus;
|
||||
import ru.spcex.platform.enumeration.*;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||
import ru.spcex.platform.utils.collection.Pair;
|
||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||
import ru.spcex.platform.utils.validation.IValidator;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.*;
|
||||
import java.util.function.Function;
|
||||
|
||||
@Service
|
||||
|
|
@ -56,6 +58,12 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
|||
private final Function<ClearingAccountNewRequest, IValidator> clearingAccountNewRequestValidator;
|
||||
private final Function<ClearingAccountUpdateRequest, IValidator> clearingAccountUpdateRequestValidator;
|
||||
|
||||
private final Imdg<SDf52> sdf52Imdg;
|
||||
private final Imdg<Account> accountImdg;
|
||||
private final Imdg<Company> companyImdg;
|
||||
private final Imdg<Relation> relationImdg;
|
||||
private final Imdg<ClearingMemberCategory> clearingMemberCategoryImdg;
|
||||
|
||||
@Autowired
|
||||
public ClearingAccountService(Consumer<String, Object> kafkaQueue,
|
||||
Producer<String, Object> kafkaResponseQueue,
|
||||
|
|
@ -78,6 +86,12 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
|||
this.requestHelper = requestHelper;
|
||||
this.clearingAccountNewRequestValidator = clearingAccountNewRequestValidator;
|
||||
this.clearingAccountUpdateRequestValidator = clearingAccountUpdateRequestValidator;
|
||||
|
||||
this.sdf52Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf52, SDf52.class);
|
||||
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
||||
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
||||
this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class);
|
||||
this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -91,6 +105,9 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
|||
callback(AccountSdf01Request.class)
|
||||
.setFunction(this::accountNewSdf01)
|
||||
.forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put);
|
||||
callback(StatementRequest.class)
|
||||
.setFunction(this::accountUpdateSdf52)
|
||||
.forDestination(Consts.ACCOUNT_PROCESS_SDF52, callbacks::put);
|
||||
|
||||
init();
|
||||
}
|
||||
|
|
@ -106,7 +123,8 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
|||
Instant now = Instant.now();
|
||||
Account account = new Account();
|
||||
account.setAccount(req.getAccount());
|
||||
account.setAccountType(AccountType.Clrn.getKey());if (req.getStatus() == null) {
|
||||
account.setAccountType(AccountType.Clrn.getKey());
|
||||
if (req.getStatus() == null) {
|
||||
account.setStatus(WorkflowStatus.Active.getKey());
|
||||
log.trace("Status not set in request. Use default: {}", account.getStatus());
|
||||
} else {
|
||||
|
|
@ -281,6 +299,130 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
|||
return null;
|
||||
}
|
||||
|
||||
String serviceTypeByCMCOfCompany(Long companyId) {
|
||||
Collection<ClearingMemberCategory> cmCategorys = clearingMemberCategoryImdg.getCollectionObjectsByFieldValues(
|
||||
Map.of("companyId", companyId));
|
||||
String serviceType = null;
|
||||
for (ClearingMemberCategory cmc : cmCategorys) {
|
||||
if (IEnumKey.contains(cmc.getClearingMemberCategory(), ClearingCategory.B, ClearingCategory.I, ClearingCategory.V)) {
|
||||
serviceType = ru.spcex.platform.enumeration.Service.MKR.getKey();
|
||||
}
|
||||
if (IEnumKey.contains(cmc.getClearingMemberCategory(), ClearingCategory.F, ClearingCategory.C)) {
|
||||
serviceType = ru.spcex.platform.enumeration.Service.MKR.getKey();
|
||||
}
|
||||
}
|
||||
return serviceType;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param status sdf.getStatus()
|
||||
* @return AccountStatus или null
|
||||
*/
|
||||
AccountStatus parseSdf52Status(Long status) {
|
||||
if (Long.valueOf(1L).equals(status)) {
|
||||
return AccountStatus.ACTIVE;
|
||||
} else if (Long.valueOf(0L).equals(status)) {
|
||||
return AccountStatus.BLOCKED;
|
||||
}
|
||||
if (Long.valueOf(2L).equals(status)) {
|
||||
return AccountStatus.CLOSE;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
public RequestInfoUpdate accountUpdateSdf52(BaseRequest<StatementRequest> systemRequest) {
|
||||
log.debug("accountUpdateSdf52 StatementRequest received, id={}", systemRequest.getId());
|
||||
|
||||
StatementRequest req = systemRequest.getRequestPayload();
|
||||
Long groupId = req.getGroupId();
|
||||
if (SdfTable.SDF_52 != req.getTable()) {
|
||||
log.warn("Unsupported table {} received on accountUpdateSdf52. Expected only {}.",
|
||||
req.getTable(), SdfTable.SDF_52);
|
||||
}
|
||||
Collection<SDf52> sdfs = sdf52Imdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"generationId", groupId
|
||||
));
|
||||
if (sdfs.isEmpty()) {
|
||||
log.info("Do not processing SDF52: S_DF52 not found by groupId={}", groupId);
|
||||
return null;
|
||||
}
|
||||
log.info("start processing {} SDF52: groupId={}", sdfs.size(), groupId);
|
||||
|
||||
List<Pair<SDf52, Account>> toUpdate = new ArrayList<>();
|
||||
{ // 1. Выборка данных
|
||||
for (SDf52 sDf52 : sdfs) {
|
||||
if (parseSdf52Status(sDf52.getStatus()) == null) {
|
||||
String msg = messageResolver.resolve(new EnumMessage(AccountError.WrongFieldValue, "status"));
|
||||
log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg);
|
||||
continue;
|
||||
}
|
||||
Company company = sDf52.getDeal() == null ? null : companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sDf52.getDeal()));
|
||||
if (company == null) {
|
||||
String msg = messageResolver.resolve(new EnumMessage(AccountError.CompanyNotFound, sDf52.getDeal()));
|
||||
log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg);
|
||||
continue;
|
||||
}
|
||||
String serviceType = serviceTypeByCMCOfCompany(company.getId());
|
||||
Relation relation = serviceType == null ? null : relationImdg.getSingleObjectByFieldValues(Map.of(
|
||||
"consumerId", company.getId(),
|
||||
"service", serviceType
|
||||
));
|
||||
if (relation == null) {
|
||||
String msg = messageResolver.resolve(new EnumMessage(AccountError.ClearingCategoryNotFound, company.getTradingCode(), serviceType));
|
||||
log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg);
|
||||
continue;
|
||||
}
|
||||
Account account = sDf52.getAccount() == null ? null : accountImdg.getFirstObjectByFieldValues(Map.of(
|
||||
"accountType", AccountType.Clrn.getKey(),
|
||||
"account", sDf52.getAccount(),
|
||||
"relationId", relation.getId()
|
||||
));
|
||||
if (account == null) {
|
||||
String msg = messageResolver.resolve(new EnumMessage(AccountError.AccountNotFound, sDf52.getAccount()));
|
||||
log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg);
|
||||
continue;
|
||||
}
|
||||
toUpdate.add(new Pair<>(sDf52, account));
|
||||
}
|
||||
}
|
||||
log.debug("Selected to update {} account's", toUpdate.size());
|
||||
int countOfUpdated = 0;
|
||||
ImdgTransaction imdgTransaction = imdgProvider.newTransaction();
|
||||
boolean txOk = false;
|
||||
imdgTransaction.beginTransaction();
|
||||
try { // 2. обновление данных, в транзакции
|
||||
Imdg<Account> accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
||||
Instant now = Instant.now();
|
||||
accountsLoop:
|
||||
for (Pair<SDf52, Account> item : toUpdate) {
|
||||
SDf52 sdf = item.getFirst();
|
||||
Account account = item.getSecond();
|
||||
AccountStatus newStatus = parseSdf52Status(sdf.getStatus());
|
||||
if (newStatus == null) {
|
||||
throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus());
|
||||
}
|
||||
if (!newStatus.equalsByKey(account.getStatus())) {
|
||||
account.setStatus(newStatus.getKey());
|
||||
account.setUpdated(now);
|
||||
accountImdg.update(account);
|
||||
log.trace("S_DF52[{}] do update status to {} for account[{}]", sdf.getId(), newStatus.getKey(), account.getId());
|
||||
countOfUpdated++;
|
||||
}
|
||||
}
|
||||
txOk = true;
|
||||
} finally {
|
||||
if (txOk) {
|
||||
imdgTransaction.commitTransaction();
|
||||
} else {
|
||||
log.debug("failed update, clearing accounts, rollback transaction. Request id={}; sdf52 groupId={}", systemRequest.getId(), groupId);
|
||||
imdgTransaction.rollbackTransaction();
|
||||
}
|
||||
}
|
||||
log.debug("successfully processed, grouping id={}. Updated {} of {} accounts.",
|
||||
groupId, countOfUpdated, toUpdate.size());
|
||||
return null;
|
||||
}
|
||||
|
||||
public void sendStatementRequestBack(Long groupingSdf01Id, Long groupSdf02Id, List<AccountSdfToStatementRequestPart> results) {
|
||||
StatementRequest request = new StatementRequest();
|
||||
request.setGroupId(groupingSdf01Id);
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ import org.apache.kafka.clients.consumer.MockConsumer;
|
|||
import org.apache.kafka.clients.producer.MockProducer;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
|
|
@ -224,6 +225,15 @@ class ClearingAccountServiceTest {
|
|||
ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount);
|
||||
}
|
||||
|
||||
@Test
|
||||
void parseSdf52Status() {
|
||||
Assertions.assertEquals(AccountStatus.BLOCKED, clearingAccountService.parseSdf52Status(0L));
|
||||
Assertions.assertEquals(AccountStatus.ACTIVE, clearingAccountService.parseSdf52Status(1L));
|
||||
Assertions.assertEquals(AccountStatus.CLOSE, clearingAccountService.parseSdf52Status(2L));
|
||||
Assertions.assertEquals(null, clearingAccountService.parseSdf52Status(3L));
|
||||
Assertions.assertEquals(null, clearingAccountService.parseSdf52Status(null));
|
||||
}
|
||||
|
||||
// /**
|
||||
// * {@link ClearingAccountService#accountNewSdf01(BaseRequest)}<br>
|
||||
// * Тест проверяет создание сущности {@link BaseRequest} в Hazelcast при передаче из Apache Kafka.<br>
|
||||
|
|
|
|||
|
|
@ -34,6 +34,7 @@ public class DbfImportKafkaMessenger implements InitializingBean {
|
|||
messengers.put(ETable.DF_06, groupId -> messageStatement(groupId, SdfTable.SDF_06, Consts.STATEMENT_PROCESS_SDF06));
|
||||
messengers.put(ETable.DF_09, groupId -> messageStatement(groupId, SdfTable.SDF_09));
|
||||
messengers.put(ETable.DF_16, groupId -> messageStatement(groupId, SdfTable.SDF_16));
|
||||
messengers.put(ETable.DF_52, groupId -> messageStatement(groupId, SdfTable.SDF_52, Consts.ACCOUNT_PROCESS_SDF52));
|
||||
messengers.put(ETable.DF_55, groupId -> messageStatement(groupId, SdfTable.SDF_55));
|
||||
messengers.put(ETable.DF_57, groupId -> messageStatement(groupId, SdfTable.SDF_57));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ public enum SdfTable implements IEnumKey {
|
|||
SDF_16("SDF_16"),
|
||||
SDF_20("SDF_20"),
|
||||
SDF_21("SDF_21"),
|
||||
SDF_52("SDF_52"),
|
||||
SDF_55("SDF_55"),
|
||||
SDF_57("SDF_57");
|
||||
|
||||
|
|
|
|||
|
|
@ -90,6 +90,7 @@ public interface Consts {
|
|||
|
||||
String ACCOUNT_NEW_SDF01 = "account-new-sdf01";
|
||||
String ACCOUNT_NEW_SDF08 = "account-new-sdf08";
|
||||
String ACCOUNT_PROCESS_SDF52 = "account-process-sdf52";
|
||||
|
||||
String DESTINATION_RELATION_NEW = "relation-new";
|
||||
String DESTINATION_RELATION_UPDATE = "relation-update";
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue