From 3bc6c83b1c797ea5f9eb76c790e609422aafdaf1 Mon Sep 17 00:00:00 2001 From: ialbert Date: Mon, 12 Dec 2022 18:54:33 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-36 parallel file execution temp fix --- .../importer/services/DBFImporterService.java | 63 +++++++++++++------ 1 file changed, 45 insertions(+), 18 deletions(-) diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java index fd3b2cdce..62a911533 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/DBFImporterService.java @@ -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 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> 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> newFilesEntry : newFiles.entrySet()) { - ETable currTable = newFilesEntry.getKey(); - List fileList = newFilesEntry.getValue(); - for (File dbfFile : fileList) { - processor.process(ResultContainer.createNewTask(currTable, dbfFile)); + Map> 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> newFilesEntry : newFiles.entrySet()) { + ETable currTable = newFilesEntry.getKey(); + List 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> getFiles(ETable specificTable) { + Map> newFiles = fileChecker.checkNewFiles(specificTable); + Map> 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> filesFromTask) { + filesFromTask.forEach((table, files) + -> files.forEach(file -> filesCurrentlyInProcess.remove(file.toPath()))); + } + private static String forLogging(Map> files) { return files .entrySet()