http://jira.mfd.msk:8088/browse/CLS-36 parallel file execution temp fix
This commit is contained in:
parent
7ddad9bd21
commit
3bc6c83b1c
1 changed files with 45 additions and 18 deletions
|
|
@ -3,7 +3,6 @@ package ru.spcex.clearing.dbf.importer.services;
|
|||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
|
@ -12,9 +11,8 @@ import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer;
|
|||
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
|
||||
|
||||
import java.io.File;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.nio.file.Path;
|
||||
import java.util.*;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
|
|
@ -26,6 +24,7 @@ public class DBFImporterService {
|
|||
private final FileChecker fileChecker;
|
||||
private final ThreadPoolTaskExecutor executorService;
|
||||
private final Processor processor;
|
||||
private final Set<Path> filesCurrentlyInProcess;
|
||||
|
||||
public DBFImporterService(@Qualifier("fileChecker") FileChecker messageListener,
|
||||
@Qualifier("executor") ThreadPoolTaskExecutor executorService,
|
||||
|
|
@ -33,6 +32,7 @@ public class DBFImporterService {
|
|||
this.fileChecker = messageListener;
|
||||
this.executorService = executorService;
|
||||
this.processor = processor;
|
||||
this.filesCurrentlyInProcess = new HashSet<>();
|
||||
}
|
||||
|
||||
@Scheduled(cron = "${import-dbf-service.scheduler.check-src-dir-cron}")
|
||||
|
|
@ -43,25 +43,52 @@ public class DBFImporterService {
|
|||
public void run(ETable specificTable) {
|
||||
log.info("adding import task {}", specificTable != null ? specificTable.name() : "");
|
||||
executorService.execute(() -> {
|
||||
log.info("checking new files... {}", specificTable != null ? specificTable.name() : "");
|
||||
Map<ETable, List<File>> newFiles = fileChecker.checkNewFiles(specificTable);
|
||||
if (newFiles.size() > 0 && log.isDebugEnabled()) {
|
||||
log.debug("following files will be processed {}", forLogging(newFiles));
|
||||
} else if (newFiles.size() > 0) {
|
||||
log.info("following files will be processed {}", forLoggingSizeOnly(newFiles));
|
||||
} else {
|
||||
log.info("no files were found");
|
||||
}
|
||||
for (Map.Entry<ETable, List<File>> newFilesEntry : newFiles.entrySet()) {
|
||||
ETable currTable = newFilesEntry.getKey();
|
||||
List<File> fileList = newFilesEntry.getValue();
|
||||
for (File dbfFile : fileList) {
|
||||
processor.process(ResultContainer.createNewTask(currTable, dbfFile));
|
||||
Map<ETable, List<File>> newFiles = null;
|
||||
try {
|
||||
log.info("checking new files... {}", specificTable != null ? specificTable.name() : "");
|
||||
newFiles = getFiles(specificTable);
|
||||
if (newFiles.size() > 0 && log.isDebugEnabled()) {
|
||||
log.debug("following files will be processed {}", forLogging(newFiles));
|
||||
} else if (newFiles.size() > 0) {
|
||||
log.info("following files will be processed {}", forLoggingSizeOnly(newFiles));
|
||||
} else {
|
||||
log.info("no files were found");
|
||||
}
|
||||
for (Map.Entry<ETable, List<File>> newFilesEntry : newFiles.entrySet()) {
|
||||
ETable currTable = newFilesEntry.getKey();
|
||||
List<File> fileList = newFilesEntry.getValue();
|
||||
for (File dbfFile : fileList) {
|
||||
processor.process(ResultContainer.createNewTask(currTable, dbfFile));
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (newFiles != null && newFiles.size() > 0) {
|
||||
cleanFiles(newFiles);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private synchronized Map<ETable, List<File>> getFiles(ETable specificTable) {
|
||||
Map<ETable, List<File>> newFiles = fileChecker.checkNewFiles(specificTable);
|
||||
Map<ETable, List<File>> newFilesFiltered = newFiles.entrySet()
|
||||
.stream()
|
||||
.map(entry ->
|
||||
new AbstractMap.SimpleEntry<>(entry.getKey(), entry.getValue()
|
||||
.stream()
|
||||
.filter(file -> !filesCurrentlyInProcess.contains(file.toPath()))
|
||||
.collect(Collectors.toList())))
|
||||
.collect(Collectors.toMap(AbstractMap.SimpleEntry::getKey, AbstractMap.SimpleEntry::getValue));
|
||||
newFilesFiltered.forEach((table, files)
|
||||
-> files.forEach(file -> filesCurrentlyInProcess.add(file.toPath())));
|
||||
return newFilesFiltered;
|
||||
}
|
||||
|
||||
private synchronized void cleanFiles(Map<ETable, List<File>> filesFromTask) {
|
||||
filesFromTask.forEach((table, files)
|
||||
-> files.forEach(file -> filesCurrentlyInProcess.remove(file.toPath())));
|
||||
}
|
||||
|
||||
private static String forLogging(Map<ETable, List<File>> files) {
|
||||
return files
|
||||
.entrySet()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue