Merge remote-tracking branch 'origin/dev' into dev

This commit is contained in:
ialbert 2022-09-21 18:40:36 +03:00
commit 21adcaa308
44 changed files with 1008 additions and 609 deletions

View file

@ -30,26 +30,22 @@
<artifactId>spring-boot-autoconfigure</artifactId>
</dependency>
<!-- JDBC -->
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-jdbc</artifactId>
</dependency>
<dependency>
<groupId>com.mchange</groupId>
<artifactId>c3p0</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
</dependency>
<!-- DBF files -->
<dependency>
<groupId>com.github.albfernandez</groupId>
<artifactId>javadbf</artifactId>
</dependency>
<!-- Own dependencies -->
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>classes</artifactId>
</dependency>
</dependencies>
<build>

View file

@ -1,21 +1,20 @@
package ru.spcex.clearing.dbf.exporter.config;
import com.mchange.v2.c3p0.ComboPooledDataSource;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
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 org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import ru.spcex.clearing.dbf.exporter.logic.stages.ExportFromDB;
import ru.spcex.clearing.dbf.exporter.logic.stages.ExportFromHazelcast;
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;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import javax.sql.DataSource;
import java.beans.PropertyVetoException;
import java.util.LinkedList;
import java.util.List;
@ -31,19 +30,26 @@ public class DBFExporterConfig {
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("taskExecutorHazelcastClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
return createThreadPoolTaskExecutor(1, true);
}
@Bean("dbfJdbcTemplate")
public JdbcTemplate jdbcTemplate(@Qualifier("dbfDataSource") DataSource dataSource) {
return new JdbcTemplate(dataSource);
@Bean("taskExecutorIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
return createThreadPoolTaskExecutor(1, false);
}
@Bean("imdgProvider")
public ImdgProvider imdgProvider(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
AProperties properties) {
HazelcastClientParams params = new HazelcastClientParams();
params.setClusterMembers(properties.getHazelcastClusterMembers());
params.setLogin(properties.getHazelcastLogin());
params.setPassword(properties.getHazelcastPassword());
params.setInstanceName("dbf-exporter");
return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
}
@Bean("pipeline")
@ -51,7 +57,7 @@ public class DBFExporterConfig {
List<Stage> pipeline = new LinkedList<>();
pipeline.add(context.getBean(PrepareDBFFile.class));
pipeline.add(context.getBean(ExportFromDB.class));
pipeline.add(context.getBean(ExportFromHazelcast.class));
return pipeline;
}
@ -68,4 +74,15 @@ public class DBFExporterConfig {
return executor;
}
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
if (maxPoolSz > 2) {
pool.setKeepAliveSeconds(60);
pool.setAllowCoreThreadTimeOut(true);
}
pool.setCorePoolSize(maxPoolSz);
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
return pool;
}
}

View file

