diff --git a/clearing-parent/dbf-loader/pom.xml b/clearing-parent/dbf-loader/pom.xml new file mode 100644 index 000000000..7f5b5e8f1 --- /dev/null +++ b/clearing-parent/dbf-loader/pom.xml @@ -0,0 +1,87 @@ + + + + 4.0.0 + + dbf-loader + dbf-loader + DBF loader for DBF files + SPCEX-1.0.0.0 + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.boot + spring-boot-autoconfigure + + + + + org.springframework + spring-jdbc + + + com.mchange + c3p0 + 0.9.5.2 + + + org.postgresql + postgresql + + + + + com.github.albfernandez + javadbf + 1.13.1 + + + + + + jar/${project.artifactId} + + + org.apache.maven.plugins + maven-jar-plugin + + + + + ru.spcex.clearing.dbf.loader.DBFLoaderApplication + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + true + exec + + + + + + \ No newline at end of file diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/DBFLoaderApplication.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/DBFLoaderApplication.java new file mode 100644 index 000000000..942869b03 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/DBFLoaderApplication.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.dbf.loader; + +import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.builder.SpringApplicationBuilder; + +@SpringBootApplication +public class DBFLoaderApplication { + public static void main(String[] args) { + try { + SpringApplicationBuilder builder = new SpringApplicationBuilder(DBFLoaderApplication.class); + builder.run(args); + } catch (Throwable e) { + LoggerFactory.getLogger(DBFLoaderApplication.class).error("DBF-Loader start failed: {} -> {}", e.getClass().getSimpleName(), e.getMessage()); + System.exit(-1); + } + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/config/DBFLoaderConfig.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/config/DBFLoaderConfig.java new file mode 100644 index 000000000..3a385d450 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/config/DBFLoaderConfig.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.dbf.loader.config; + +import com.mchange.v2.c3p0.ComboPooledDataSource; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.ApplicationContext; +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 ru.spcex.clearing.dbf.loader.logic.stages.ImportToDB; +import ru.spcex.clearing.dbf.loader.logic.stages.Stage; +import ru.spcex.clearing.dbf.loader.logic.stages.ValidateFields; +import ru.spcex.clearing.dbf.loader.properties.AProperties; + +import javax.sql.DataSource; +import java.beans.PropertyVetoException; +import java.util.LinkedList; +import java.util.List; + +@Configuration +@ComponentScan(basePackages = {"ru.spcex.clearing.dbf.loader"}) +public class DBFLoaderConfig { + private final AProperties properties; + private final ApplicationContext context; + + public DBFLoaderConfig(@Qualifier("dbfLoaderProperties") AProperties properties, ApplicationContext context) { + this.properties = properties; + this.context = context; + } + + @Bean("dbfDataSource") + public DataSource dataSource() throws PropertyVetoException { + ComboPooledDataSource dataSource = new ComboPooledDataSource(); + dataSource.setDriverClass(properties.getDbDriver()); + dataSource.setJdbcUrl(properties.getJdbcUrl()); + dataSource.setUser(properties.getDbLogin()); + dataSource.setPassword(properties.getDbPassword()); + return dataSource; + } + + @Bean("dbfJdbcTemplate") + public JdbcTemplate jdbcTemplate(@Qualifier("dbfDataSource") DataSource dataSource) { + return new JdbcTemplate(dataSource); + } + + @Bean("pipeline") + public List pipeline() { + List pipeline = new LinkedList<>(); + + pipeline.add(context.getBean(ValidateFields.class)); + pipeline.add(context.getBean(ImportToDB.class)); + + return pipeline; + } + +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/exceptions/ConfigException.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/exceptions/ConfigException.java new file mode 100644 index 000000000..95aa73185 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/exceptions/ConfigException.java @@ -0,0 +1,6 @@ +package ru.spcex.clearing.dbf.loader.exceptions; + +public class ConfigException extends RuntimeException { + public ConfigException(String msg) { super(msg); } + public ConfigException(String msg, Throwable cause) { super(msg, cause); } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ColumnStructureDBF.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ColumnStructureDBF.java new file mode 100644 index 000000000..78e051c77 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ColumnStructureDBF.java @@ -0,0 +1,51 @@ +package ru.spcex.clearing.dbf.loader.logic.data; + +import ru.spcex.clearing.dbf.loader.logic.data.enums.ColumnType; + +public class ColumnStructureDBF { + private String name; + private ColumnType type; + private int length; + private int decimalDigits; + private String comment; + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + + public ColumnType getType() { + return type; + } + + public void setType(ColumnType type) { + this.type = type; + } + + public int getLength() { + return length; + } + + public void setLength(int length) { + this.length = length; + } + + public int getDecimalDigits() { + return decimalDigits; + } + + public void setDecimalDigits(int decimalDigits) { + this.decimalDigits = decimalDigits; + } + + public String getComment() { + return comment; + } + + public void setComment(String comment) { + this.comment = comment; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/DBStructureDBF.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/DBStructureDBF.java new file mode 100644 index 000000000..b384fb518 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/DBStructureDBF.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.dbf.loader.logic.data; + +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; + +import java.util.HashMap; +import java.util.Map; + +public class DBStructureDBF { + private Map tables = new HashMap<>(); + + public Map getTables() { + return tables; + } + + public void setTables(Map tables) { + this.tables = tables; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ResultContainer.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ResultContainer.java new file mode 100644 index 000000000..f6c1264a2 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/ResultContainer.java @@ -0,0 +1,36 @@ +package ru.spcex.clearing.dbf.loader.logic.data; + +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; + +/** + * Контейнер для передачи результата между стадиями + */ +public class ResultContainer { + private Table dbfTable; + private byte[] dbfSource; + + protected ResultContainer() {} + + public static ResultContainer createNewTask(Table dbfTable, byte[] dbfSource) { + ResultContainer container = new ResultContainer(); + container.dbfTable = dbfTable; + container.dbfSource = dbfSource; + return container; + } + + public Table getDbfTable() { + return dbfTable; + } + + public void setDbfTable(Table dbfTable) { + this.dbfTable = dbfTable; + } + + public byte[] getDbfSource() { + return dbfSource; + } + + public void setDbfSource(byte[] dbfSource) { + this.dbfSource = dbfSource; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/TableStructureDBF.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/TableStructureDBF.java new file mode 100644 index 000000000..cb5c20728 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/TableStructureDBF.java @@ -0,0 +1,25 @@ +package ru.spcex.clearing.dbf.loader.logic.data; + +import java.util.HashMap; +import java.util.Map; + +public class TableStructureDBF { + private Map columns = new HashMap<>(); + + public Map getColumns() { + return columns; + } + + public void setColumns(Map columns) { + this.columns = columns; + } + + public ColumnStructureDBF getIgnoreCase(String columnName) { + for (Map.Entry entry : columns.entrySet()) { + String currColumnName = entry.getKey(); + if (currColumnName.equalsIgnoreCase(columnName)) return entry.getValue(); + } + return null; + } + +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/ColumnType.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/ColumnType.java new file mode 100644 index 000000000..2fab3a8cb --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/ColumnType.java @@ -0,0 +1,46 @@ +package ru.spcex.clearing.dbf.loader.logic.data.enums; + +import com.linuxense.javadbf.DBFDataType; + +import java.sql.Types; + +/** + * Типы данных в таблицах. + * Ставит в соответствие типы PostGRE и DBF + */ +public enum ColumnType { + VARCHAR(DBFDataType.CHARACTER, Types.VARCHAR), + CHARACTER(DBFDataType.CHARACTER, Types.CHAR), + NUMERIC(DBFDataType.NUMERIC, Types.NUMERIC), + DATE(DBFDataType.DATE, Types.DATE); + + private final DBFDataType dbfType; + private final int sqlType; + + ColumnType(DBFDataType dbfType, int postgreSqlType) { + this.dbfType = dbfType; + this.sqlType = postgreSqlType; + } + + public DBFDataType getDbfType() { + return dbfType; + } + + public int getSqlType() { + return sqlType; + } + + public static ColumnType getForSQLType(int sqlType) { + for (ColumnType columnType : values()) { + if (columnType.sqlType == sqlType) return columnType; + } + return null; + } + + public static ColumnType getForDBFType(DBFDataType dbfType) { + for (ColumnType columnType : values()) { + if (columnType.dbfType == dbfType) return columnType; + } + return null; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/StageResult.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/StageResult.java new file mode 100644 index 000000000..9558147d2 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/StageResult.java @@ -0,0 +1,7 @@ +package ru.spcex.clearing.dbf.loader.logic.data.enums; + +public enum StageResult { + OK, + ERROR, + COMPLETE +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/Table.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/Table.java new file mode 100644 index 000000000..536e48625 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/data/enums/Table.java @@ -0,0 +1,35 @@ +package ru.spcex.clearing.dbf.loader.logic.data.enums; + +public enum Table { + DF_01("DF-01"), + DF_02("DF-02"), + DF_03("DF-03"), + DF_04("DF-04"), + DF_05("DF-05"), + DF_06("DF-06"); + + private final String prefix; + + Table(String prefix) { + this.prefix = prefix; + } + + public boolean fileForThisTable(String filename) { + return filename != null && filename.startsWith(prefix); + } + + public static Table getTableForFilename(String filename) { + for (Table table : Table.values()) { + if (table.fileForThisTable(filename)) return table; + } + return null; + } + + public static Table tableForName(String name) { + for (Table table : values()) { + if (name.equalsIgnoreCase(table.name())) + return table; + } + return null; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ImportToDB.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ImportToDB.java new file mode 100644 index 000000000..82245314c --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ImportToDB.java @@ -0,0 +1,107 @@ +package ru.spcex.clearing.dbf.loader.logic.stages; + +import com.linuxense.javadbf.DBFField; +import com.linuxense.javadbf.DBFReader; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.jdbc.core.BatchPreparedStatementSetter; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.dbf.loader.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.loader.logic.data.enums.ColumnType; +import ru.spcex.clearing.dbf.loader.logic.data.enums.StageResult; +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; +import ru.spcex.clearing.dbf.loader.properties.AProperties; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.Charset; +import java.sql.PreparedStatement; +import java.sql.SQLException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Objects; + +/** + * Заливка проверенных данных в базу + */ +@Component +public class ImportToDB extends Stage { + private final AProperties properties; + private final JdbcTemplate dbfJdbcTemplate; + + public ImportToDB(@Qualifier("dbfLoaderProperties") AProperties properties, + @Qualifier("dbfJdbcTemplate") JdbcTemplate dbfJdbcTemplate) { + this.properties = properties; + this.dbfJdbcTemplate = dbfJdbcTemplate; + } + + @Override + public StageResult process(ResultContainer resultContainer) { + Objects.requireNonNull(resultContainer.getDbfTable()); + Objects.requireNonNull(resultContainer.getDbfSource()); + + Table currTable = resultContainer.getDbfTable(); + byte[] source = resultContainer.getDbfSource(); + Charset sourceCharset = Charset.forName(properties.getDbfEncoding()); + + try (InputStream is = new ByteArrayInputStream(source); + DBFReader dbfReader = new DBFReader(is, sourceCharset)) { + int columnCount = dbfReader.getFieldCount(); + + List columnNames = new ArrayList<>(columnCount); + List columnTypes = new ArrayList<>(columnCount); + for (int columnIdx = 0; columnIdx < columnCount; columnIdx++) { + DBFField currField = dbfReader.getField(columnIdx); + String fieldName = currField.getName(); + ColumnType columnType = ColumnType.getForDBFType(currField.getType()); + columnNames.add(fieldName); + columnTypes.add(columnType.getSqlType()); + } + + int recordCount = dbfReader.getRecordCount(); + List values = new ArrayList<>(recordCount); + for (int rowIdx = 0; rowIdx < recordCount; rowIdx++) { + Object[] dbfRow = dbfReader.nextRecord(); + values.add(dbfRow); + } + + int maxBatchSize = properties.getInsertBatchSize(); + int batchCount = recordCount / maxBatchSize; + int batchRem = recordCount % maxBatchSize; + for (int batchIdx = 0; batchIdx < batchCount + 1; batchIdx++) { + if (batchIdx == batchCount && batchRem == 0) continue; + final Integer currentBatchIdx = batchIdx; + + int[] rowsInserted = dbfJdbcTemplate.batchUpdate( + "insert into " + currTable.name() + + " (" + String.join(",", columnNames) + ")" + + " values (" + String.join(",", Collections.nCopies(columnCount, "?")) + ")", + new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + Object[] currValuesRow = values.get(currentBatchIdx * maxBatchSize + i); + int parameterIndex = 1; + for (Object object : currValuesRow) { + ps.setObject(parameterIndex, object, columnTypes.get(parameterIndex - 1)); + parameterIndex++; + } + } + + @Override + public int getBatchSize() { + return currentBatchIdx == batchCount ? batchRem : maxBatchSize; + } + }); + log.info("Imported records: " + rowsInserted.length); + } + + } catch (IOException e) { + e.printStackTrace(); + } + + + return StageResult.OK; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/Stage.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/Stage.java new file mode 100644 index 000000000..978cc1f84 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/Stage.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.dbf.loader.logic.stages; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.dbf.loader.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.loader.logic.data.enums.StageResult; + +public abstract class Stage { + protected Logger log = LoggerFactory.getLogger(getClass()); + + public abstract StageResult process(ResultContainer resultContainer); +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ValidateFields.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ValidateFields.java new file mode 100644 index 000000000..d3b18cc34 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/logic/stages/ValidateFields.java @@ -0,0 +1,189 @@ +package ru.spcex.clearing.dbf.loader.logic.stages; + +import com.linuxense.javadbf.DBFDataType; +import com.linuxense.javadbf.DBFField; +import com.linuxense.javadbf.DBFReader; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.support.JdbcUtils; +import org.springframework.jdbc.support.MetaDataAccessException; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.dbf.loader.exceptions.ConfigException; +import ru.spcex.clearing.dbf.loader.logic.data.ColumnStructureDBF; +import ru.spcex.clearing.dbf.loader.logic.data.DBStructureDBF; +import ru.spcex.clearing.dbf.loader.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.loader.logic.data.TableStructureDBF; +import ru.spcex.clearing.dbf.loader.logic.data.enums.ColumnType; +import ru.spcex.clearing.dbf.loader.logic.data.enums.StageResult; +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; +import ru.spcex.clearing.dbf.loader.properties.AProperties; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.Charset; +import java.sql.ResultSet; +import java.util.LinkedList; +import java.util.List; +import java.util.Objects; + +/** + * Фильтрация источника на предмет соответствия полей + */ +@Component +public class ValidateFields extends Stage implements InitializingBean { + + private final JdbcTemplate dbfJdbcTemplate; + private final AProperties properties; + + private DBStructureDBF dbStructureDBF; + private Charset dbfCharset; + + public ValidateFields(@Qualifier("dbfJdbcTemplate") JdbcTemplate dbfJdbcTemplate, AProperties properties) { + this.dbfJdbcTemplate = dbfJdbcTemplate; + this.properties = properties; + } + + @Override + public StageResult process(ResultContainer resultContainer) { + Objects.requireNonNull(resultContainer.getDbfTable()); + Objects.requireNonNull(resultContainer.getDbfSource()); + Table currTable = resultContainer.getDbfTable(); + byte[] source = resultContainer.getDbfSource(); + + TableStructureDBF currTableStructure = dbStructureDBF.getTables().get(currTable); + if (currTableStructure == null) { + log.warn("Unknown table {}, table was skipped.", currTable); + return StageResult.COMPLETE; + } + + log.info("Check source " + currTable); + + try (InputStream is = new ByteArrayInputStream(source); + DBFReader dbfReader = new DBFReader(is, dbfCharset)) { + + int columnCount = dbfReader.getFieldCount(); + if (columnCount != currTableStructure.getColumns().size()) { + log.warn("Source for table {} doesn't match DB columns count. Source column count: {}. DB column count: {}. Source was skipped.", + currTable, + columnCount, + currTableStructure.getColumns().size()); + return StageResult.COMPLETE; + } + + boolean valid = true; + for (int columnIdx = 0; columnIdx < columnCount; columnIdx++) { + DBFField currField = dbfReader.getField(columnIdx); + + String currFieldName = currField.getName(); + ColumnStructureDBF currColumnStructure = currTableStructure.getIgnoreCase(currFieldName); + if (currColumnStructure == null) { + log.warn("Source for table {} contain unknown field {}. ", currTable, currFieldName); + valid = false; + continue; + } + + DBFDataType dbfDataType = currField.getType(); + if (!currColumnStructure.getType().getDbfType().equals(dbfDataType)) { + log.warn("Field {} from source {} has wrong data type (actual {}, expected {}).", + currFieldName, + currTable, + dbfDataType.name(), + currColumnStructure.getType()); + valid = false; + continue; + } + + int fieldLength = currField.getLength(); + if (fieldLength != currColumnStructure.getLength()) { + + if (!dbfDataType.equals(DBFDataType.DATE)) { + log.warn("Field {} from source {} has wrong length (actual {}, expected {}).", + currFieldName, + currTable, + fieldLength, + currColumnStructure.getLength()); + valid = false; + continue; + } else { + if (fieldLength > currColumnStructure.getLength()) { + log.warn("..."); + } + } + } + + int decimalCount = currField.getDecimalCount(); + if (decimalCount != currColumnStructure.getDecimalDigits()) { + log.warn("Field {} from source {} has wrong decimal count (actual {}, expected {}).", + currFieldName, + currTable, + decimalCount, + currColumnStructure.getDecimalDigits()); + valid = false; + } + + if (!valid) { + log.warn("Source was skipped."); + return StageResult.COMPLETE; + } + } + } catch (IOException e) { + log.error("Can't read source for table {}. Source was skipped.", currTable); + return StageResult.ERROR; + } + + return StageResult.OK; + } + + + /** + * 1. Собирает из БД названия столбцов для дальнейшей валидации + * 2. Проверяет кодировку из настройки dbf.encoding-source + */ + @Override + public void afterPropertiesSet() throws Exception { + initDBStructure(); + initDBFCharset(); + } + + private void initDBStructure() throws MetaDataAccessException { + dbStructureDBF = JdbcUtils.extractDatabaseMetaData(Objects.requireNonNull(dbfJdbcTemplate.getDataSource()), + dbMeta -> { + DBStructureDBF dbStructureDBF = new DBStructureDBF(); + List tableNamesFromResultSet = new LinkedList<>(); + ResultSet tableNamesRS = dbMeta.getTables(null, null, "%", new String[]{"TABLE"}); + while (tableNamesRS.next()) { + tableNamesFromResultSet.add(tableNamesRS.getString("TABLE_NAME")); + } + for (String tableNameFromResultSet : tableNamesFromResultSet) { + Table currTable = Table.tableForName(tableNameFromResultSet); + if (currTable == null) continue; + TableStructureDBF tableStructureDBF = new TableStructureDBF(); + ResultSet columnNamesRS = dbMeta.getColumns(null, null, tableNameFromResultSet, null); + while (columnNamesRS.next()) { + ColumnStructureDBF columnStructureDBF = new ColumnStructureDBF(); + columnStructureDBF.setName(columnNamesRS.getString("COLUMN_NAME")); + columnStructureDBF.setType(ColumnType.getForSQLType(columnNamesRS.getInt("DATA_TYPE"))); + columnStructureDBF.setLength(columnNamesRS.getInt("COLUMN_SIZE")); + columnStructureDBF.setDecimalDigits(columnNamesRS.getInt("DECIMAL_DIGITS")); + columnStructureDBF.setComment(columnNamesRS.getString("REMARKS")); + tableStructureDBF.getColumns().put(columnStructureDBF.getName(), columnStructureDBF); + } + assert !dbStructureDBF.getTables().containsKey(currTable); + + dbStructureDBF.getTables().put(currTable, tableStructureDBF); + } + return dbStructureDBF; + }); + } + + private void initDBFCharset() { + try { + dbfCharset = Charset.forName(properties.getDbfEncoding()); + } catch (Exception e) { + throw new ConfigException("В properties файле содержится неизвестная кодировка: " + properties.getDbfEncoding(), e); + } + } + +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/properties/AProperties.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/properties/AProperties.java new file mode 100644 index 000000000..687c9fef1 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/properties/AProperties.java @@ -0,0 +1,121 @@ +package ru.spcex.clearing.dbf.loader.properties; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; + +@Component("dbfLoaderProperties") +@PropertySource(value = {"classpath:dbf_loader.properties"}) +@PropertySource(value = {"file:dbf_loader.properties"}, ignoreResourceNotFound = true) +public class AProperties { + + @Value("${db.jdbc-url}") + private String jdbcUrl; + + @Value("${db.driver}") + private String dbDriver; + + @Value("${db.login}") + private String dbLogin; + + @Value("${db.password}") + private String dbPassword; + + @Value("${dbf.src-dir}") + private String srcDir; + + @Value("${dbf.delete-src-files}") + private boolean deleteSrcFiles = true; + + @Value("${dbf.check-source-timeout}") + private long checkSourceTimeout; + + @Value("${dbf.execute-timeout}") + private long executeTimeout; + + @Value("${dbf.encoding-source}") + private String dbfEncoding; + + @Value("${dbf.insert-batch-size}") + private int insertBatchSize; + + public String getJdbcUrl() { + return jdbcUrl; + } + + public void setJdbcUrl(String jdbcUrl) { + this.jdbcUrl = jdbcUrl; + } + + public String getDbLogin() { + return dbLogin; + } + + public void setDbLogin(String dbLogin) { + this.dbLogin = dbLogin; + } + + public String getDbPassword() { + return dbPassword; + } + + public void setDbPassword(String dbPassword) { + this.dbPassword = dbPassword; + } + + public String getDbDriver() { + return dbDriver; + } + + public void setDbDriver(String dbDriver) { + this.dbDriver = dbDriver; + } + + public String getSrcDir() { + return srcDir; + } + + public void setSrcDir(String srcDir) { + this.srcDir = srcDir; + } + + public long getCheckSourceTimeout() { + return checkSourceTimeout; + } + + public void setCheckSourceTimeout(long checkSourceTimeout) { + this.checkSourceTimeout = checkSourceTimeout; + } + + public boolean deleteSrcFiles() { + return deleteSrcFiles; + } + + public void setDeleteSrcFiles(boolean deleteSrcFiles) { + this.deleteSrcFiles = deleteSrcFiles; + } + + public String getDbfEncoding() { + return dbfEncoding; + } + + public void setDbfEncoding(String dbfEncoding) { + this.dbfEncoding = dbfEncoding; + } + + public long getExecuteTimeout() { + return executeTimeout; + } + + public void setExecuteTimeout(long executeTimeout) { + this.executeTimeout = executeTimeout; + } + + public int getInsertBatchSize() { + return insertBatchSize; + } + + public void setInsertBatchSize(int insertBatchSize) { + this.insertBatchSize = insertBatchSize; + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/DBFLoaderService.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/DBFLoaderService.java new file mode 100644 index 000000000..359a4b3dd --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/DBFLoaderService.java @@ -0,0 +1,41 @@ +package ru.spcex.clearing.dbf.loader.services; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.dbf.loader.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.loader.logic.data.enums.StageResult; +import ru.spcex.clearing.dbf.loader.logic.stages.Stage; +import ru.spcex.clearing.dbf.loader.properties.AProperties; + +import java.util.Arrays; +import java.util.List; + +@Service +@EnableScheduling +public class DBFLoaderService { + private final AProperties properties; + private final MessageListener messageListener; + private final List pipeline; + + public DBFLoaderService(AProperties properties, + @Qualifier("fileMessageListener") MessageListener messageListener, + @Qualifier("pipeline") List pipeline) { + this.properties = properties; + this.messageListener = messageListener; + this.pipeline = pipeline; + } + + @Scheduled(fixedDelayString = "${dbf.execute-timeout}") + public void run() { + List newTaskList = messageListener.getTasks(); + for (ResultContainer newTask : newTaskList) { + for (Stage stage : pipeline) { + StageResult result = stage.process(newTask); + if (Arrays.asList(StageResult.ERROR, StageResult.COMPLETE).contains(result)) + break; + } + } + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/FileMessageListener.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/FileMessageListener.java new file mode 100644 index 000000000..394e680c4 --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/FileMessageListener.java @@ -0,0 +1,87 @@ +package ru.spcex.clearing.dbf.loader.services; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.dbf.loader.exceptions.ConfigException; +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; +import ru.spcex.clearing.dbf.loader.properties.AProperties; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Paths; +import java.util.Arrays; +import java.util.LinkedList; +import java.util.List; + +@Service("fileMessageListener") +@EnableScheduling +public class FileMessageListener extends MessageListener implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final AProperties properties; + + public FileMessageListener(AProperties properties) { + this.properties = properties; + } + + @Override + @Scheduled(fixedDelayString = "${dbf.check-source-timeout}") + public void checkMessage() { + String srcDir = properties.getSrcDir(); + List dbfFiles = lsDBF(srcDir); + if (dbfFiles.isEmpty()) return; + + for (File dbfFile : dbfFiles) { + Table currTable = Table.getTableForFilename(dbfFile.getName()); + if (currTable == null) continue; + + byte[] fileBytes; + try { + fileBytes = Files.readAllBytes(Paths.get(dbfFile.getPath())); + if (properties.deleteSrcFiles()) { + boolean deleteOk = dbfFile.delete(); + if (!deleteOk) { + log.warn("Can't remove source file {}. File was skipped.", dbfFile.getPath()); + continue; + } + } + } catch (IOException e) { + log.error("Can't read file {}. File was skipped.", dbfFile.getPath()); + continue; + } + + List currList = sources.computeIfAbsent(currTable, k -> new LinkedList<>()); + currList.add(fileBytes); + + } + } + + private List lsDBF(String dbfDirPath) { + File dbfDir = new File(dbfDirPath); + File[] dbfFiles = dbfDir.listFiles((dir, name) -> { + int formatPosition = name.lastIndexOf("."); + if (formatPosition == -1 || formatPosition == name.length() - 1) return false; + return "dbf".equalsIgnoreCase(name.substring(formatPosition + 1)); + }); + + List resultFiles = new LinkedList<>(); + if (dbfFiles != null && dbfFiles.length >= 1) { + resultFiles.addAll(Arrays.asList(dbfFiles)); + } + + return resultFiles; + } + + @Override + public void afterPropertiesSet() { + String dbfSourceDirPath = properties.getSrcDir(); + File dbfSourceDir = new File(dbfSourceDirPath); + + if (!dbfSourceDir.exists()) throw new ConfigException(String.format("DBF directory (%s) doesn't exist", dbfSourceDirPath)); + if (!dbfSourceDir.isDirectory()) throw new ConfigException(String.format("DBF path (%s) isn't directory", dbfSourceDirPath)); + } +} diff --git a/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/MessageListener.java b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/MessageListener.java new file mode 100644 index 000000000..fb93941de --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/java/ru/spcex/clearing/dbf/loader/services/MessageListener.java @@ -0,0 +1,35 @@ +package ru.spcex.clearing.dbf.loader.services; + +import org.springframework.stereotype.Service; +import ru.spcex.clearing.dbf.loader.logic.data.ResultContainer; +import ru.spcex.clearing.dbf.loader.logic.data.enums.Table; + +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +@Service +public abstract class MessageListener { + protected Map> sources = new ConcurrentHashMap<>(); + + public List getTasks() { + List taskList = new LinkedList<>(); + for (Map.Entry> sourceEntry : sources.entrySet()) { + Table currTable = sourceEntry.getKey(); + List sourceList = sourceEntry.getValue(); + for (byte[] sourceBytes : sourceList) { + ResultContainer currResultContainer = ResultContainer.createNewTask(currTable, sourceBytes); + taskList.add(currResultContainer); + } + sources.remove(currTable); + } + return taskList; + } + + /** + * Перегрузить для поставки sources в очередь + */ + public abstract void checkMessage(); + +} diff --git a/clearing-parent/dbf-loader/src/main/resources/dbf_loader.properties b/clearing-parent/dbf-loader/src/main/resources/dbf_loader.properties new file mode 100644 index 000000000..dd9c2002a --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/resources/dbf_loader.properties @@ -0,0 +1,12 @@ +db.jdbc-url=jdbc:postgresql://10.200.200.133:5432/postgres +db.driver=org.postgresql.Driver +db.login=clearing +db.password=Aa111111 + +dbf.check-source-timeout=300000000 +dbf.execute-timeout=300000000 +dbf.encoding-source=cp866 +dbf.insert-batch-size=100 + +dbf.delete-src-files=false +dbf.src-dir=D:\\dbf\\ diff --git a/clearing-parent/dbf-loader/src/main/resources/logback.xml b/clearing-parent/dbf-loader/src/main/resources/logback.xml new file mode 100644 index 000000000..560b5061c --- /dev/null +++ b/clearing-parent/dbf-loader/src/main/resources/logback.xml @@ -0,0 +1,37 @@ + + + + + UTF-8 + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + + ../logs/dbf_loader.log + + UTF-8 + %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n + + + + ../logs/dbf_loader.%i.log + + 1 + 10 + + + 500MB + + + + + + + + + + + + + \ No newline at end of file diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 4f1813f10..2fc9ebc13 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -21,6 +21,7 @@ backend-api storage db-scripts + dbf-loader