From 27c2510588237259228bfa318e564eb1b479df72 Mon Sep 17 00:00:00 2001 From: Ivan Nikolaev-Axenov Date: Wed, 7 Aug 2024 16:41:27 +0300 Subject: [PATCH] notification sending on duplicate transaction number added in swt-importer http://jira.mfd.msk:8088/browse/CLS-726 --- .../logic/data/tables/AbstractTable.java | 8 ++++ .../logic/data/tables/SDf20Table.java | 22 ++++++++-- .../logic/data/tables/SDf21Table.java | 28 ++++--------- .../data/tables/WithTransactionNumber.java | 4 ++ .../swt/importer/logic/stages/ImportToDB.java | 42 ++++++++++++------- .../logic/stages/SWTImportKafkaMessenger.java | 2 +- 6 files changed, 66 insertions(+), 40 deletions(-) create mode 100644 clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/WithTransactionNumber.java diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/AbstractTable.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/AbstractTable.java index e675ac155..07d896ef5 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/AbstractTable.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/AbstractTable.java @@ -34,6 +34,10 @@ public abstract class AbstractTable { public abstract T getEntity(SWTRecord record); + public String getPrefix() { + return prefix; + } + public void injectEntity(T obj) { if (map == null) { bootMap(); @@ -41,6 +45,10 @@ public abstract class AbstractTable { map.insert(obj); } + public boolean checkOnExisting(T obj){ + return true; + } + protected void bootMap() { map = (ImdgHazelcast) hazelcastService.getImdg(nameOfMap, clazz); } diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf20Table.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf20Table.java index e31825c51..e0914d527 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf20Table.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf20Table.java @@ -1,14 +1,17 @@ package ru.spcex.clearing.swt.importer.logic.data.tables; +import java.time.Instant; +import java.util.Collection; +import java.util.Map; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import ru.clearing.classes.statics.data.sdf.SDf20; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.swt.importer.logic.data.enums.ETable; import ru.spcex.clearing.swt.importer.readers.SWTRecord; -import java.time.Instant; - -public class SDf20Table extends AbstractTable { - +public class SDf20Table extends AbstractTable implements WithTransactionNumber { + private final Logger log = LoggerFactory.getLogger(getClass()); private static final String PREFIX = ETable.S_DF_20.name(); private static final Class CLAZZ = SDf20.class; private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf20; @@ -41,4 +44,15 @@ public class SDf20Table extends AbstractTable { return result; } + @Override + public boolean checkOnExisting(SDf20 obj) { + if (map == null) { + bootMap(); + } + Collection sDf20s = map.getCollectionObjectsByFieldValues(Map.of("transactionNumber", obj.getTransactionNumber())); + if (!sDf20s.isEmpty()) { + log.warn("Skip insert by {}, sdf57.dbfId: {}", PREFIX, obj.getTransactionNumber()); + } + return sDf20s.isEmpty(); + } } \ No newline at end of file diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf21Table.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf21Table.java index cb8038b50..076868aa6 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf21Table.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/SDf21Table.java @@ -1,19 +1,17 @@ package ru.spcex.clearing.swt.importer.logic.data.tables; +import java.time.Instant; +import java.util.Collection; +import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import ru.clearing.classes.statics.data.sdf.SDf21; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.swt.importer.logic.data.enums.ETable; import ru.spcex.clearing.swt.importer.readers.SWTRecord; -import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; -import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; -import java.time.Instant; - -public class SDf21Table extends AbstractTable { +public class SDf21Table extends AbstractTable implements WithTransactionNumber { private final Logger log = LoggerFactory.getLogger(getClass()); - private static final String PREFIX = ETable.S_DF_21.name(); private static final Class CLAZZ = SDf21.class; private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf21; @@ -46,22 +44,14 @@ public class SDf21Table extends AbstractTable { } @Override - public void injectEntity(SDf21 obj) { + public boolean checkOnExisting(SDf21 obj) { if (map == null) { bootMap(); } - ImdgPredicateBuilder pb = map.predicateBuilder(); - ImdgPredicate sameTransactionNumberAndOperationCode = pb.and( - pb.equals("transactionNumber", obj.getTransactionNumber()), - pb.equals("operationCode", obj.getOperationCode()) - ); - SDf21 firstObjectBySQL = map.getFirstObjectByPredicate( - sameTransactionNumberAndOperationCode - ); - if (firstObjectBySQL != null) { - log.warn("Skip insert by {}, sdf21.transactionNumber: {}", PREFIX, obj.getTransactionNumber()); - } else { - map.insert(obj); + Collection sDf21s = map.getCollectionObjectsByFieldValues(Map.of("transactionNumber", obj.getTransactionNumber())); + if (!sDf21s.isEmpty()) { + log.warn("Skip insert by {}, sdf57.dbfId: {}", PREFIX, obj.getTransactionNumber()); } + return sDf21s.isEmpty(); } } \ No newline at end of file diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/WithTransactionNumber.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/WithTransactionNumber.java new file mode 100644 index 000000000..42a6d4993 --- /dev/null +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/data/tables/WithTransactionNumber.java @@ -0,0 +1,4 @@ +package ru.spcex.clearing.swt.importer.logic.data.tables; + +public interface WithTransactionNumber { +} diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java index 9895d7fb2..aadb9dee0 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java @@ -1,17 +1,6 @@ package ru.spcex.clearing.swt.importer.logic.stages; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.stereotype.Component; -import ru.spcex.clearing.swt.importer.config.settings.ImportSWTServiceSettings; -import ru.spcex.clearing.swt.importer.logic.data.ResultContainer; -import ru.spcex.clearing.swt.importer.logic.data.enums.ETable; -import ru.spcex.clearing.swt.importer.logic.data.enums.StageResult; -import ru.spcex.clearing.swt.importer.logic.data.tables.AbstractTable; -import ru.spcex.clearing.swt.importer.readers.SWTReader; -import ru.spcex.clearing.swt.importer.readers.SWTRecord; -import ru.spcex.clearing.swt.importer.readers.ValidationResult; -import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; -import ru.spcex.platform.utils.log.ExceptionUtils; +import static ru.spcex.clearing.swt.importer.readers.ValidationResult.SUCSESS; import java.io.ByteArrayInputStream; import java.io.IOException; @@ -20,8 +9,22 @@ import java.nio.charset.Charset; import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; - -import static ru.spcex.clearing.swt.importer.readers.ValidationResult.SUCSESS; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.swt.importer.config.settings.ImportSWTServiceSettings; +import ru.spcex.clearing.swt.importer.logic.data.ResultContainer; +import ru.spcex.clearing.swt.importer.logic.data.enums.ETable; +import ru.spcex.clearing.swt.importer.logic.data.enums.StageResult; +import ru.spcex.clearing.swt.importer.logic.data.tables.AbstractTable; +import ru.spcex.clearing.swt.importer.logic.data.tables.WithTransactionNumber; +import ru.spcex.clearing.swt.importer.readers.SWTReader; +import ru.spcex.clearing.swt.importer.readers.SWTRecord; +import ru.spcex.clearing.swt.importer.readers.ValidationResult; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.enumeration.ObjectType; +import ru.spcex.platform.enumeration.Priority; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; +import ru.spcex.platform.utils.log.ExceptionUtils; /** * Заливка проверенных данных в базу @@ -68,8 +71,15 @@ public class ImportToDB extends Stage { List records = swtReader.getGroupRecords(); records.forEach(record -> { if (record.isNotEmpty()) { - table.injectEntity(table.getEntity(record)); - log.debug("{} entity stored.", counter.incrementAndGet()); + SpcexObjectBase entityTable = table.getEntity(record); + if (!table.checkOnExisting(entityTable) && table instanceof WithTransactionNumber) { + kafkaMessenger.sendUserNotification(ObjectType.rgst, + "Номер транзакции " + record.getValFor23Tag() + " в полученном " + table.getPrefix() + " уже был обработан ранее", + Priority.HIGH); + } else { + table.injectEntity(entityTable); + log.debug("{} entity stored.", counter.incrementAndGet()); + } } }); } while (swtReader.hasNextRecord()); diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/SWTImportKafkaMessenger.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/SWTImportKafkaMessenger.java index 952469534..5433f3915 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/SWTImportKafkaMessenger.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/SWTImportKafkaMessenger.java @@ -84,7 +84,7 @@ public class SWTImportKafkaMessenger implements InitializingBean { // } } - private void sendUserNotification(ObjectType objectType, String comment, Priority priority) { + public void sendUserNotification(ObjectType objectType, String comment, Priority priority) { final String destination = Consts.NOTIFICATION_NEW; NotificationNewRequest request = new NotificationNewRequest(); //request.setObjectId();