@ -10,7 +10,7 @@ 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.IFilter;
import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table;
import ru.spcex.clearing.dbf.exporter.services.DBFExportService;
@ -37,7 +37,7 @@ public class DefaultController implements InitializingBean {
@ResponseBody
public String exportTables() {
log.info("Call export method for exporter controller");
Map<Table, ISqlFilter> tablesForExport = new EnumMap<>(Table.class);
Map<Table, IFilter> tablesForExport = new EnumMap<>(Table.class);
for (Table table : Table.values()) tablesForExport.put(table, null);
dbfExportService.run(tablesForExport);
return "export done, see log";

View file

@ -3,6 +3,6 @@ package ru.spcex.clearing.dbf.exporter.logic.data;
/**
* Фильтр записей для экспорта
*/
public interface ISqlFilter {
String getSqlCondition();
public interface IFilter {
}

View file

@ -11,12 +11,12 @@ import java.util.UUID;
public class ResultContainer {
private UUID uuid;
private Table tableForExport;
private ISqlFilter filter;
private IFilter filter;
private File fileForExport;
protected ResultContainer() {}
public static ResultContainer createNewTask(Table tableForExport, ISqlFilter filter) {
public static ResultContainer createNewTask(Table tableForExport, IFilter filter) {
ResultContainer container = new ResultContainer();
container.tableForExport = tableForExport;
container.filter = filter;
@ -32,11 +32,11 @@ public class ResultContainer {
this.tableForExport = tableForExport;
}
public ISqlFilter getFilter() {
public IFilter getFilter() {
return filter;
}
public void setFilter(ISqlFilter filter) {
public void setFilter(IFilter filter) {
this.filter = filter;
}

View file

@ -2,37 +2,38 @@ package ru.spcex.clearing.dbf.exporter.logic.data.enums;
import com.linuxense.javadbf.DBFDataType;
import java.sql.Types;
import java.math.BigDecimal;
import java.time.LocalDate;
/**
* Типы данных в таблицах.
* Ставит в соответствие типы 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);
VARCHAR(DBFDataType.CHARACTER, String.class),
CHARACTER(DBFDataType.CHARACTER, Character.class),
NUMERIC(DBFDataType.NUMERIC, BigDecimal.class),
DATE(DBFDataType.DATE, LocalDate.class);
private final DBFDataType dbfType;
private final int sqlType;
private final Class<?> javaType;
ColumnType(DBFDataType dbfType, int postgreSqlType) {
ColumnType(DBFDataType dbfType, Class<?> javaType) {
this.dbfType = dbfType;
this.sqlType = postgreSqlType;
this.javaType = javaType;
}
public DBFDataType getDbfType() {
return dbfType;
}
public int getSqlType() {
return sqlType;
public Class<?> getJavaType() {
return javaType;
}
public static ColumnType getForSQLType(int sqlType) {
public static ColumnType getForSQLType(Class<?> sqlType) {
for (ColumnType columnType : values()) {
if (columnType.sqlType == sqlType) return columnType;
if (columnType.javaType == sqlType) return columnType;
}
return null;
}

View file

@ -1,39 +1,49 @@
package ru.spcex.clearing.dbf.exporter.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");
import ru.clearing.classes.statics.data.sdf.*;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.classes.base.SpcexObjectBase;
public enum Table {
S_DF01("DF-01", IMDGDistributedNames.Map_SDf01, SDf01.class),
S_DF02("DF-02", IMDGDistributedNames.Map_SDf02, SDf02.class),
S_DF08("DF-08", IMDGDistributedNames.Map_SDf08, SDf08.class),
S_DF12("DF-12", IMDGDistributedNames.Map_SDf12, SDf12.class),
S_DF16("DF-16", IMDGDistributedNames.Map_SDf16, SDf16.class),
S_DF18("DF-18", IMDGDistributedNames.Map_SDf18, SDf18.class);
/**
* Префикс имени файла для экспорта
*/
private final String filePrefix;
Table(String prefix) {
this.filePrefix = prefix;
/**
* Имя мапы hazelcast
*/
private final String hazelcastMapName;
/**
* Класс объекта
*/
private final Class<? extends SpcexObjectBase> entityClass;
Table(String filePrefix, String hazelcastMapName, Class<? extends SpcexObjectBase> entityClass) {
this.filePrefix = filePrefix;
this.hazelcastMapName = hazelcastMapName;
this.entityClass = entityClass;
}
public boolean fileForThisTable(String filename) {
return filename != null && filename.startsWith(filePrefix);
}
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;
public String getHazelcastMapName() {
return hazelcastMapName;
}
public String getFilePrefix() {
return filePrefix;
}
public Class<? extends SpcexObjectBase> getEntityClass() {
return entityClass;
}
}

View file

@ -1,139 +0,0 @@
package ru.spcex.clearing.dbf.exporter.logic.stages;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import com.linuxense.javadbf.DBFWriter;
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.exporter.exceptions.ConfigException;
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.ColumnType;
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.math.BigDecimal;
import java.nio.charset.Charset;
import java.sql.ResultSet;
import java.time.LocalDate;
import java.time.ZoneOffset;
import java.util.*;
/**
* Выгрузка данных из базы и их запись
*/
@Component
public class ExportFromDB extends Stage implements InitializingBean {
private final AProperties properties;
private final JdbcTemplate dbfJdbcTemplate;
private final Map<Table, DBFField[]> dbfFieldsForTable = new HashMap<>();
private Charset dbfCharset;
public ExportFromDB(@Qualifier("dbfExporterProperties") AProperties properties,
@Qualifier("dbfJdbcTemplate") JdbcTemplate dbfJdbcTemplate) {
this.properties = properties;
this.dbfJdbcTemplate = dbfJdbcTemplate;
}
@Override
public StageResult process(ResultContainer resultContainer) {
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<Map<String, Object>> recordsFromDB = dbfJdbcTemplate.queryForList(sql);
File dbfFile = resultContainer.getFileForExport();
DBFField[] dbfFields = dbfFieldsForTable.get(table);
try (DBFWriter dbfWriter = new DBFWriter(dbfFile, dbfCharset)) {
dbfWriter.setFields(dbfFields);
for (Map<String, Object> 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);
}
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.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 файлов
* 2. Проверяет кодировку из настройки dbf.encoding
*/
@Override
public void afterPropertiesSet() throws Exception {
initDBStructure();
initDBFCharset();
}
private void initDBStructure() throws MetaDataAccessException {
boolean ok = JdbcUtils.extractDatabaseMetaData(Objects.requireNonNull(dbfJdbcTemplate.getDataSource()),
dbMeta -> {
List<String> 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;
List<DBFField> dbfFields = new ArrayList<>();
ResultSet columnNamesRS = dbMeta.getColumns(null, null, tableNameFromResultSet, null);
while (columnNamesRS.next()) {
String name = columnNamesRS.getString("COLUMN_NAME");
ColumnType type = ColumnType.getForSQLType(columnNamesRS.getInt("DATA_TYPE"));
if (type == null) throw new ConfigException("Unsupported ColumnType from DB");
int length = columnNamesRS.getInt("COLUMN_SIZE");
if (type == ColumnType.DATE) length = 8;
int decimalCount = columnNamesRS.getInt("DECIMAL_DIGITS");
DBFField dbfField = new DBFField(name.toUpperCase(), type.getDbfType(), length, decimalCount);
dbfFields.add(dbfField);
}
dbfFieldsForTable.put(currTable, dbfFields.toArray(DBFField[]::new));
}
return true;
});
}
private void initDBFCharset() {
try {
dbfCharset = Charset.forName(properties.getDbfEncoding());
} catch (Exception e) {
throw new ConfigException("В properties файле содержится неизвестная кодировка: " + properties.getDbfEncoding(), e);
}
}
}

View file

@ -0,0 +1,140 @@
package ru.spcex.clearing.dbf.exporter.logic.stages;
import com.linuxense.javadbf.DBFField;
import com.linuxense.javadbf.DBFWriter;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.sdf.*;
import ru.spcex.clearing.dbf.exporter.exceptions.ConfigException;
import ru.spcex.clearing.dbf.exporter.logic.data.IFilter;
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 ru.spcex.clearing.dbf.exporter.services.converters.*;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.io.File;
import java.nio.charset.Charset;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
/**
* Выгрузка данных из мапы hazelcast и их запись в файлы
*/
@Component
public class ExportFromHazelcast extends Stage implements InitializingBean {
private final AProperties properties;
private final ImdgProvider imdgProvider;
private final S_DF01_Converter s_df01_converter;
private final S_DF02_Converter s_df02_converter;
private final S_DF08_Converter s_df08_converter;
private final S_DF12_Converter s_df12_converter;
private final S_DF16_Converter s_df16_converter;
private final S_DF18_Converter s_df18_converter;
private final Map<Table, DBFField[]> dbfFieldsForTable = new HashMap<>();
private Charset dbfCharset;
public ExportFromHazelcast(@Qualifier("dbfExporterProperties") AProperties properties,
@Qualifier("imdgProvider") ImdgProvider imdgProvider,
S_DF01_Converter s_df01_converter,
S_DF02_Converter s_df02_converter,
S_DF08_Converter s_df08_converter,
S_DF12_Converter s_df12_converter,
S_DF16_Converter s_df16_converter,
S_DF18_Converter s_df18_converter) {
this.properties = properties;
this.imdgProvider = imdgProvider;
this.s_df01_converter = s_df01_converter;
this.s_df02_converter = s_df02_converter;
this.s_df08_converter = s_df08_converter;
this.s_df12_converter = s_df12_converter;
this.s_df16_converter = s_df16_converter;
this.s_df18_converter = s_df18_converter;
}
@Override
public StageResult process(ResultContainer resultContainer) {
Objects.requireNonNull(resultContainer.getTableForExport());
Objects.requireNonNull(resultContainer.getFileForExport());
Table table = resultContainer.getTableForExport();
IFilter filter = resultContainer.getFilter();
Imdg<? extends SpcexObjectBase> map = imdgProvider.getImdg(table.getHazelcastMapName(), table.getEntityClass());
File dbfFile = resultContainer.getFileForExport();
boolean writeOk = false;
boolean emptyMap = true;
try (DBFWriter dbfWriter = new DBFWriter(dbfFile, dbfCharset)) {
dbfWriter.setFields(dbfFieldsForTable.get(table));
Collection<? extends SpcexObjectBase> allValues = map.getAllValues();
if (allValues.isEmpty()) {
log.info("uuid {}. Map {} is empty.", resultContainer.getUuid(), resultContainer.getTableForExport().getHazelcastMapName());
return StageResult.COMPLETE;
}
for (SpcexObjectBase value : allValues) {
Object[] values;
if (value instanceof SDf01 sdf01Value) values = s_df01_converter.toObjectArray(sdf01Value);
else if (value instanceof SDf02 sDf02Value) values = s_df02_converter.toObjectArray(sDf02Value);
else if (value instanceof SDf08 sDf08Value) values = s_df08_converter.toObjectArray(sDf08Value);
else if (value instanceof SDf12 sDf12Value) values = s_df12_converter.toObjectArray(sDf12Value);
else if (value instanceof SDf16 sDf16Value) values = s_df16_converter.toObjectArray(sDf16Value);
else if (value instanceof SDf18 sDf18Value) values = s_df18_converter.toObjectArray(sDf18Value);
else throw new Exception("Get unknown object from imdg. Class: " + value.getClass().getSimpleName());
dbfWriter.addRecord(values);
}
writeOk = true;
emptyMap = allValues.isEmpty();
} 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;
} finally {
if (!writeOk || emptyMap) {
boolean deleteOk = resultContainer.getFileForExport().delete();
if (!deleteOk) {
log.warn("Can't delete file {}", resultContainer.getFileForExport().getName());
}
}
}
return StageResult.COMPLETE;
}
/**
* 1. Конструирует структуру DBF файлов для дальнейшей записи
* 2. Проверяет кодировку из настройки dbf.encoding
*/
@Override
public void afterPropertiesSet() {
initExportFileStructure();
initDBFCharset();
}
/**
* Инициализация структуры выходных файлов (заголовки столбцов)
*/
private void initExportFileStructure() {
dbfFieldsForTable.put(Table.S_DF01, s_df01_converter.getDBFHeaders());
dbfFieldsForTable.put(Table.S_DF02, s_df02_converter.getDBFHeaders());
dbfFieldsForTable.put(Table.S_DF08, s_df08_converter.getDBFHeaders());
dbfFieldsForTable.put(Table.S_DF12, s_df12_converter.getDBFHeaders());
dbfFieldsForTable.put(Table.S_DF16, s_df16_converter.getDBFHeaders());
dbfFieldsForTable.put(Table.S_DF18, s_df18_converter.getDBFHeaders());
}
private void initDBFCharset() {
try {
dbfCharset = Charset.forName(properties.getDbfEncoding());
} catch (Exception e) {
throw new ConfigException("В properties файле содержится неизвестная кодировка: " + properties.getDbfEncoding(), e);
}
}
}

View file

@ -6,17 +6,14 @@ import org.springframework.stereotype.Component;
@Component("dbfExporterProperties")
public class AProperties {
@Value("${db.jdbc-url}")
private String jdbcUrl;
@Value("${hazelcast.cluster-members}")
private String hazelcastClusterMembers;
@Value("${db.driver}")
private String dbDriver;
@Value("${hazelcast.login}")
private String hazelcastLogin;
@Value("${db.login}")
private String dbLogin;
@Value("${db.password}")
private String dbPassword;
@Value("${hazelcast.password}")
private String hazelcastPassword;
@Value("${dbf.out-dir}")
private String outDir;
@ -27,36 +24,28 @@ public class AProperties {
@Value("${dbf.threads-count}")
private int threadsCount;
public String getJdbcUrl() {
return jdbcUrl;
public String getHazelcastClusterMembers() {
return hazelcastClusterMembers;
}
public void setJdbcUrl(String jdbcUrl) {
this.jdbcUrl = jdbcUrl;
public void setHazelcastClusterMembers(String hazelcastClusterMembers) {
this.hazelcastClusterMembers = hazelcastClusterMembers;
}
public String getDbLogin() {
return dbLogin;
public String getHazelcastLogin() {
return hazelcastLogin;
}
public void setDbLogin(String dbLogin) {
this.dbLogin = dbLogin;
public void setHazelcastLogin(String hazelcastLogin) {
this.hazelcastLogin = hazelcastLogin;
}
public String getDbPassword() {
return dbPassword;
public String getHazelcastPassword() {
return hazelcastPassword;
}
public void setDbPassword(String dbPassword) {
this.dbPassword = dbPassword;
}
public String getDbDriver() {
return dbDriver;
}
public void setDbDriver(String dbDriver) {
this.dbDriver = dbDriver;
public void setHazelcastPassword(String hazelcastPassword) {
this.hazelcastPassword = hazelcastPassword;
}
public String getOutDir() {

View file

@ -3,7 +3,7 @@ package ru.spcex.clearing.dbf.exporter.services;
import org.springframework.beans.factory.annotation.Qualifier;
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.IFilter;
import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer;
import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table;
import ru.spcex.clearing.dbf.exporter.logic.stages.Processor;
@ -21,10 +21,10 @@ public class DBFExportService {
this.processor = processor;
}
public void run(Map<Table, ISqlFilter> tablesForExport) {
for (Map.Entry<Table, ISqlFilter> tableForExport : tablesForExport.entrySet()) {
public void run(Map<Table, IFilter> tablesForExport) {
for (Map.Entry<Table, IFilter> tableForExport : tablesForExport.entrySet()) {
Table table = tableForExport.getKey();
ISqlFilter filter = tableForExport.getValue();
IFilter filter = tableForExport.getValue();
executor.submit(() -> processor.process(ResultContainer.createNewTask(table, filter)));
}
}

View file

@ -0,0 +1,25 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFField;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import java.sql.Date;
import java.sql.Timestamp;
import java.time.Instant;
import java.time.LocalDate;
public abstract class DFConverter<T extends SpcexObjectBase> {
public abstract Object[] toObjectArray(T entity);
public abstract DBFField[] getDBFHeaders();
/**
* Приведение типов в соответствие (для записи)
* @param src исходный объект
* @return объект другого типа, если тип не поддерживается dbfWriter
*/
protected Object typeMatch(Object src) {
if (src instanceof Instant srcInstant) return Timestamp.from(srcInstant);
if (src instanceof LocalDate srcLocalDate) return Date.valueOf(srcLocalDate);
return src;
}
}

View file

@ -0,0 +1,49 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf01;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF01_Converter extends DFConverter<SDf01> {
@Override
public Object[] toObjectArray(SDf01 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getCurr_code()));
values.add(typeMatch(entity.getAccount()));
values.add(typeMatch(entity.getRemainder()));
values.add(typeMatch(entity.getDeal()));
values.add(typeMatch(entity.getAcc_code()));
values.add(typeMatch(entity.getDat()));
values.add(typeMatch(entity.getMarket()));
values.add(typeMatch(entity.getAcc_name()));
values.add(typeMatch(entity.getAcc_type()));
values.add(typeMatch(entity.getSumengage()));
values.add(typeMatch(entity.getSumunblock()));
values.add(typeMatch(entity.getFile_type()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("curr_code", DBFDataType.CHARACTER, 12));
dbfFields.add(new DBFField("account", DBFDataType.CHARACTER, 35));
dbfFields.add(new DBFField("reminder", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("deal", DBFDataType.CHARACTER, 10));
dbfFields.add(new DBFField("acc_code", DBFDataType.CHARACTER, 5));
dbfFields.add(new DBFField("dat", DBFDataType.CHARACTER, 8));
dbfFields.add(new DBFField("market", DBFDataType.CHARACTER, 1));
dbfFields.add(new DBFField("acc_name", DBFDataType.CHARACTER, 30));
dbfFields.add(new DBFField("acc_type", DBFDataType.CHARACTER, 2));
dbfFields.add(new DBFField("sumengage", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("sumunblock", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("file_type", DBFDataType.CHARACTER, 22));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -0,0 +1,57 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf02;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF02_Converter extends DFConverter<SDf02> {
@Override
public Object[] toObjectArray(SDf02 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getCurr_code()));
values.add(typeMatch(entity.getAccount()));
values.add(typeMatch(entity.getRemainder()));
values.add(typeMatch(entity.getDeal()));
values.add(typeMatch(entity.getAcc_code()));
values.add(typeMatch(entity.getDat()));
values.add(typeMatch(entity.getMarket()));
values.add(typeMatch(entity.getAcc_name()));
values.add(typeMatch(entity.getAcc_type()));
values.add(typeMatch(entity.getSumengage()));
values.add(typeMatch(entity.getSumunblock()));
values.add(typeMatch(entity.getFile_type()));
values.add(typeMatch(entity.getResult()));
values.add(typeMatch(entity.getGenerationTime()));
values.add(typeMatch(entity.getGenerationId()));
values.add(typeMatch(entity.getInSDf01Id()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("curr_code", DBFDataType.CHARACTER, 12));
dbfFields.add(new DBFField("account", DBFDataType.CHARACTER, 35));
dbfFields.add(new DBFField("reminder", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("deal", DBFDataType.CHARACTER, 10));
dbfFields.add(new DBFField("acc_code", DBFDataType.CHARACTER, 5));
dbfFields.add(new DBFField("dat", DBFDataType.CHARACTER, 8));
dbfFields.add(new DBFField("market", DBFDataType.CHARACTER, 1));
dbfFields.add(new DBFField("acc_name", DBFDataType.CHARACTER, 30));
dbfFields.add(new DBFField("acc_type", DBFDataType.CHARACTER, 2));
dbfFields.add(new DBFField("sumengage", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("sumunblock", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("file_type", DBFDataType.CHARACTER, 22));
dbfFields.add(new DBFField("result", DBFDataType.CHARACTER, 3));
dbfFields.add(new DBFField("gen_time", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_id", DBFDataType.NUMERIC));
dbfFields.add(new DBFField("insdf01_id", DBFDataType.NUMERIC));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf08;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF08_Converter extends DFConverter<SDf08> {
@Override
public Object[] toObjectArray(SDf08 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getNumber()));
values.add(typeMatch(entity.getDatetime()));
values.add(typeMatch(entity.getGenerationTime()));
values.add(typeMatch(entity.getGenerationId()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("number", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("datetime", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_time", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_id", DBFDataType.NUMERIC));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -0,0 +1,37 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf12;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF12_Converter extends DFConverter<SDf12> {
@Override
public Object[] toObjectArray(SDf12 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getAccount()));
values.add(typeMatch(entity.getDeal()));
values.add(typeMatch(entity.getStatus()));
values.add(typeMatch(entity.getFileName()));
values.add(typeMatch(entity.getGenerationTime()));
values.add(typeMatch(entity.getGenerationId()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("account", DBFDataType.CHARACTER, 25));
dbfFields.add(new DBFField("deal", DBFDataType.CHARACTER, 25));
dbfFields.add(new DBFField("status", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("file_name", DBFDataType.CHARACTER, 254));
dbfFields.add(new DBFField("gen_time", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_id", DBFDataType.NUMERIC));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -0,0 +1,43 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf16;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF16_Converter extends DFConverter<SDf16> {
@Override
public Object[] toObjectArray(SDf16 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getDate()));
values.add(typeMatch(entity.getAccount()));
values.add(typeMatch(entity.getSum()));
values.add(typeMatch(entity.getMarket()));
values.add(typeMatch(entity.getType()));
values.add(typeMatch(entity.getNumber()));
values.add(typeMatch(entity.getFileName()));
values.add(typeMatch(entity.getGenerationTime()));
values.add(typeMatch(entity.getGenerationId()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("date", DBFDataType.DATE));
dbfFields.add(new DBFField("account", DBFDataType.CHARACTER, 20));
dbfFields.add(new DBFField("sum", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("market", DBFDataType.CHARACTER, 1));
dbfFields.add(new DBFField("type", DBFDataType.CHARACTER, 1));
dbfFields.add(new DBFField("number", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("file_name", DBFDataType.CHARACTER, 254));
dbfFields.add(new DBFField("gen_time", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_id", DBFDataType.NUMERIC));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -0,0 +1,39 @@
package ru.spcex.clearing.dbf.exporter.services.converters;
import com.linuxense.javadbf.DBFDataType;
import com.linuxense.javadbf.DBFField;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.sdf.SDf18;
import java.util.LinkedList;
import java.util.List;
@Service
public class S_DF18_Converter extends DFConverter<SDf18> {
@Override
public Object[] toObjectArray(SDf18 entity) {
List<Object> values = new LinkedList<>();
values.add(typeMatch(entity.getAccount()));
values.add(typeMatch(entity.getDeal()));
values.add(typeMatch(entity.getStatus()));
values.add(typeMatch(entity.getResult()));
values.add(typeMatch(entity.getGenerationTime()));
values.add(typeMatch(entity.getGenerationId()));
values.add(typeMatch(entity.getInSDf12Id()));
return values.toArray(Object[]::new);
}
@Override
public DBFField[] getDBFHeaders() {
List<DBFField> dbfFields = new LinkedList<>();
dbfFields.add(new DBFField("account", DBFDataType.CHARACTER, 25));
dbfFields.add(new DBFField("deal", DBFDataType.CHARACTER, 4));
dbfFields.add(new DBFField("status", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("result", DBFDataType.NUMERIC, 32));
dbfFields.add(new DBFField("gen_time", DBFDataType.DATE));
dbfFields.add(new DBFField("gen_id", DBFDataType.NUMERIC));
dbfFields.add(new DBFField("insdf12_id", DBFDataType.NUMERIC));
return dbfFields.toArray(DBFField[]::new);
}
}

View file

@ -2,10 +2,9 @@ server.port=8080
server.servlet.context-path=/exporter
spring.main.web-application-type=servlet
db.jdbc-url=jdbc:postgresql://10.200.200.133:5432/postgres
db.driver=org.postgresql.Driver
db.login=clearing
db.password=Aa111111
hazelcast.cluster-members=127.0.0.1:5701
hazelcast.login=dev
hazelcast.password=dev-pass
dbf.encoding=cp866
dbf.threads-count=10

View file

@ -1,6 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
<project xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns="http://maven.apache.org/POM/4.0.0"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
@ -50,9 +50,30 @@
<groupId>com.github.albfernandez</groupId>
<artifactId>javadbf</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-imdg-api</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>classes</artifactId>
</dependency>
</dependencies>
<build>
<resources>
<resource>
<directory>src/main/resources</directory>
<excludes>
<exclude>application.properties</exclude>
</excludes>
<filtering>false</filtering>
</resource>
</resources>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
@ -69,6 +90,7 @@
</configuration>
</plugin>
</plugins>
</build>
</project>

View file

@ -1,13 +1,11 @@
package ru.spcex.clearing.dbf.importer.config;
import com.mchange.v2.c3p0.ComboPooledDataSource;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
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 org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import ru.spcex.clearing.dbf.importer.logic.stages.ImportToDB;
import ru.spcex.clearing.dbf.importer.logic.stages.LoadFileFromDisk;
@ -15,8 +13,6 @@ import ru.spcex.clearing.dbf.importer.logic.stages.Stage;
import ru.spcex.clearing.dbf.importer.logic.stages.ValidateFields;
import ru.spcex.clearing.dbf.importer.properties.AProperties;
import javax.sql.DataSource;
import java.beans.PropertyVetoException;
import java.util.LinkedList;
import java.util.List;
@ -27,26 +23,12 @@ public class DBFImporterConfig {
private final AProperties properties;
private final ApplicationContext context;
public DBFImporterConfig(@Qualifier("dbfImporterProperties") 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<Stage> pipeline() {
List<Stage> pipeline = new LinkedList<>();

View file

@ -0,0 +1,46 @@
package ru.spcex.clearing.dbf.importer.config;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@Configuration
public class ImporterImdgConfig {
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
if (maxPoolSz > 2) {
pool.setKeepAliveSeconds(60);
pool.setAllowCoreThreadTimeOut(true);
}
pool.setCorePoolSize(maxPoolSz);
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
return pool;
}
@Bean(name = "taskExecutorHazelcastClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
return createThreadPoolTaskExecutor(1, true);
}
@Bean(name = "taskExecutorIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
return createThreadPoolTaskExecutor(1, false);
}
@Autowired
@Bean
public HazelcastService imdgProvider(
@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
ImportDBFServiceSettings settings
) {
return new HazelcastService(taskExecutorHazelcastClientInitializer,
taskExecutorIdGeneratorAwaiter,
settings.getHazelcast());
}
}

View file

@ -0,0 +1,21 @@
package ru.spcex.clearing.dbf.importer.config.settings;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.PropertySource;
import org.springframework.stereotype.Component;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
@Component
@PropertySource("file:${spring.config.location}/application.properties")
@ConfigurationProperties("import-dbf-service")
public class ImportDBFServiceSettings {
private HazelcastClientParams hazelcast;
public HazelcastClientParams getHazelcast() {
return hazelcast;
}
public void setHazelcast(HazelcastClientParams hazelcast) {
this.hazelcast = hazelcast;
}
}

View file

@ -1,18 +1,18 @@
package ru.spcex.clearing.dbf.importer.logic.data;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import java.util.EnumMap;
import java.util.Map;
public class DBStructureDBF {
private Map<Table, TableStructureDBF> tables = new EnumMap<>(Table.class);
private Map<ETable, TableStructureDBF> tables = new EnumMap<>(ETable.class);
public Map<Table, TableStructureDBF> getTables() {
public Map<ETable, TableStructureDBF> getTables() {
return tables;
}
public void setTables(Map<Table, TableStructureDBF> tables) {
public void setTables(Map<ETable, TableStructureDBF> tables) {
this.tables = tables;
}
}

View file

@ -1,6 +1,6 @@
package ru.spcex.clearing.dbf.importer.logic.data;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import java.io.File;
import java.util.UUID;
@ -10,14 +10,15 @@ import java.util.UUID;
*/
public class ResultContainer {
private UUID uuid;
private Table dbfTable;
private ETable dbfTable;
private File dbfFile;
private byte[] dbfSource;
protected ResultContainer() {}
protected ResultContainer() {
}
public static ResultContainer createNewTask(Table dbfTable, File dbfFile) {
public static ResultContainer createNewTask(ETable dbfTable, File dbfFile) {
ResultContainer container = new ResultContainer();
container.dbfTable = dbfTable;
container.dbfFile = dbfFile;
@ -25,11 +26,11 @@ public class ResultContainer {
return container;
}
public Table getDbfTable() {
public ETable getDbfTable() {
return dbfTable;
}
public void setDbfTable(Table dbfTable) {
public void setDbfTable(ETable dbfTable) {
this.dbfTable = dbfTable;
}

View file

@ -1,35 +1,41 @@
package ru.spcex.clearing.dbf.importer.logic.data.enums;
public enum Table {
public enum ETable {
DF_01("DF-01"),
DF_02("DF-02"),
DF_03("DF-03"),
DF_04("DF-04"),
DF_05("DF-05"),
DF_06("DF-06");
DF_08("DF-08"),
DF_09("DF-09"),
DF_10("DF-10"),
DF_11("DF-11"),
DF_12("DF-12"),
DF_16("DF-16"),
DF_17("DF-17"),
DF_18("DF-18");
private final String prefix;
Table(String prefix) {
ETable(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()) {
public static ETable getTableForFilename(String filename) {
for (ETable table : ETable.values()) {
if (table.fileForThisTable(filename)) return table;
}
return null;
}
public static Table tableForName(String name) {
for (Table table : values()) {
public static ETable tableForName(String name) {
for (ETable table : values()) {
if (name.equalsIgnoreCase(table.name()))
return table;
}
return null;
}
public boolean fileForThisTable(String filename) {
return filename != null && filename.startsWith(prefix);
}
}

View file

@ -0,0 +1,36 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
public abstract class AbstractTable<T extends SpcexObjectBase> {
private final String prefix;
private final Class<T> clazz;
private final String nameOfMap;
protected HazelcastService hazelcastService;
protected ImdgHazelcast<T> map;
protected AbstractTable(String prefix, Class<T> clazz, String nameOfMap) {
this.prefix = prefix;
this.clazz = clazz;
this.nameOfMap = nameOfMap;
}
public void setHazelcastService(HazelcastService hazelcastService) {
this.hazelcastService = hazelcastService;
}
public abstract T getEntity(Object[] entity);
public void injectEntity(T obj) {
if (map == null) {
bootMap();
}
map.insert(obj);
}
private void bootMap() {
map = (ImdgHazelcast<T>) hazelcastService.getImdg(nameOfMap, clazz);
}
}

View file

@ -0,0 +1,36 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf01;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
public class SDf01Table extends AbstractTable<SDf01> {
private static final String PREFIX = ETable.DF_01.name();
private static final Class<SDf01> CLAZZ = SDf01.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf01;
public SDf01Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf01 getEntity(Object[] entity) {
SDf01 result = new SDf01();
result.setCurr_code((String) entity[0]);
result.setAccount((String) entity[1]);
result.setRemainder((String) entity[2]);
result.setDeal((String) entity[3]);
result.setAcc_code((String) entity[4]);
result.setDat((String) entity[5]);
result.setMarket((String) entity[6]);
result.setAcc_name((String) entity[7]);
result.setAcc_type((String) entity[8]);
result.setSumengage((String) entity[9]);
result.setSumunblock((String) entity[10]);
result.setFile_type((String) entity[11]);
return result;
}
}

View file

@ -0,0 +1,42 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf02;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.utils.time.TimeUtil;
public class SDf02Table extends AbstractTable<SDf02> {
private static final String PREFIX = ETable.DF_02.name();
private static final Class<SDf02> CLAZZ = SDf02.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf02;
public SDf02Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf02 getEntity(Object[] entity) {
SDf02 result = new SDf02();
result.setCurr_code((String) entity[0]);
result.setAccount((String) entity[1]);
result.setRemainder((String) entity[2]);
result.setDeal((String) entity[3]);
result.setAcc_code((String) entity[4]);
result.setDat((String) entity[5]);
result.setMarket((String) entity[6]);
result.setAcc_name((String) entity[7]);
result.setAcc_type((String) entity[8]);
result.setSumengage((String) entity[9]);
result.setSumunblock((String) entity[10]);
result.setFile_type((String) entity[11]);
result.setResult((String) entity[12]);
result.setGenerationTime(TimeUtil.strToInstant((String) entity[13]));
result.setGenerationId((Long) entity[14]);
result.setInSDf01Id((Long) entity[15]);
return result;
}
}

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf08;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
public class SDf08Table extends AbstractTable<SDf08> {
private static final String PREFIX = ETable.DF_08.name();
private static final Class<SDf08> CLAZZ = SDf08.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf08;
public SDf08Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf08 getEntity(Object[] entity) {
SDf08 result = new SDf08();
result.setNumber((BigDecimal) entity[0]);
result.setDatetime(TimeUtil.strToInstant((String) entity[1]));
result.setGenerationTime(TimeUtil.strToInstant((String) entity[2]));
result.setGenerationId((Long) entity[3]);
return result;
}
}

View file

@ -0,0 +1,34 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf12;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
public class SDf12Table extends AbstractTable<SDf12> {
private static final String PREFIX = ETable.DF_12.name();
private static final Class<SDf12> CLAZZ = SDf12.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf12;
public SDf12Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf12 getEntity(Object[] entity) {
SDf12 result = new SDf12();
result.setAccount((String) entity[0]);
result.setDeal((String) entity[1]);
result.setStatus((BigDecimal) entity[2]);
result.setFileName((String) entity[3]);
result.setGenerationTime(TimeUtil.strToInstant((String) entity[4]));
result.setGenerationId((Long) entity[5]);
return result;
}
}

View file

@ -0,0 +1,39 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf16;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
public class SDf16Table extends AbstractTable<SDf16> {
private static final String PREFIX = ETable.DF_16.name();
private static final Class<SDf16> CLAZZ = SDf16.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf16;
public SDf16Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf16 getEntity(Object[] entity) {
SDf16 result = new SDf16();
result.setDate(TimeUtil.strToInstant((String) entity[0]));
result.setAccount((String) entity[1]);
result.setSum((BigDecimal) entity[2]);
result.setMarket((String) entity[3]);
result.setType((String) entity[4]);
result.setINN((BigDecimal) entity[5]);
result.setBIC((BigDecimal) entity[6]);
result.setSPEC((String) entity[7]);
result.setNumber((BigDecimal) entity[8]);
result.setFileName((String) entity[9]);
result.setGenerationTime(TimeUtil.strToInstant((String) entity[10]));
result.setGenerationId((Long) entity[11]);
return result;
}
}

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.dbf.importer.logic.data.tables;
import ru.clearing.classes.statics.data.sdf.SDf18;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
public class SDf18Table extends AbstractTable<SDf18> {
private static final String PREFIX = ETable.DF_18.name();
private static final Class<SDf18> CLAZZ = SDf18.class;
private static final String NAME_OF_HZ_MAP = IMDGDistributedNames.Map_SDf18;
public SDf18Table() {
super(PREFIX, CLAZZ, NAME_OF_HZ_MAP);
}
@Override
public SDf18 getEntity(Object[] entity) {
SDf18 result = new SDf18();
result.setAccount((String) entity[0]);
result.setDeal((String) entity[1]);
result.setStatus((BigDecimal) entity[2]);
result.setResult((BigDecimal) entity[3]);
result.setGenerationTime(TimeUtil.strToInstant((String) entity[4]));
result.setGenerationId((Long) entity[5]);
result.setInSDf12Id((Long) entity[6]);
return result;
}
}

View file

@ -1,28 +1,21 @@
package ru.spcex.clearing.dbf.importer.logic.stages;
import com.linuxense.javadbf.DBFField;
import com.linuxense.javadbf.DBFReader;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ColumnType;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.dbf.importer.logic.data.enums.StageResult;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.logic.data.tables.AbstractTable;
import ru.spcex.clearing.dbf.importer.properties.AProperties;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
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;
import static ru.spcex.clearing.dbf.importer.utils.Const.MAPPING_ENUM_TABLE_OBJECT_TABLE;
/**
* Заливка проверенных данных в базу
@ -30,81 +23,32 @@ import java.util.Objects;
@Component
public class ImportToDB extends Stage {
private final AProperties properties;
private final JdbcTemplate dbfJdbcTemplate;
private final HazelcastService hazelcastService;
public ImportToDB(@Qualifier("dbfImporterProperties") AProperties properties,
@Qualifier("dbfJdbcTemplate") JdbcTemplate dbfJdbcTemplate) {
HazelcastService hazelcastService) {
this.properties = properties;
this.dbfJdbcTemplate = dbfJdbcTemplate;
this.hazelcastService = hazelcastService;
}
@Override
public StageResult process(ResultContainer resultContainer) {
Objects.requireNonNull(resultContainer.getDbfTable());
Objects.requireNonNull(resultContainer.getDbfSource());
Table currTable = resultContainer.getDbfTable();
ETable 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<String> columnNames = new ArrayList<>(columnCount);
List<Integer> 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());
AbstractTable table = MAPPING_ENUM_TABLE_OBJECT_TABLE.get(currTable);
table.setHazelcastService(hazelcastService);
for (int i = 0; i < dbfReader.getRecordCount(); i++) {
Object[] entity = dbfReader.nextRecord();
table.injectEntity(table.getEntity(entity));
}
int recordCount = dbfReader.getRecordCount();
List<Object[]> 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("uuid {}. Imported records: {}", resultContainer.getUuid(), rowsInserted.length);
}
} catch (IOException e) {
log.error(String.format("uuid %s. Can't read source file %s.",
resultContainer.getUuid(),
resultContainer.getDbfFile().getName()), e);
return StageResult.ERROR;
} catch (DataAccessException e) {
log.error(String.format("uuid %s. Can't insert records", resultContainer.getUuid()), e);
} catch (IOException exception) {
log.warn(exception.getMessage());
return StageResult.ERROR;
}

View file

@ -32,7 +32,7 @@ public class LoadFileFromDisk extends Stage {
log.debug("uuid {}. Read all bytes from source file {}", taskUuid, resultContainer.getDbfFile().getName());
fileBytes = Files.readAllBytes(Paths.get(dbfFile.getAbsolutePath()));
if (fileBytes.length == 0) throw new IOException("Empty file");
if (properties.deleteSrcFiles()) Files.delete(dbfFile.toPath());
if (properties.isDeleteSrcFiles()) Files.delete(dbfFile.toPath());
} catch (IOException e) {
log.error(String.format("uuid %s. Can't read file %s", taskUuid, resultContainer.getDbfFile().getName()), e);
return StageResult.ERROR;

View file

@ -1,194 +1,39 @@
package ru.spcex.clearing.dbf.importer.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.importer.exceptions.ConfigException;
import ru.spcex.clearing.dbf.importer.logic.data.ColumnStructureDBF;
import ru.spcex.clearing.dbf.importer.logic.data.DBStructureDBF;
import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer;
import ru.spcex.clearing.dbf.importer.logic.data.TableStructureDBF;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ColumnType;
import ru.spcex.clearing.dbf.importer.logic.data.enums.StageResult;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.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;
import java.util.UUID;
/**
* Фильтрация источника на предмет соответствия полей
*/
@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;
public ValidateFields(AProperties properties) {
this.properties = properties;
}
@Override
public StageResult process(ResultContainer resultContainer) {
Objects.requireNonNull(resultContainer.getDbfTable());
Objects.requireNonNull(resultContainer.getDbfSource());
UUID taskUuid = resultContainer.getUuid();
Table currTable = resultContainer.getDbfTable();
byte[] source = resultContainer.getDbfSource();
TableStructureDBF currTableStructure = dbStructureDBF.getTables().get(currTable);
if (currTableStructure == null) {
log.error("uuid {}. Unknown table {}.", taskUuid, currTable);
return StageResult.ERROR;
}
log.info("uuid {}. Check source {}", taskUuid, currTable);
try (InputStream is = new ByteArrayInputStream(source);
DBFReader dbfReader = new DBFReader(is, dbfCharset)) {
int columnCount = dbfReader.getFieldCount();
if (columnCount != currTableStructure.getColumns().size()) {
log.error("uuid {}. Source for table {} doesn't match DB columns count. Source column count: {}. DB column count: {}.",
taskUuid,
currTable,
columnCount,
currTableStructure.getColumns().size());
return StageResult.ERROR;
}
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.error("uuid {}. Source for table {} contain unknown field {}. ",
taskUuid,
currTable,
currFieldName);
valid = false;
continue;
}
DBFDataType dbfDataType = currField.getType();
if (!currColumnStructure.getType().getDbfType().equals(dbfDataType)) {
log.error("uuid {}. Field {} from source {} has wrong data type (actual {}, expected {}).",
taskUuid,
currFieldName,
currTable,
dbfDataType != null ? dbfDataType.name() : "null",
currColumnStructure.getType());
valid = false;
continue;
}
int fieldLength = currField.getLength();
if (fieldLength != currColumnStructure.getLength()) {
if (!dbfDataType.equals(DBFDataType.DATE)) {
log.error("uuid {}. Field {} from source {} has wrong length (actual {}, expected {}).",
taskUuid,
currFieldName,
currTable,
fieldLength,
currColumnStructure.getLength());
valid = false;
continue;
} else {
if (fieldLength > currColumnStructure.getLength()) {
log.error("uuid {}. Field {} with DATE type has wrong length (length from file > length db field)",
taskUuid,
currFieldName);
valid = false;
continue;
}
}
}
int decimalCount = currField.getDecimalCount();
if (decimalCount != currColumnStructure.getDecimalDigits()) {
log.error("uuid {}. Field {} from source {} has wrong decimal count (actual {}, expected {}).",
taskUuid,
currFieldName,
currTable,
decimalCount,
currColumnStructure.getDecimalDigits());
valid = false;
}
if (!valid) {
return StageResult.ERROR;
}
}
} catch (IOException e) {
log.error(String.format("uuid %s. Can't read source for table %s.", taskUuid, currTable), e);
return StageResult.ERROR;
}
return StageResult.OK;
}
/**
* 1. Собирает из БД названия столбцов для дальнейшей валидации
* 2. Проверяет кодировку из настройки dbf.encoding-source
* 1. Проверяет кодировку из настройки 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<String> 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());
@ -196,5 +41,4 @@ public class ValidateFields extends Stage implements InitializingBean {
throw new ConfigException("Unknown encoding from properties: " + properties.getDbfEncoding(), e);
}
}
}

View file

@ -6,18 +6,6 @@ import org.springframework.stereotype.Component;
@Component("dbfImporterProperties")
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;
@ -33,75 +21,23 @@ public class AProperties {
@Value("${dbf.threads-count}")
private int threadsCount;
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 boolean deleteSrcFiles() {
public boolean isDeleteSrcFiles() {
return deleteSrcFiles;
}
public void setDeleteSrcFiles(boolean deleteSrcFiles) {
this.deleteSrcFiles = deleteSrcFiles;
}
public String getDbfEncoding() {
return dbfEncoding;
}
public void setDbfEncoding(String dbfEncoding) {
this.dbfEncoding = dbfEncoding;
}
public int getInsertBatchSize() {
return insertBatchSize;
}
public void setInsertBatchSize(int insertBatchSize) {
this.insertBatchSize = insertBatchSize;
}
public int getThreadsCount() {
return threadsCount;
}
public void setThreadsCount(int threadsCount) {
this.threadsCount = threadsCount;
}
}

View file

@ -6,7 +6,7 @@ import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.dbf.importer.logic.data.ResultContainer;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.dbf.importer.logic.stages.Processor;
import java.io.File;
@ -30,9 +30,9 @@ public class DBFImporterService {
@Scheduled(cron = "${dbf.check-src-dir-cron}")
public void run() {
Map<Table, List<File>> newFiles = fileChecker.checkNewFiles();
for (Map.Entry<Table, List<File>> newFilesEntry : newFiles.entrySet()) {
Table currTable = newFilesEntry.getKey();
Map<ETable, List<File>> newFiles = fileChecker.checkNewFiles();
for (Map.Entry<ETable, List<File>> newFilesEntry : newFiles.entrySet()) {
ETable currTable = newFilesEntry.getKey();
List<File> fileList = newFilesEntry.getValue();
for (File dbfFile : fileList) {
executorService.execute(() -> processor.process(ResultContainer.createNewTask(currTable, dbfFile)));

View file

@ -1,7 +1,7 @@
package ru.spcex.clearing.dbf.importer.services;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.dbf.importer.logic.data.enums.Table;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.dbf.importer.properties.AProperties;
import java.io.File;
@ -15,8 +15,8 @@ public class FileChecker {
this.properties = properties;
}
public Map<Table, List<File>> checkNewFiles() {
Map<Table, List<File>> newFiles = new EnumMap<>(Table.class);
public Map<ETable, List<File>> checkNewFiles() {
Map<ETable, List<File>> newFiles = new EnumMap<>(ETable.class);
String srcDir = properties.getSrcDir();
List<File> dbfFiles = lsDBF(srcDir);
@ -25,7 +25,7 @@ public class FileChecker {
for (File dbfFile : dbfFiles) {
if (dbfFile.isDirectory()) continue;
Table currTable = Table.getTableForFilename(dbfFile.getName());
ETable currTable = ETable.getTableForFilename(dbfFile.getName());
if (currTable == null) continue;
List<File> currList = newFiles.computeIfAbsent(currTable, list -> new LinkedList<>());

View file

@ -0,0 +1,23 @@
package ru.spcex.clearing.dbf.importer.utils;
import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable;
import ru.spcex.clearing.dbf.importer.logic.data.tables.*;
import java.util.HashMap;
import java.util.Map;
public class Const {
public static final Map<ETable, AbstractTable> MAPPING_ENUM_TABLE_OBJECT_TABLE = new HashMap<>();
static {
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_01, new SDf01Table());
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_02, new SDf02Table());
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_08, new SDf08Table());
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_12, new SDf12Table());
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_16, new SDf16Table());
MAPPING_ENUM_TABLE_OBJECT_TABLE.put(ETable.DF_18, new SDf18Table());
}
}

View file

@ -1,17 +1,16 @@
server.port=8080
server.servlet.context-path=/importer
spring.main.web-application-type=servlet
db.jdbc-url=jdbc:postgresql://10.200.200.133:5432/postgres
db.driver=org.postgresql.Driver
db.login=clearing
db.password=Aa111111
dbf.check-src-dir-cron=* * * * 1 ?
dbf.encoding-source=cp866
dbf.insert-batch-size=100
dbf.delete-src-files=false
dbf.src-dir=D:\\dbf\\
dbf.threads-count=10
dbf.threads-count=10
import-dbf-service.hazelcast.cluster-members=127.0.0.1
import-dbf-service.hazelcast.login=dev
import-dbf-service.hazelcast.password=dev-pass

View file

@ -13,6 +13,7 @@ import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.utils.log.ExceptionUtils;
import java.util.*;
public class ImdgHazelcast<T extends SpcexObjectBase> implements Imdg<T> {
private final Logger log = LoggerFactory.getLogger(getClass());
private IdGenerator idGenerator;
@ -109,7 +110,7 @@ public class ImdgHazelcast<T extends SpcexObjectBase> implements Imdg<T> {
Set<Long> ids = map.keySet(or);
Iterator<Long> idIterator = ids.iterator();
Collection<T> searchResult = new ArrayList<>();
while(idIterator.hasNext()) {
while (idIterator.hasNext()) {
T element = map.get(idIterator.next());
if (element != null) {
searchResult.add(element);
@ -118,6 +119,11 @@ public class ImdgHazelcast<T extends SpcexObjectBase> implements Imdg<T> {
return searchResult;
}
@Override
public Long nextIDSequenceFor() {
return idGenerator.newId();
}
public IMap<Long, T> getMap() {
return map;
}

View file

@ -105,7 +105,7 @@ public class ImdgTransactionalHazelcast<T extends SpcexObjectBase> implements Im
Set<Long> ids = map.keySet(or);
Iterator<Long> idIterator = ids.iterator();
Collection<T> searchResult = new ArrayList<>();
while(idIterator.hasNext()) {
while (idIterator.hasNext()) {
T element = map.get(idIterator.next());
if (element != null) {
searchResult.add(element);
@ -145,4 +145,9 @@ public class ImdgTransactionalHazelcast<T extends SpcexObjectBase> implements Im
public void setMapName(String mapName) {
this.mapName = mapName;
}
@Override
public Long nextIDSequenceFor() {
return idGenerator.newId();
}
}

View file

@ -4,8 +4,10 @@ import java.time.Instant;
import java.time.LocalDate;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.Objects;
public class TimeUtil {
public static final DateTimeFormatter PROPERTY_DATE_FORMATTER = DateTimeFormatter.ofPattern("dd.MM.yyyy");
public static ZoneId zone = ZoneId.systemDefault();
/**
@ -27,4 +29,17 @@ public class TimeUtil {
public static Instant localDateToInstant(LocalDate date) {
return date.atStartOfDay(zone).toInstant();
}
public static LocalDate strToLocalDate(String date) {
try {
LocalDate parsed = LocalDate.parse(date, PROPERTY_DATE_FORMATTER);
return parsed;
} catch (Throwable e) {
return null;
}
}
public static Instant strToInstant(String date) {
return localDateToInstant(Objects.requireNonNull(strToLocalDate(date)));
}
}