From 4256c053451404d20b2e9f02eb3eaf66b054db79 Mon Sep 17 00:00:00 2001 From: Ivan Nikolaev-Axenov Date: Mon, 2 Sep 2024 10:44:49 +0300 Subject: [PATCH 1/3] http://jira.mfd.msk:8088/browse/CLS-742 --- .../account/service/ClearingAccountService.java | 2 +- .../v2/listeners/ClientCodeMessageListener.java | 2 +- .../clearing/backendapi/config/MetaConfiguration.java | 2 +- .../payment/PaymentInstructionOutboundService.java | 11 ++++++++--- .../logic/stages/DbfImportKafkaMessenger.java | 2 +- .../trade/importer/services/TradeImporterService.java | 3 +-- .../platform/messaging/service/QueueConsumer.java | 3 ++- .../platform/messaging/service/QueueConsumerV2.java | 3 ++- 8 files changed, 17 insertions(+), 11 deletions(-) diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java index e92187e32..f95716e4a 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java @@ -793,7 +793,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin } if (returnTrigger) { log.debug("Remove from waiting list accountId={} (list groupId={})", accountId, inWaitingLst.getGroupId()); - waitingList.remove(inWaitingLst); + waitingList.remove(accountId); } if (returnTrigger) return inWaitingLst; diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java index ee8a81f81..a2d20ec51 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java @@ -132,7 +132,7 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali clientCodeNewRequest.setCompanyId(company.getId()); clientCodeNewRequest.setCode(tkrAccount.getClientCode()); clientCodeNewRequest.setDepoAccountId(depoAccountId); - Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB); + Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB.getKey()); clientCodeNewRequest.setMoneyAccountId(moneyAccount.getId()); List foreignCurrencyList = accountByCurrency.entrySet() diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java index 515707844..639edd9fb 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java @@ -17,7 +17,7 @@ import java.nio.file.Files; @Configuration public class MetaConfiguration { - @Value("file:${spring.config.location}/meta.json") + @Value("file:${spring.config.location}/meta/meta.json") private Resource meta; @Bean("metaJson") diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java index fc181c258..80b59e008 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java @@ -188,9 +188,14 @@ public class PaymentInstructionOutboundService { sDf54.setId(idGenerator.nextId()); sDf54.setGenerationId(pmt.getId()); if (time.isTradingTime()) { - Supplier builder = () -> GatewayRequestCreator.from(pmt, - assets.get().a__b().getTradingCode(), - tcr.getCode()); + Supplier builder = () -> { + if (tcr != null) { + return GatewayRequestCreator.from(pmt, + assets.get().a__b().getTradingCode(), + tcr.getCode()); + } + return null; + }; Optional gatewayOk = gateway.gatewayRequestAndWait(builder); if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) { String gtwErrMsg = msgs.resolve(mapError(gatewayOk)); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java index 597e99404..89f21e00f 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -60,7 +60,7 @@ public class DbfImportKafkaMessenger implements InitializingBean { } public void notifyUserAboutErrorParsing(Throwable error, ResultContainer resultContainer) { - if (resultContainer == null || resultContainer.getDbfFile() == null || resultContainer.getDbfFile() == null) + if (resultContainer == null || resultContainer.getDbfFile() == null || resultContainer.getDbfTable() == null) return; ETable currTable = resultContainer.getDbfTable(); String fileName = resultContainer.getDbfFile().getName(); diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java index 054038f4c..363d0cda6 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java @@ -213,8 +213,7 @@ public class TradeImporterService { private boolean isValidTrades(STrades trades) { return trades.getTradeDate() != null && trades.getTradeNum() != null - && StringUtils.hasText(trades.getOperation()) - && trades.getTradeNum() != null; + && StringUtils.hasText(trades.getOperation()); } public STrades readSTrades(ResultSet resultSet) throws SQLException { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index bfe4c6cb4..ca4322b92 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -2,6 +2,7 @@ package ru.spcex.clearing.platform.messaging.service; import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.ObjectMapper; +import java.util.concurrent.ExecutionException; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -196,7 +197,7 @@ public class QueueConsumer implements AutoCloseable { } send = producer.send(respRec); send.get(); - } catch (Exception e) { + } catch (InterruptedException | ExecutionException e) { log.error(ExceptionUtils.getStackTrace(e)); } }); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java index 38d084745..ef9a50f6c 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java @@ -6,6 +6,7 @@ import java.util.HashMap; import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -215,7 +216,7 @@ public class QueueConsumerV2 implements AutoCloseable { } send = producer.send(respRec); send.get(); - } catch (Exception e) { + } catch (InterruptedException | ExecutionException e) { log.error(ExceptionUtils.getStackTrace(e)); } }); From 3a43598223e96059f6362c50418c90b70d968228 Mon Sep 17 00:00:00 2001 From: Ivan Nikolaev-Axenov Date: Mon, 2 Sep 2024 12:43:07 +0300 Subject: [PATCH 2/3] http://jira.mfd.msk:8088/browse/CLS-742 HazelcastServiceBase --- .../imdg/iml/hazelcast/service/HazelcastServiceBase.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java index 3ae1ca517..844985cd0 100644 --- a/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java +++ b/platform-parent/platform-imdg-api-hazelcast-impl/src/main/java/ru/spcex/platform/imdg/iml/hazelcast/service/HazelcastServiceBase.java @@ -112,7 +112,11 @@ public abstract class HazelcastServiceBase break; } }); - } catch (Throwable e) { + } catch (InterruptedException e) { + log.info("Hazelcast: event interrupted"); + Thread.currentThread().interrupt(); + } + catch (Throwable e) { log.error("reinitializeHazelcastClient() has error: {}", ExceptionUtils.getStackTrace(e)); } } From 426d608a445eb7b9df728d75b9f11a54a35238b8 Mon Sep 17 00:00:00 2001 From: Ivan Nikolaev-Axenov Date: Tue, 10 Sep 2024 17:51:43 +0300 Subject: [PATCH 3/3] http://jira.mfd.msk:8088/browse/CLS-742 bug fixes --- .../backendapi/config/MetaConfiguration.java | 7 +++-- .../PaymentInstructionOutboundService.java | 16 +++++------ .../messaging/service/QueueConsumer.java | 27 ++++++++++--------- .../messaging/service/QueueConsumerV2.java | 12 ++++++--- 4 files changed, 32 insertions(+), 30 deletions(-) diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java index 639edd9fb..e5527b883 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/config/MetaConfiguration.java @@ -4,6 +4,8 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.DeserializationFeature; import com.fasterxml.jackson.databind.MapperFeature; import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.IOException; +import java.nio.file.Files; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; @@ -12,12 +14,9 @@ import org.springframework.context.annotation.Configuration; import org.springframework.core.io.Resource; import ru.spcex.clearing.backendapi.meta.MetaServer; -import java.io.IOException; -import java.nio.file.Files; - @Configuration public class MetaConfiguration { - @Value("file:${spring.config.location}/meta/meta.json") + @Value("file:${spring.config.location}/meta.json") private Resource meta; @Bean("metaJson") diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java index 80b59e008..96e2c35fc 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/payment/PaymentInstructionOutboundService.java @@ -1,5 +1,8 @@ package ru.spcex.clearing.service.payment; +import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError; +import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; + import java.math.BigDecimal; import java.time.Instant; import java.util.Locale; @@ -40,7 +43,6 @@ import ru.spcex.clearing.service.registry.RegistryManager; import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.validation.ValidationStored; import ru.spcex.clearing.session.stage.impl.GatewayRequester; -import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError; import ru.spcex.clearing.util.LocaleUtil; import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.platform.enumeration.CurrencyCode; @@ -55,7 +57,6 @@ import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.localization.SpringPropertiesLocalizer; -import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; import ru.spcex.platform.utils.text.TextUtil; import ru.spcex.platform.utils.validation.IValidator; @@ -188,14 +189,9 @@ public class PaymentInstructionOutboundService { sDf54.setId(idGenerator.nextId()); sDf54.setGenerationId(pmt.getId()); if (time.isTradingTime()) { - Supplier builder = () -> { - if (tcr != null) { - return GatewayRequestCreator.from(pmt, - assets.get().a__b().getTradingCode(), - tcr.getCode()); - } - return null; - }; + Supplier builder = () -> GatewayRequestCreator.from(pmt, + assets.get().a__b().getTradingCode(), + tcr == null ? null : tcr.getCode()); Optional gatewayOk = gateway.gatewayRequestAndWait(builder); if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) { String gtwErrMsg = msgs.resolve(mapError(gatewayOk)); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index ca4322b92..1e86f5c0b 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -2,7 +2,17 @@ package ru.spcex.clearing.platform.messaging.service; import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.ObjectMapper; +import java.time.Duration; +import java.time.temporal.ChronoUnit; +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Function; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -25,17 +35,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.log.ExceptionUtils; import ru.spcex.platform.utils.validation.IValidator; -import java.time.Duration; -import java.time.temporal.ChronoUnit; -import java.util.HashMap; -import java.util.Map; -import java.util.Optional; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.function.Function; - /** * утилитный класс для обработки сообщений из очереди */ @@ -130,6 +129,7 @@ public class QueueConsumer implements AutoCloseable { try { Thread.sleep(1000L); } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(ie)); break; } @@ -197,7 +197,10 @@ public class QueueConsumer implements AutoCloseable { } send = producer.send(respRec); send.get(); - } catch (InterruptedException | ExecutionException e) { + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.error(ExceptionUtils.getStackTrace(e)); + } catch (ExecutionException e) { log.error(ExceptionUtils.getStackTrace(e)); } }); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java index ef9a50f6c..03a5d1403 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumerV2.java @@ -1,5 +1,7 @@ package ru.spcex.clearing.platform.messaging.service; +import com.fasterxml.jackson.databind.JavaType; +import com.fasterxml.jackson.databind.ObjectMapper; import java.time.Duration; import java.time.temporal.ChronoUnit; import java.util.HashMap; @@ -13,8 +15,6 @@ import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Function; -import com.fasterxml.jackson.databind.JavaType; -import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -104,7 +104,8 @@ public class QueueConsumerV2 implements AutoCloseable { log.warn("Too many error at row, {}. Sleep.", lastErrors); try { Thread.sleep(1000L); - } catch (InterruptedException e){ + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(e)); break; } @@ -216,7 +217,10 @@ public class QueueConsumerV2 implements AutoCloseable { } send = producer.send(respRec); send.get(); - } catch (InterruptedException | ExecutionException e) { + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.error(ExceptionUtils.getStackTrace(e)); + } catch (ExecutionException e) { log.error(ExceptionUtils.getStackTrace(e)); } });