diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExportConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExportConfig.java index 8343b0c20..8be4a7799 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExportConfig.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExportConfig.java @@ -8,8 +8,9 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import ru.spcex.clearing.dbf.exporter.logic.stages.ExportFromDB; -import ru.spcex.clearing.dbf.exporter.logic.stages.PrepareDBFFiles; +import ru.spcex.clearing.dbf.exporter.logic.stages.PrepareDBFFile; import ru.spcex.clearing.dbf.exporter.logic.stages.Stage; import ru.spcex.clearing.dbf.exporter.properties.AProperties; @@ -49,10 +50,18 @@ public class DBFExportConfig { public List pipeline() { List pipeline = new LinkedList<>(); - pipeline.add(context.getBean(PrepareDBFFiles.class)); + pipeline.add(context.getBean(PrepareDBFFile.class)); pipeline.add(context.getBean(ExportFromDB.class)); return pipeline; } + @Bean("executor") + public ThreadPoolTaskExecutor executor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setMaxPoolSize(properties.getThreadsCount()); + executor.setThreadNamePrefix("dbf-exporter"); + return executor; + } + } diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/controller/DefaultController.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/controller/DefaultController.java index f9497f339..5889d1604 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/controller/DefaultController.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/controller/DefaultController.java @@ -4,15 +4,27 @@ package ru.spcex.clearing.dbf.exporter.controller; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.http.MediaType; import org.springframework.stereotype.Controller; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; import org.springframework.web.bind.annotation.ResponseBody; +import ru.spcex.clearing.dbf.exporter.logic.data.ISqlFilter; +import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; +import ru.spcex.clearing.dbf.exporter.services.DBFExportService; + +import java.util.EnumMap; +import java.util.Map; @Controller("/") public class DefaultController implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); + private final DBFExportService dbfExportService; + + public DefaultController(@Qualifier("dbfExportService") DBFExportService dbfExportService) { + this.dbfExportService = dbfExportService; + } @RequestMapping(method = RequestMethod.GET, path = "/test", produces = MediaType.TEXT_PLAIN_VALUE) @ResponseBody @@ -21,8 +33,19 @@ public class DefaultController implements InitializingBean { return "exporter controller test method"; } + @RequestMapping(method = RequestMethod.GET, path = "/export", produces = MediaType.TEXT_PLAIN_VALUE) + @ResponseBody + public String exportTables() { + log.info("Call export method for exporter controller"); + Map tablesForExport = new EnumMap<>(Table.class); + for (Table table : Table.values()) tablesForExport.put(table, null); + dbfExportService.run(tablesForExport); + return "export done, see log"; + } + @Override public void afterPropertiesSet() throws Exception { log.info("controller started"); + } } diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/data/ResultContainer.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/data/ResultContainer.java index 99d8e8e30..ecdc1b43e 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/data/ResultContainer.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/data/ResultContainer.java @@ -3,33 +3,56 @@ package ru.spcex.clearing.dbf.exporter.logic.data; import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; import java.io.File; -import java.util.Map; +import java.util.UUID; /** * Контейнер для передачи результата между стадиями */ public class ResultContainer { - private Map tablesForExport; - private Map filesForTables; + private UUID uuid; + private Table tableForExport; + private ISqlFilter filter; + private File fileForExport; protected ResultContainer() {} - public static ResultContainer createNewTask(Map tablesForExport) { + public static ResultContainer createNewTask(Table tableForExport, ISqlFilter filter) { ResultContainer container = new ResultContainer(); - container.tablesForExport = tablesForExport; + container.tableForExport = tableForExport; + container.filter = filter; + container.uuid = UUID.randomUUID(); return container; } - public Map getTablesForExport() { - return tablesForExport; + public Table getTableForExport() { + return tableForExport; } - public Map getFilesForTables() { - return filesForTables; + public void setTableForExport(Table tableForExport) { + this.tableForExport = tableForExport; } - public void setFilesForTables(Map filesForTables) { - this.filesForTables = filesForTables; + public ISqlFilter getFilter() { + return filter; } + public void setFilter(ISqlFilter filter) { + this.filter = filter; + } + + public File getFileForExport() { + return fileForExport; + } + + public void setFileForExport(File fileForExport) { + this.fileForExport = fileForExport; + } + + public UUID getUuid() { + return uuid; + } + + public void setUuid(UUID uuid) { + this.uuid = uuid; + } } diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromDB.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromDB.java index 07142a7b8..a95f588ca 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromDB.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromDB.java @@ -1,7 +1,6 @@ package ru.spcex.clearing.dbf.exporter.logic.stages; import com.linuxense.javadbf.DBFDataType; -import com.linuxense.javadbf.DBFException; import com.linuxense.javadbf.DBFField; import com.linuxense.javadbf.DBFWriter; import org.springframework.beans.factory.InitializingBean; @@ -19,7 +18,6 @@ import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; import ru.spcex.clearing.dbf.exporter.properties.AProperties; import java.io.File; -import java.io.IOException; import java.math.BigDecimal; import java.nio.charset.Charset; import java.sql.ResultSet; @@ -46,59 +44,51 @@ public class ExportFromDB extends Stage implements InitializingBean { @Override public StageResult process(ResultContainer resultContainer) { - for (Table table : resultContainer.getTablesForExport().keySet()) { - ISqlFilter sqlFilter = resultContainer.getTablesForExport().get(table); - String sqlCondition = sqlFilter != null ? " where " + sqlFilter : ""; + Objects.requireNonNull(resultContainer.getTableForExport()); + Objects.requireNonNull(resultContainer.getFileForExport()); + + Table table = resultContainer.getTableForExport(); + ISqlFilter sqlFilter = resultContainer.getFilter(); + String sqlCondition = sqlFilter != null ? " where " + sqlFilter : ""; - String sql = "select * from " + table + sqlCondition; - List> recordsFromDB = dbfJdbcTemplate.queryForList(sql); + String sql = "select * from " + table + sqlCondition; + List> recordsFromDB = dbfJdbcTemplate.queryForList(sql); - File dbfFile = resultContainer.getFilesForTables().get(table); - DBFField[] dbfFields = dbfFieldsForTable.get(table); + File dbfFile = resultContainer.getFileForExport(); + DBFField[] dbfFields = dbfFieldsForTable.get(table); - try { - boolean createOk = dbfFile.createNewFile(); - if (!createOk) throw new IOException("Can't create " + dbfFile.getName() + " file"); - } catch (IOException e) { - log.error("Can't create file " + dbfFile + " for table " + table + ". Table was skipped.", e); - continue; - } - - try (DBFWriter dbfWriter = new DBFWriter(dbfFile, dbfCharset)) { - dbfWriter.setFields(dbfFields); - for (Map record : recordsFromDB) { - int columnCount = dbfFields.length; - Object[] values = new Object[columnCount]; - for (int columnIdx = 0; columnIdx < columnCount; columnIdx++) { - DBFDataType currDBFType = dbfFields[columnIdx].getType(); - Object currColumn = record.get(dbfFields[columnIdx].getName()); - - if (currDBFType == DBFDataType.LOGICAL) { - values[columnIdx] = Boolean.parseBoolean(String.valueOf(currColumn)); - - } else if (currDBFType == DBFDataType.DATE) { - LocalDate date = LocalDate.parse(String.valueOf(currColumn)); - values[columnIdx] = new Date(date.atStartOfDay(ZoneOffset.UTC).toInstant().toEpochMilli()); - - } else if (currDBFType == DBFDataType.NUMERIC || currDBFType == DBFDataType.FLOATING_POINT) { - BigDecimal bigDecimal = new BigDecimal(String.valueOf(currColumn)); - values[columnIdx] = bigDecimal; - - } else { - values[columnIdx] = String.valueOf(currColumn); - - } - } - dbfWriter.addRecord(values); + try (DBFWriter dbfWriter = new DBFWriter(dbfFile, dbfCharset)) { + dbfWriter.setFields(dbfFields); + for (Map record : recordsFromDB) { + int columnCount = dbfFields.length; + Object[] values = new Object[columnCount]; + for (int columnIdx = 0; columnIdx < columnCount; columnIdx++) { + DBFDataType currDBFType = dbfFields[columnIdx].getType(); + Object currColumn = record.get(dbfFields[columnIdx].getName()); + values[columnIdx] = convertToDBFValue(currDBFType, currColumn); } - } catch (DBFException e) { - log.error("Can't export table " + table + " to file " + dbfFile + ". Table was skipped.", e); + dbfWriter.addRecord(values); } + } catch (Exception e) { + log.error(String.format("uuid %s. Can't export table %s to file %s. Table was skipped.", resultContainer.getUuid(), table, dbfFile), e); + return StageResult.ERROR; } - return StageResult.OK; + return StageResult.COMPLETE; } + private Object convertToDBFValue(DBFDataType dbfType, Object valueFromDB) throws Exception { + if (dbfType == DBFDataType.LOGICAL) { + return Boolean.parseBoolean(String.valueOf(valueFromDB)); + } else if (dbfType == DBFDataType.DATE) { + LocalDate date = LocalDate.parse(String.valueOf(valueFromDB)); + return new Date(date.atStartOfDay(ZoneOffset.UTC).toInstant().toEpochMilli()); + } else if (dbfType == DBFDataType.NUMERIC || dbfType == DBFDataType.FLOATING_POINT) { + return new BigDecimal(String.valueOf(valueFromDB)); + } else { + return String.valueOf(valueFromDB); + } + } /** * 1. Собирает структуру БД для дальнейшей записи DBF файлов diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFile.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFile.java new file mode 100644 index 000000000..cae1d707a --- /dev/null +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFile.java @@ -0,0 +1,60 @@ +package ru.spcex.clearing.dbf.exporter.logic.stages; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult; +import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; +import ru.spcex.clearing.dbf.exporter.properties.AProperties; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Objects; + +/** + * Создание DBF файла + */ +@Component +public class PrepareDBFFile extends Stage implements InitializingBean { + private final AProperties properties; + private String outDir; + + public PrepareDBFFile(@Qualifier("dbfExporterProperties") AProperties properties) { + this.properties = properties; + } + + @Override + public StageResult process(ResultContainer resultContainer) { + Objects.requireNonNull(resultContainer.getTableForExport()); + + Table table = resultContainer.getTableForExport(); + File dbfFile = new File(String.format("%s%s%s.dbf", + outDir, + File.separator, + table.getFilePrefix())); + try { + Path dbfFilePath = dbfFile.toPath(); + Files.deleteIfExists(dbfFilePath); + Files.createFile(dbfFilePath); + } catch (IOException e) { + log.error(String.format("uuid %s. Can't create file %s", resultContainer.getUuid(), dbfFile.getName()), e); + return StageResult.ERROR; + } + resultContainer.setFileForExport(dbfFile); + return StageResult.OK; + } + + + @Override + public void afterPropertiesSet() throws Exception { + String outDirPath = properties.getOutDir(); + File outDirFile = new File(outDirPath); + if (outDirFile.exists() && !outDirFile.isDirectory()) + throw new IOException("Output directory " + outDirPath + " is file."); + if (!outDirFile.exists()) Files.createDirectories(outDirFile.toPath()); + this.outDir = outDirPath; + } +} diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFiles.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFiles.java deleted file mode 100644 index 99195cff4..000000000 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/PrepareDBFFiles.java +++ /dev/null @@ -1,69 +0,0 @@ -package ru.spcex.clearing.dbf.exporter.logic.stages; - -import org.springframework.beans.factory.InitializingBean; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.stereotype.Component; -import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; -import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult; -import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; -import ru.spcex.clearing.dbf.exporter.properties.AProperties; - -import java.io.File; -import java.io.IOException; -import java.time.LocalDateTime; -import java.time.format.DateTimeFormatter; -import java.util.HashMap; -import java.util.Map; - -/** - * Создание DBF файлов - */ -@Component -public class PrepareDBFFiles extends Stage implements InitializingBean { - private final AProperties properties; - private String outDirWithDateTime; - - public PrepareDBFFiles(@Qualifier("dbfExporterProperties") AProperties properties) { - this.properties = properties; - } - - @Override - public StageResult process(ResultContainer resultContainer) { - Map dbfFiles = new HashMap<>(); - for (Table table : resultContainer.getTablesForExport().keySet()) { - File dbfFile = new File(String.format("%s%s%s.dbf", - outDirWithDateTime, - File.separator, - table.getFilePrefix())); - try { - if (dbfFile.exists()) { - boolean deleteOk = dbfFile.delete(); - if (!deleteOk) throw new IOException("Can't overwrite " + dbfFile.getName() + " file"); - } - - dbfFiles.put(table, dbfFile); - } catch (IOException e) { - log.error("Can't create " + dbfFile + ". Table was skipped", e); - resultContainer.getTablesForExport().remove(table); - } - } - resultContainer.setFilesForTables(dbfFiles); - return StageResult.OK; - } - - - @Override - public void afterPropertiesSet() throws Exception { - String outDirPath = properties.getOutDir(); - File outDir = new File(String.format("%s%s%s", - outDirPath, - File.separator, - LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd_HHmmss")))); - if (outDir.exists()) throw new IOException("Output dir " + outDir + " already exist."); - - boolean createDirOk = outDir.mkdirs(); - if (!createDirOk) throw new IOException("Can't create output dir " + outDir); - - outDirWithDateTime = outDir.getAbsolutePath(); - } -} diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/Processor.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/Processor.java new file mode 100644 index 000000000..5fc31168e --- /dev/null +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/Processor.java @@ -0,0 +1,37 @@ +package ru.spcex.clearing.dbf.exporter.logic.stages; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult; + +import java.util.Arrays; +import java.util.List; + +@Component("processor") +public class Processor { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final List pipeline; + + public Processor(@Qualifier("pipeline") List pipeline) { + this.pipeline = pipeline; + } + + public void process(ResultContainer task) { + log.info("uuid {}. Task started", task.getUuid()); + StageResult result = null; + long startMills = System.currentTimeMillis(); + for (Stage currStage : pipeline) { + log.info("uuid {}. Stage: {}", task.getUuid(), currStage.getClass().getSimpleName()); + result = currStage.process(task); + if (Arrays.asList(StageResult.ERROR, StageResult.COMPLETE).contains(result)) break; + } + long endMills = System.currentTimeMillis(); + log.info("uuid {}. Task completed, result: {}, time working: {}", + task.getUuid(), + result, + endMills - startMills); + } +} diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/properties/AProperties.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/properties/AProperties.java index c2260f04e..26e589067 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/properties/AProperties.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/properties/AProperties.java @@ -24,6 +24,9 @@ public class AProperties { @Value("${dbf.encoding}") private String dbfEncoding; + @Value("${dbf.threads-count}") + private int threadsCount; + public String getJdbcUrl() { return jdbcUrl; } @@ -71,4 +74,12 @@ public class AProperties { public void setDbfEncoding(String dbfEncoding) { this.dbfEncoding = dbfEncoding; } + + public int getThreadsCount() { + return threadsCount; + } + + public void setThreadsCount(int threadsCount) { + this.threadsCount = threadsCount; + } } diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java index b6071d575..e199d31f9 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java @@ -1,42 +1,31 @@ package ru.spcex.clearing.dbf.exporter.services; 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; import ru.spcex.clearing.dbf.exporter.logic.data.ISqlFilter; import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; -import ru.spcex.clearing.dbf.exporter.logic.data.enums.StageResult; import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; -import ru.spcex.clearing.dbf.exporter.logic.stages.Stage; -import ru.spcex.clearing.dbf.exporter.properties.AProperties; +import ru.spcex.clearing.dbf.exporter.logic.stages.Processor; -import java.util.Arrays; -import java.util.HashMap; -import java.util.List; import java.util.Map; -@Service -@EnableScheduling +@Service("dbfExportService") public class DBFExportService { - private final AProperties properties; - private final List pipeline; + private final ThreadPoolTaskExecutor executor; + private final Processor processor; - public DBFExportService(AProperties properties, - @Qualifier("pipeline") List pipeline) { - this.properties = properties; - this.pipeline = pipeline; + public DBFExportService(@Qualifier("executor") ThreadPoolTaskExecutor executor, + @Qualifier("processor") Processor processor) { + this.executor = executor; + this.processor = processor; } - @Scheduled(fixedDelay = 9999999999999999L) - public void run() { - Map tablesForExport = new HashMap<>(); - for (Table table : Table.values()) tablesForExport.put(table, null); - ResultContainer resultContainer = ResultContainer.createNewTask(tablesForExport); - for (Stage stage : pipeline) { - StageResult result = stage.process(resultContainer); - if (Arrays.asList(StageResult.ERROR, StageResult.COMPLETE).contains(result)) - break; + public void run(Map tablesForExport) { + for (Map.Entry tableForExport : tablesForExport.entrySet()) { + Table table = tableForExport.getKey(); + ISqlFilter filter = tableForExport.getValue(); + executor.execute(() -> processor.process(ResultContainer.createNewTask(table, filter))); } } } diff --git a/clearing-parent/dbf-exporter/src/main/resources/application.properties b/clearing-parent/dbf-exporter/src/main/resources/application.properties index 600848cda..2c2a05216 100644 --- a/clearing-parent/dbf-exporter/src/main/resources/application.properties +++ b/clearing-parent/dbf-exporter/src/main/resources/application.properties @@ -8,5 +8,6 @@ db.login=clearing db.password=Aa111111 dbf.encoding=cp866 +dbf.threads-count=1 dbf.out-dir=D:\\dbf\\out