notification sending on duplicate transaction number added in swt-importer http://jira.mfd.msk:8088/browse/CLS-726
This commit is contained in:
parent
6ccd8d3a3e
commit
27c2510588
6 changed files with 66 additions and 40 deletions
|
|
@ -34,6 +34,10 @@ public abstract class AbstractTable<T extends SpcexObjectBase> {
|
|||
|
||||
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<T extends SpcexObjectBase> {
|
|||
map.insert(obj);
|
||||
}
|
||||
|
||||
public boolean checkOnExisting(T obj){
|
||||
return true;
|
||||
}
|
||||
|
||||
protected void bootMap() {
|
||||
map = (ImdgHazelcast<T>) hazelcastService.getImdg(nameOfMap, clazz);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<SDf20> {
|
||||
|
||||
public class SDf20Table extends AbstractTable<SDf20> implements WithTransactionNumber {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private static final String PREFIX = ETable.S_DF_20.name();
|
||||
private static final Class<SDf20> CLAZZ = SDf20.class;
|
||||
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf20;
|
||||
|
|
@ -41,4 +44,15 @@ public class SDf20Table extends AbstractTable<SDf20> {
|
|||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean checkOnExisting(SDf20 obj) {
|
||||
if (map == null) {
|
||||
bootMap();
|
||||
}
|
||||
Collection<SDf20> sDf20s = map.getCollectionObjectsByFieldValues(Map.of("transactionNumber", obj.getTransactionNumber()));
|
||||
if (!sDf20s.isEmpty()) {
|
||||
log.warn("Skip insert by {}, sdf57.dbfId: {}", PREFIX, obj.getTransactionNumber());
|
||||
}
|
||||
return sDf20s.isEmpty();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<SDf21> {
|
||||
public class SDf21Table extends AbstractTable<SDf21> implements WithTransactionNumber {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
||||
private static final String PREFIX = ETable.S_DF_21.name();
|
||||
private static final Class<SDf21> CLAZZ = SDf21.class;
|
||||
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf21;
|
||||
|
|
@ -46,22 +44,14 @@ public class SDf21Table extends AbstractTable<SDf21> {
|
|||
}
|
||||
|
||||
@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<SDf21> sDf21s = map.getCollectionObjectsByFieldValues(Map.of("transactionNumber", obj.getTransactionNumber()));
|
||||
if (!sDf21s.isEmpty()) {
|
||||
log.warn("Skip insert by {}, sdf57.dbfId: {}", PREFIX, obj.getTransactionNumber());
|
||||
}
|
||||
return sDf21s.isEmpty();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,4 @@
|
|||
package ru.spcex.clearing.swt.importer.logic.data.tables;
|
||||
|
||||
public interface WithTransactionNumber {
|
||||
}
|
||||
|
|
@ -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<SWTRecord> 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());
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue