Pipelines removed, minor code changes in ObjectTags
This commit is contained in:
parent
b958dd0e5c
commit
54722afb2c
10 changed files with 73 additions and 128 deletions
|
|
@ -8,8 +8,6 @@ import com.fasterxml.jackson.dataformat.xml.XmlMapper;
|
|||
import com.fasterxml.jackson.dataformat.xml.ser.ToXmlGenerator;
|
||||
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import javax.xml.stream.XMLInputFactory;
|
||||
import javax.xml.stream.XMLOutputFactory;
|
||||
|
|
@ -30,9 +28,6 @@ import ru.spcex.clearing.xml.importer.logic.data.tags.objects.DF52ObjectTag;
|
|||
import ru.spcex.clearing.xml.importer.logic.data.tags.objects.DF55ObjectTag;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.tags.objects.DF57ObjectTag;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.tags.objects.ObjectTag;
|
||||
import ru.spcex.clearing.xml.importer.logic.stages.ChangeDirOfFileStage;
|
||||
import ru.spcex.clearing.xml.importer.logic.stages.ImportToDB;
|
||||
import ru.spcex.clearing.xml.importer.logic.stages.Stage;
|
||||
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||
|
||||
@Configuration
|
||||
|
|
@ -48,17 +43,6 @@ public class XMLImporterConfig {
|
|||
this.context = context;
|
||||
}
|
||||
|
||||
@Bean("pipeline")
|
||||
public List<Stage> pipeline() {
|
||||
List<Stage> pipeline = new LinkedList<>();
|
||||
log.info("XMLImporterConfig.pipeline");
|
||||
|
||||
pipeline.add(context.getBean(ImportToDB.class));
|
||||
pipeline.add(context.getBean(ChangeDirOfFileStage.class));
|
||||
|
||||
return pipeline;
|
||||
}
|
||||
|
||||
@Bean("executor")
|
||||
public ThreadPoolTaskExecutor executor() {
|
||||
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
|
||||
|
|
|
|||
|
|
@ -1,43 +0,0 @@
|
|||
package ru.spcex.clearing.xml.importer.logic;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.stereotype.Component;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult;
|
||||
import ru.spcex.clearing.xml.importer.logic.stages.Stage;
|
||||
|
||||
@Component("processor")
|
||||
public class Processor {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final List<Stage> pipeline;
|
||||
|
||||
public Processor(@Qualifier("pipeline") List<Stage> pipeline) {
|
||||
this.pipeline = pipeline;
|
||||
}
|
||||
|
||||
public void process(ResultContainer task) {
|
||||
log.info("uuid {}. Task started", task.getUuid());
|
||||
long startMills = System.currentTimeMillis();
|
||||
for (Stage currStage : pipeline) {
|
||||
if (statusIsFinal(task.getLastStageStatus()) && currStage.skipCompleted()) {
|
||||
continue;
|
||||
}
|
||||
log.info("uuid {}. Stage: {}", task.getUuid(), currStage.getClass().getSimpleName());
|
||||
task.setLastStageStatus(currStage.process(task));
|
||||
log.info("uuid {}. Stage {} finished with status {}", task.getUuid(), currStage.getClass().getSimpleName(), task.getLastStageStatus());
|
||||
}
|
||||
long endMills = System.currentTimeMillis();
|
||||
log.info("uuid {}. Task completed, result: {}, time working: {} ms",
|
||||
task.getUuid(),
|
||||
task,
|
||||
endMills - startMills);
|
||||
}
|
||||
|
||||
private boolean statusIsFinal(StageResult previousStageStatus) {
|
||||
return Arrays.asList(StageResult.ERROR, StageResult.COMPLETE).contains(previousStageStatus);
|
||||
}
|
||||
}
|
||||
|
|
@ -76,16 +76,16 @@ public class DF04ObjectTag extends ObjectTag<SDf04, DF04ObjectTag> {
|
|||
result.setDocnmprev(entity.getDocnmprev());
|
||||
result.setC_acc_deb(entity.getcAccDeb());
|
||||
result.setSbanknam1(entity.getSbanknam1());
|
||||
result.setSbanknam2(entity.getSbanknam2());
|
||||
result.setSbanknam3(entity.getSbanknam3());
|
||||
result.setSbanknam4(entity.getSbanknam4());
|
||||
result.setSbanknam5(entity.getSbanknam5());
|
||||
result.setSbanknam2(entity.getSbanknam1());
|
||||
result.setSbanknam3(entity.getSbanknam1());
|
||||
result.setSbanknam4(entity.getSbanknam1());
|
||||
result.setSbanknam5(entity.getSbanknam1());
|
||||
result.setC_acc_cred(entity.getcAccCred());
|
||||
result.setRbanknam1(entity.getRbanknam1());
|
||||
result.setRbanknam2(entity.getRbanknam2());
|
||||
result.setRbanknam3(entity.getRbanknam3());
|
||||
result.setRbanknam4(entity.getRbanknam4());
|
||||
result.setRbanknam5(entity.getRbanknam5());
|
||||
result.setRbanknam2(entity.getRbanknam1());
|
||||
result.setRbanknam3(entity.getRbanknam1());
|
||||
result.setRbanknam4(entity.getRbanknam1());
|
||||
result.setRbanknam5(entity.getRbanknam1());
|
||||
result.setPay_date(entity.getPayDate().toString());
|
||||
result.setPay_val(entity.getPayVal());
|
||||
result.setSum_deb(entity.getSumDeb());
|
||||
|
|
|
|||
|
|
@ -139,33 +139,33 @@ public class DF55ObjectTag extends ObjectTag<SDf55, DF55ObjectTag> {
|
|||
result.setSbankcode(entity.getSbankcode());
|
||||
result.setC_acc_deb(entity.getcAccDeb());
|
||||
result.setSbanknam1(entity.getSbanknam1());
|
||||
result.setSbanknam2(entity.getSbanknam2());
|
||||
result.setSbanknam3(entity.getSbanknam3());
|
||||
result.setSbanknam4(entity.getSbanknam4());
|
||||
result.setSbanknam5(entity.getSbanknam5());
|
||||
result.setSbanknam2(entity.getSbanknam1());
|
||||
result.setSbanknam3(entity.getSbanknam1());
|
||||
result.setSbanknam4(entity.getSbanknam1());
|
||||
result.setSbanknam5(entity.getSbanknam1());
|
||||
result.setRbankcode(entity.getRbankcode());
|
||||
result.setC_acc_cred(entity.getcAccCred());
|
||||
result.setRbanknam1(entity.getRbanknam1());
|
||||
result.setRbanknam2(entity.getRbanknam2());
|
||||
result.setRbanknam3(entity.getRbanknam3());
|
||||
result.setRbanknam4(entity.getRbanknam4());
|
||||
result.setRbanknam5(entity.getRbanknam5());
|
||||
result.setRbanknam2(entity.getRbanknam1());
|
||||
result.setRbanknam3(entity.getRbanknam1());
|
||||
result.setRbanknam4(entity.getRbanknam1());
|
||||
result.setRbanknam5(entity.getRbanknam1());
|
||||
result.setOp_type(entity.getOpType());
|
||||
result.setOp_order(entity.getOpOrder());
|
||||
result.setPay_date(entity.getPayDate().toString());
|
||||
result.setPay_val(entity.getPayVal());
|
||||
result.setSum_deb(entity.getSumDeb());
|
||||
result.setSclientn1(entity.getSclientn1());
|
||||
result.setSclientn2(entity.getSclientn2());
|
||||
result.setSclientn3(entity.getSclientn3());
|
||||
result.setSclientn4(entity.getSclientn4());
|
||||
result.setSclientn2(entity.getSclientn1());
|
||||
result.setSclientn3(entity.getSclientn1());
|
||||
result.setSclientn4(entity.getSclientn1());
|
||||
result.setInn_deb(entity.getInnCred());
|
||||
result.setKpp_deb(entity.getKppCred());
|
||||
result.setAcc_deb(entity.getAccDeb());
|
||||
result.setRclientn1(entity.getRclientn1());
|
||||
result.setRclientn2(entity.getRclientn2());
|
||||
result.setRclientn3(entity.getRclientn3());
|
||||
result.setRclientn4(entity.getRclientn4());
|
||||
result.setRclientn2(entity.getRclientn1());
|
||||
result.setRclientn3(entity.getRclientn1());
|
||||
result.setRclientn4(entity.getRclientn1());
|
||||
result.setInn_cred(entity.getInnCred());
|
||||
result.setKpp_cred(entity.getKppCred());
|
||||
result.setAcc_kr_1(entity.getAccKr1());
|
||||
|
|
|
|||
|
|
@ -131,33 +131,33 @@ public class DF57ObjectTag extends ObjectTag<SDf57, DF57ObjectTag> {
|
|||
result.setSbankcode(entity.getSbankcode());
|
||||
result.setC_acc_deb(entity.getcAccDeb());
|
||||
result.setSbanknam1(entity.getSbanknam1());
|
||||
result.setSbanknam2(entity.getSbanknam2());
|
||||
result.setSbanknam3(entity.getSbanknam3());
|
||||
result.setSbanknam4(entity.getSbanknam4());
|
||||
result.setSbanknam5(entity.getSbanknam5());
|
||||
result.setSbanknam2(entity.getSbanknam1());
|
||||
result.setSbanknam3(entity.getSbanknam1());
|
||||
result.setSbanknam4(entity.getSbanknam1());
|
||||
result.setSbanknam5(entity.getSbanknam1());
|
||||
result.setRbankcode(entity.getRbankcode());
|
||||
result.setC_acc_cred(entity.getcAccCred());
|
||||
result.setRbanknam1(entity.getRbanknam1());
|
||||
result.setRbanknam2(entity.getRbanknam2());
|
||||
result.setRbanknam3(entity.getRbanknam3());
|
||||
result.setRbanknam4(entity.getRbanknam4());
|
||||
result.setRbanknam5(entity.getRbanknam5());
|
||||
result.setRbanknam2(entity.getRbanknam1());
|
||||
result.setRbanknam3(entity.getRbanknam1());
|
||||
result.setRbanknam4(entity.getRbanknam1());
|
||||
result.setRbanknam5(entity.getRbanknam1());
|
||||
result.setOp_type(entity.getOpType());
|
||||
result.setPay_date(entity.getPayDate().toString());
|
||||
result.setExt_date(entity.getExtDate().toString());
|
||||
result.setPay_val(entity.getPayVal());
|
||||
result.setSum_deb(entity.getSumDeb());
|
||||
result.setSclientn1(entity.getSclientn1());
|
||||
result.setSclientn2(entity.getSclientn2());
|
||||
result.setSclientn3(entity.getSclientn3());
|
||||
result.setSclientn4(entity.getSclientn4());
|
||||
result.setSclientn2(entity.getSclientn1());
|
||||
result.setSclientn3(entity.getSclientn1());
|
||||
result.setSclientn4(entity.getSclientn1());
|
||||
result.setInn_deb(entity.getInnDeb());
|
||||
result.setKpp_deb(entity.getKppDeb());
|
||||
result.setAcc_deb(entity.getAccDeb());
|
||||
result.setRclientn1(entity.getRclientn1());
|
||||
result.setRclientn2(entity.getRclientn2());
|
||||
result.setRclientn3(entity.getRclientn3());
|
||||
result.setRclientn4(entity.getRclientn4());
|
||||
result.setRclientn2(entity.getRclientn1());
|
||||
result.setRclientn3(entity.getRclientn1());
|
||||
result.setRclientn4(entity.getRclientn1());
|
||||
result.setInn_cred(entity.getInnDeb());
|
||||
result.setKpp_cred(entity.getKppDeb());
|
||||
result.setAcc_kr(entity.getAccKr());
|
||||
|
|
|
|||
|
|
@ -1,16 +0,0 @@
|
|||
package ru.spcex.clearing.xml.importer.logic.stages;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult;
|
||||
|
||||
public abstract class Stage {
|
||||
protected final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
||||
public abstract StageResult process(ResultContainer resultContainer);
|
||||
|
||||
public boolean skipCompleted() {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package ru.spcex.clearing.xml.importer.logic.stages;
|
||||
package ru.spcex.clearing.xml.importer.logic.steps;
|
||||
|
||||
import static java.nio.file.StandardCopyOption.REPLACE_EXISTING;
|
||||
|
||||
|
|
@ -7,6 +7,8 @@ import java.io.IOException;
|
|||
import java.nio.file.Files;
|
||||
import java.time.Instant;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Component;
|
||||
import ru.spcex.clearing.xml.importer.config.settings.ImportXMLServiceSettings;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
|
||||
|
|
@ -14,7 +16,8 @@ import ru.spcex.clearing.xml.importer.logic.data.enums.StageResult;
|
|||
import ru.spcex.platform.utils.time.TimeUtil;
|
||||
|
||||
@Component
|
||||
public class ChangeDirOfFileStage extends Stage {
|
||||
public class ChangeDirOfFileStage {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy.MM.dd HH.mm.ss");
|
||||
private final ImportXMLServiceSettings settings;
|
||||
|
||||
|
|
@ -22,12 +25,6 @@ public class ChangeDirOfFileStage extends Stage {
|
|||
this.settings = settings;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean skipCompleted() {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public StageResult process(ResultContainer resultContainer) {
|
||||
File srcDir = new File(settings.getStore().getSrcDir());
|
||||
File outDir;
|
||||
|
|
@ -1,9 +1,11 @@
|
|||
package ru.spcex.clearing.xml.importer.logic.stages;
|
||||
package ru.spcex.clearing.xml.importer.logic.steps;
|
||||
|
||||
import com.fasterxml.jackson.dataformat.xml.XmlMapper;
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
import java.util.Map;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.stereotype.Component;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
|
||||
|
|
@ -16,7 +18,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
|
|||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||
|
||||
@Component
|
||||
public class ImportToDB extends Stage {
|
||||
public class ImportToDB {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final XmlMapper xmlMapper;
|
||||
private final HazelcastService hazelcastService;
|
||||
private final Map<ETable, ObjectTag<? extends SpcexObjectBase, ?>> objectTagMap;
|
||||
|
|
@ -32,7 +35,6 @@ public class ImportToDB extends Stage {
|
|||
this.kafkaMessenger = kafkaMessenger;
|
||||
}
|
||||
|
||||
@Override
|
||||
public StageResult process(ResultContainer resultContainer) {
|
||||
ETable currTable = resultContainer.getXmlTable();
|
||||
File xmlFile = resultContainer.getXmlFile();
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package ru.spcex.clearing.xml.importer.logic.stages;
|
||||
package ru.spcex.clearing.xml.importer.logic.steps;
|
||||
|
||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.PAIR_SDF;
|
||||
|
||||
|
|
@ -20,9 +20,10 @@ import org.springframework.scheduling.annotation.EnableScheduling;
|
|||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.spcex.clearing.xml.importer.logic.Processor;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.ResultContainer;
|
||||
import ru.spcex.clearing.xml.importer.logic.data.enums.ETable;
|
||||
import ru.spcex.clearing.xml.importer.logic.steps.ChangeDirOfFileStage;
|
||||
import ru.spcex.clearing.xml.importer.logic.steps.ImportToDB;
|
||||
import ru.spcex.platform.utils.collection.Pair;
|
||||
|
||||
@Service("xmlImporterService")
|
||||
|
|
@ -31,15 +32,19 @@ public class XMLImporterService {
|
|||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final FileChecker fileChecker;
|
||||
private final ThreadPoolTaskExecutor executorService;
|
||||
private final Processor processor;
|
||||
private final Set<Path> filesCurrentlyInProcess;
|
||||
|
||||
private final ImportToDB importToDB;
|
||||
private final ChangeDirOfFileStage changeDirOfFileStage;
|
||||
|
||||
public XMLImporterService(@Qualifier("fileChecker") FileChecker fileChecker,
|
||||
@Qualifier("executor") ThreadPoolTaskExecutor executorService,
|
||||
@Qualifier("processor") Processor processor) {
|
||||
ImportToDB importToDB,
|
||||
ChangeDirOfFileStage changeDirOfFileStage) {
|
||||
this.fileChecker = fileChecker;
|
||||
this.executorService = executorService;
|
||||
this.processor = processor;
|
||||
this.importToDB = importToDB;
|
||||
this.changeDirOfFileStage = changeDirOfFileStage;
|
||||
this.filesCurrentlyInProcess = new HashSet<>();
|
||||
}
|
||||
|
||||
|
|
@ -82,7 +87,23 @@ public class XMLImporterService {
|
|||
for (Pair<File, ETable> newFilesEntry : orderByTimeFiles) {
|
||||
ETable currTable = newFilesEntry.getSecond();
|
||||
File xmlFile = newFilesEntry.getFirst();
|
||||
processor.process(ResultContainer.createNewTask(currTable, xmlFile));
|
||||
ResultContainer task = ResultContainer.createNewTask(currTable, xmlFile);
|
||||
|
||||
log.info("uuid {}. Task started", task.getUuid());
|
||||
long startMills = System.currentTimeMillis();
|
||||
|
||||
log.info("uuid {}. Stage: Import to DB", task.getUuid());
|
||||
task.setLastStageStatus(importToDB.process(task));
|
||||
log.info("uuid {}. Stage import to DB finished with status {}", task.getUuid(), task.getLastStageStatus());
|
||||
log.info("uuid {}. Stage: Change directory of file", task.getUuid());
|
||||
task.setLastStageStatus(changeDirOfFileStage.process(task));
|
||||
log.info("uuid {}. Stage change directory of file finished with status {}", task.getUuid(), task.getLastStageStatus());
|
||||
|
||||
long endMills = System.currentTimeMillis();
|
||||
log.info("uuid {}. Task completed, result: {}, time working: {} ms",
|
||||
task.getUuid(),
|
||||
task,
|
||||
endMills - startMills);
|
||||
}
|
||||
} finally {
|
||||
if (newFiles != null && !newFiles.isEmpty()) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue