From 73d3250ce0e2bd3168144f1e57786b37eaa99d9d Mon Sep 17 00:00:00 2001 From: AKurakin Date: Tue, 25 Mar 2025 14:37:44 +0300 Subject: [PATCH] trade-importer http://jira.mfd.msk:8088/browse/CLS-820 --- clearing-parent/trade-importer/pom.xml | 4 ++ .../services/TradeImporterService.java | 65 ++++++++++++++++--- .../src/main/resources/application.properties | 15 +++-- 3 files changed, 70 insertions(+), 14 deletions(-) diff --git a/clearing-parent/trade-importer/pom.xml b/clearing-parent/trade-importer/pom.xml index 615b5c5e1..5d7466a06 100644 --- a/clearing-parent/trade-importer/pom.xml +++ b/clearing-parent/trade-importer/pom.xml @@ -43,6 +43,10 @@ com.microsoft.sqlserver mssql-jdbc + + org.postgresql + postgresql + com.mchange c3p0 diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java index 427140412..fc3eed521 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java @@ -78,8 +78,8 @@ public class TradeImporterService { } log.info("Start process import STrades from DB. byCommand={}", byCommand); if (log.isDebugEnabled()) { - String countQuery = String.format("SELECT count(*) FROM %s.Trades", schema); - log.trace("Query: {}", countQuery); + String countQuery = String.format("SELECT count(*) FROM \"%s\".\"Trades\"", schema); + log.debug("Query: {}", countQuery); Collection tradesCountFromDB = jdbcTemplate.query(countQuery, (resultSet, i) -> resultSet.getObject(1, Long.class)); log.debug("Count of Trades in db: {}. Will be loading.", tradesCountFromDB); @@ -91,7 +91,7 @@ public class TradeImporterService { AtomicLong minCreatedId = new AtomicLong(Long.MAX_VALUE); final long logProgressTime = 60_000L; // интервал вывода в лог каждую минуту AtomicLong logTime = new AtomicLong(System.currentTimeMillis() + logProgressTime); - String query = String.format("SELECT * FROM %s.Trades", schema); + String query = String.format("SELECT * FROM \"%s\".\"Trades\"", schema); log.debug("Load from DB STrades. Query: {}", query); Instant currentInstant = Instant.now(); jdbcTemplate.query(query, resultSet -> { @@ -216,9 +216,56 @@ public class TradeImporterService { && StringUtils.hasText(trades.getOperation()); } + /** + * В случае, если в таблицах postgresql будет numeric + */ + protected Long safeReadAsLong(ResultSet rs, String column) throws SQLException { + int columnType = rs.getMetaData().getColumnType(rs.findColumn(column)); + switch (columnType) { + case java.sql.Types.NUMERIC: + BigDecimal d = rs.getBigDecimal(column); + if (d == null) + return null; + if (d.remainder(BigDecimal.ONE).compareTo(BigDecimal.ZERO) != 0) { + log.warn("Column {} is double value {}. Expected Long value.", column, d); + } + return d.longValueExact(); + case java.sql.Types.INTEGER: + int i = rs.getInt(column); + if (rs.wasNull()) + return null; + return (long)i; + case java.sql.Types.BIGINT: + long l = rs.getLong(column); + if (rs.wasNull()) + return null; + return l; + default: + return rs.getLong(column); // try by JDBC convert + } + } + + protected String safeReadAsString(ResultSet rs, String column) throws SQLException { + int columnType = rs.getMetaData().getColumnType(rs.findColumn(column)); + switch (columnType) { + case java.sql.Types.INTEGER: + int i = rs.getInt(column); + if (rs.wasNull()) + return null; + return String.valueOf(i); + case java.sql.Types.BIGINT: + long l = rs.getLong(column); + if (rs.wasNull()) + return null; + return String.valueOf(l); + default: + return rs.getString(column); + } + } + public STrades readSTrades(ResultSet resultSet) throws SQLException { STrades sTrades = new STrades(); - sTrades.setTradeNum(resultSet.getObject("TradeNum", Long.class)); + sTrades.setTradeNum(safeReadAsLong(resultSet,"TradeNum")); String operation = resultSet.getObject("Operation", String.class); sTrades.setOperation(mapOperation(operation)); sTrades.setClassCode(resultSet.getObject("ClassCode", String.class)); @@ -227,7 +274,7 @@ public class TradeImporterService { sTrades.setAccruedint(resultSet.getObject("Accruedint", BigDecimal.class)); sTrades.setAccruedint2(resultSet.getObject("Accruedint2", BigDecimal.class)); sTrades.setLowerDiscount(resultSet.getObject("LowerDiscount", BigDecimal.class)); - sTrades.setOrderNum(resultSet.getObject("OrderNum", Long.class)); + sTrades.setOrderNum(safeReadAsLong(resultSet,"OrderNum")); sTrades.setPrice(resultSet.getObject("Price", BigDecimal.class)); sTrades.setPrice2(resultSet.getObject("Price2", BigDecimal.class)); sTrades.setRepoRate(resultSet.getObject("RepoRate", BigDecimal.class)); @@ -241,7 +288,7 @@ public class TradeImporterService { sTrades.setQty(resultSet.getObject("Qty", BigDecimal.class)); sTrades.setQtyPcs(resultSet.getObject("Qty_pcs", BigDecimal.class)); sTrades.setTradeDateTime(getInstantFromTimestamp(resultSet, "TradeDateTime")); - sTrades.setRepoTerm(resultSet.getObject("RepoTerm", Long.class)); + sTrades.setRepoTerm(safeReadAsLong(resultSet, "RepoTerm")); sTrades.setClearingCommission(resultSet.getObject("ClearingCommission", BigDecimal.class)); sTrades.setExchangeCommission(resultSet.getObject("ExchangeCommission", BigDecimal.class)); sTrades.setTechCenterCommission(resultSet.getObject("TechCenterCommission", BigDecimal.class)); @@ -259,10 +306,10 @@ public class TradeImporterService { sTrades.setSecName(resultSet.getObject("SecName", String.class)); sTrades.setSettleDate(getLocalDateFromSqlDate(resultSet, "SettleDate")); sTrades.setSettleCurrency(resultSet.getObject("SettleCurrency", String.class)); - sTrades.setTradeTimeMs(resultSet.getObject("TradeTimeMs", Long.class)); + sTrades.setTradeTimeMs(safeReadAsLong(resultSet, "TradeTimeMs")); sTrades.setBankAccId(resultSet.getObject("BankAccId", String.class)); - sTrades.setKind(resultSet.getObject("Kind", String.class)); - sTrades.setLinkedTrade(resultSet.getObject("LinkedTrade", Long.class)); + sTrades.setKind(safeReadAsString(resultSet,"Kind")); + sTrades.setLinkedTrade(safeReadAsLong(resultSet, "LinkedTrade")); return sTrades; } diff --git a/clearing-parent/trade-importer/src/main/resources/application.properties b/clearing-parent/trade-importer/src/main/resources/application.properties index 7b7f8f668..f5facf6d0 100644 --- a/clearing-parent/trade-importer/src/main/resources/application.properties +++ b/clearing-parent/trade-importer/src/main/resources/application.properties @@ -3,11 +3,16 @@ spring.main.web-application-type=none #trade-importer.cron.load-from-db-cron=0 0/5 * * * ? - каждые 5 минут trade-importer.cron.load-from-db-cron=0 0/5 * * * ? -trade-importer.database.login=sa -trade-importer.database.password=Aa123456 -trade-importer.database.schema=SPVB_TS -trade-importer.database.url=jdbc:sqlserver://10.200.200.144:1433;database=ni; -trade-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver +#trade-importer.database.login=sa +#trade-importer.database.password=Aa123456 +#trade-importer.database.schema=SPVB_TS +#trade-importer.database.url=jdbc:sqlserver://10.200.200.144:1433;database=ni; +#trade-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver +trade-importer.database.login=clearing +trade-importer.database.password=Aa111111 +trade-importer.database.schema=CBRInfo +trade-importer.database.url=jdbc:postgresql://10.200.200.133:5432/clearing?currentSchema=CBRInfo +trade-importer.database.driver=org.postgresql.Driver trade-importer.hazelcast.cluster-members=127.0.0.1:5701 trade-importer.hazelcast.login=dev