Browse Source

refactored sqlUpgradeService implementation (#2395)

* refactored sqlUpgradeService implementation

* fix typo

* change string constant name

* add ability to re-init chunks for upgrade timescale
pull/2412/head
ShvaykaD 7 years ago
committed by Andrew Shvayka
parent
commit
bdee8951c4
  1. 52
      application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql
  2. 52
      application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql
  3. 124
      application/src/main/java/org/thingsboard/server/service/install/AbstractSqlTsDatabaseUpgradeService.java
  4. 127
      application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/install/SqlAbstractDatabaseSchemaService.java
  6. 31
      application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java
  7. 147
      application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseUpgradeService.java
  8. 68
      application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseSchemaService.java
  9. 125
      application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java
  10. 8
      application/src/main/resources/thingsboard.yml
  11. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  12. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  13. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/InsertLatestTsRepository.java
  14. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java
  15. 6
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java
  16. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java
  17. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java
  18. 8
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java
  19. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java
  20. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java
  21. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  22. 4
      dao/src/main/resources/sql/schema-timescale.sql

52
application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql

@ -160,18 +160,18 @@ DECLARE
insert_size CONSTANT integer := 10000;
insert_counter integer DEFAULT 0;
insert_record RECORD;
insert_cursor CURSOR FOR SELECT CONCAT(first, '-', second, '-1', third, '-', fourth, '-', fifth)::uuid AS entity_id,
substrings.key AS key,
substrings.ts AS ts,
substrings.bool_v AS bool_v,
substrings.str_v AS str_v,
substrings.long_v AS long_v,
substrings.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(entity_id, 8, 8) AS first,
SUBSTRING(entity_id, 4, 4) AS second,
SUBSTRING(entity_id, 1, 3) AS third,
SUBSTRING(entity_id, 16, 4) AS fourth,
SUBSTRING(entity_id, 20) AS fifth,
insert_cursor CURSOR FOR SELECT CONCAT(first_part_uuid, '-', second_part_uuid, '-1', third_part_uuid, '-', fourth_part_uuid, '-', fifth_part_uuid)::uuid AS entity_id,
ts_kv_records.key AS key,
ts_kv_records.ts AS ts,
ts_kv_records.bool_v AS bool_v,
ts_kv_records.str_v AS str_v,
ts_kv_records.long_v AS long_v,
ts_kv_records.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(entity_id, 8, 8) AS first_part_uuid,
SUBSTRING(entity_id, 4, 4) AS second_part_uuid,
SUBSTRING(entity_id, 1, 3) AS third_part_uuid,
SUBSTRING(entity_id, 16, 4) AS fourth_part_uuid,
SUBSTRING(entity_id, 20) AS fifth_part_uuid,
key_id AS key,
ts,
bool_v,
@ -179,7 +179,7 @@ DECLARE
long_v,
dbl_v
FROM ts_kv_old
INNER JOIN ts_kv_dictionary ON (ts_kv_old.key = ts_kv_dictionary.key)) AS substrings;
INNER JOIN ts_kv_dictionary ON (ts_kv_old.key = ts_kv_dictionary.key)) AS ts_kv_records;
BEGIN
OPEN insert_cursor;
LOOP
@ -208,18 +208,18 @@ DECLARE
insert_size CONSTANT integer := 10000;
insert_counter integer DEFAULT 0;
insert_record RECORD;
insert_cursor CURSOR FOR SELECT CONCAT(first, '-', second, '-1', third, '-', fourth, '-', fifth)::uuid AS entity_id,
substrings.key AS key,
substrings.ts AS ts,
substrings.bool_v AS bool_v,
substrings.str_v AS str_v,
substrings.long_v AS long_v,
substrings.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(entity_id, 8, 8) AS first,
SUBSTRING(entity_id, 4, 4) AS second,
SUBSTRING(entity_id, 1, 3) AS third,
SUBSTRING(entity_id, 16, 4) AS fourth,
SUBSTRING(entity_id, 20) AS fifth,
insert_cursor CURSOR FOR SELECT CONCAT(first_part_uuid, '-', second_part_uuid, '-1', third_part_uuid, '-', fourth_part_uuid, '-', fifth_part_uuid)::uuid AS entity_id,
ts_kv_latest_records.key AS key,
ts_kv_latest_records.ts AS ts,
ts_kv_latest_records.bool_v AS bool_v,
ts_kv_latest_records.str_v AS str_v,
ts_kv_latest_records.long_v AS long_v,
ts_kv_latest_records.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(entity_id, 8, 8) AS first_part_uuid,
SUBSTRING(entity_id, 4, 4) AS second_part_uuid,
SUBSTRING(entity_id, 1, 3) AS third_part_uuid,
SUBSTRING(entity_id, 16, 4) AS fourth_part_uuid,
SUBSTRING(entity_id, 20) AS fifth_part_uuid,
key_id AS key,
ts,
bool_v,
@ -227,7 +227,7 @@ DECLARE
long_v,
dbl_v
FROM ts_kv_latest_old
INNER JOIN ts_kv_dictionary ON (ts_kv_latest_old.key = ts_kv_dictionary.key)) AS substrings;
INNER JOIN ts_kv_dictionary ON (ts_kv_latest_old.key = ts_kv_dictionary.key)) AS ts_kv_latest_records;
BEGIN
OPEN insert_cursor;
LOOP

52
application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql

@ -38,9 +38,9 @@ BEGIN
END;
$$ LANGUAGE 'plpgsql';
-- select create_tenant_ts_kv_table_copy();
-- select create_new_tenant_ts_kv_table();
CREATE OR REPLACE FUNCTION create_tenant_ts_kv_table_copy() RETURNS VOID AS $$
CREATE OR REPLACE FUNCTION create_new_tenant_ts_kv_table() RETURNS VOID AS $$
BEGIN
ALTER TABLE tenant_ts_kv
@ -59,7 +59,7 @@ BEGIN
ADD CONSTRAINT tenant_ts_kv_pkey PRIMARY KEY(tenant_id, entity_id, key, ts);
ALTER INDEX idx_tenant_ts_kv RENAME TO idx_tenant_ts_kv_old;
ALTER INDEX tenant_ts_kv_ts_idx RENAME TO tenant_ts_kv_ts_idx_old;
PERFORM create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => 86400000, if_not_exists => true);
-- PERFORM create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => 86400000, if_not_exists => true);
CREATE INDEX IF NOT EXISTS idx_tenant_ts_kv ON tenant_ts_kv(tenant_id, entity_id, key, ts);
END;
$$ LANGUAGE 'plpgsql';
@ -132,24 +132,24 @@ DECLARE
insert_size CONSTANT integer := 10000;
insert_counter integer DEFAULT 0;
insert_record RECORD;
insert_cursor CURSOR FOR SELECT CONCAT(tenant_id_first, '-', tenant_id_second, '-1', tenant_id_third, '-', tenant_id_fourth, '-', tenant_id_fifth)::uuid AS tenant_id,
CONCAT(entity_id_first, '-', entity_id_second, '-1', entity_id_third, '-', entity_id_fourth, '-', entity_id_fifth)::uuid AS entity_id,
substrings.key AS key,
substrings.ts AS ts,
substrings.bool_v AS bool_v,
substrings.str_v AS str_v,
substrings.long_v AS long_v,
substrings.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(tenant_id, 8, 8) AS tenant_id_first,
SUBSTRING(tenant_id, 4, 4) AS tenant_id_second,
SUBSTRING(tenant_id, 1, 3) AS tenant_id_third,
SUBSTRING(tenant_id, 16, 4) AS tenant_id_fourth,
SUBSTRING(tenant_id, 20) AS tenant_id_fifth,
SUBSTRING(entity_id, 8, 8) AS entity_id_first,
SUBSTRING(entity_id, 4, 4) AS entity_id_second,
SUBSTRING(entity_id, 1, 3) AS entity_id_third,
SUBSTRING(entity_id, 16, 4) AS entity_id_fourth,
SUBSTRING(entity_id, 20) AS entity_id_fifth,
insert_cursor CURSOR FOR SELECT CONCAT(tenant_id_first_part_uuid, '-', tenant_id_second_part_uuid, '-1', tenant_id_third_part_uuid, '-', tenant_id_fourth_part_uuid, '-', tenant_id_fifth_part_uuid)::uuid AS tenant_id,
CONCAT(entity_id_first_part_uuid, '-', entity_id_second_part_uuid, '-1', entity_id_third_part_uuid, '-', entity_id_fourth_part_uuid, '-', entity_id_fifth_part_uuid)::uuid AS entity_id,
tenant_ts_kv_records.key AS key,
tenant_ts_kv_records.ts AS ts,
tenant_ts_kv_records.bool_v AS bool_v,
tenant_ts_kv_records.str_v AS str_v,
tenant_ts_kv_records.long_v AS long_v,
tenant_ts_kv_records.dbl_v AS dbl_v
FROM (SELECT SUBSTRING(tenant_id, 8, 8) AS tenant_id_first_part_uuid,
SUBSTRING(tenant_id, 4, 4) AS tenant_id_second_part_uuid,
SUBSTRING(tenant_id, 1, 3) AS tenant_id_third_part_uuid,
SUBSTRING(tenant_id, 16, 4) AS tenant_id_fourth_part_uuid,
SUBSTRING(tenant_id, 20) AS tenant_id_fifth_part_uuid,
SUBSTRING(entity_id, 8, 8) AS entity_id_first_part_uuid,
SUBSTRING(entity_id, 4, 4) AS entity_id_second_part_uuid,
SUBSTRING(entity_id, 1, 3) AS entity_id_third_part_uuid,
SUBSTRING(entity_id, 16, 4) AS entity_id_fourth_part_uuid,
SUBSTRING(entity_id, 20) AS entity_id_fifth_part_uuid,
key_id AS key,
ts,
bool_v,
@ -157,7 +157,7 @@ DECLARE
long_v,
dbl_v
FROM tenant_ts_kv_old
INNER JOIN ts_kv_dictionary ON (tenant_ts_kv_old.key = ts_kv_dictionary.key)) AS substrings;
INNER JOIN ts_kv_dictionary ON (tenant_ts_kv_old.key = ts_kv_dictionary.key)) AS tenant_ts_kv_records;
BEGIN
OPEN insert_cursor;
LOOP
@ -188,10 +188,10 @@ DECLARE
latest_record RECORD;
insert_record RECORD;
insert_cursor CURSOR FOR SELECT
latest.key AS key,
latest.entity_id AS entity_id,
latest.ts AS ts
FROM (SELECT DISTINCT key AS key, entity_id AS entity_id, MAX(ts) AS ts FROM tenant_ts_kv GROUP BY key, entity_id) AS latest;
latest_records.key AS key,
latest_records.entity_id AS entity_id,
latest_records.ts AS ts
FROM (SELECT DISTINCT key AS key, entity_id AS entity_id, MAX(ts) AS ts FROM tenant_ts_kv GROUP BY key, entity_id) AS latest_records;
BEGIN
OPEN insert_cursor;
LOOP

124
application/src/main/java/org/thingsboard/server/service/install/AbstractSqlTsDatabaseUpgradeService.java

@ -0,0 +1,124 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.CallableStatement;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Types;
@Slf4j
public abstract class AbstractSqlTsDatabaseUpgradeService {
protected static final String CALL_REGEX = "call ";
protected static final String CHECK_VERSION = "check_version()";
protected static final String DROP_TABLE = "DROP TABLE ";
protected static final String DROP_FUNCTION_IF_EXISTS = "DROP FUNCTION IF EXISTS ";
private static final String CALL_CHECK_VERSION = CALL_REGEX + CHECK_VERSION;
private static final String FUNCTION = "function: {}";
private static final String DROP_STATEMENT = "drop statement: {}";
private static final String QUERY = "query: {}";
private static final String SUCCESSFULLY_EXECUTED = "Successfully executed ";
private static final String FAILED_TO_EXECUTE = "Failed to execute ";
private static final String FAILED_DUE_TO = " due to: {}";
protected static final String SUCCESSFULLY_EXECUTED_FUNCTION = SUCCESSFULLY_EXECUTED + FUNCTION;
protected static final String FAILED_TO_EXECUTE_FUNCTION_DUE_TO = FAILED_TO_EXECUTE + FUNCTION + FAILED_DUE_TO;
protected static final String SUCCESSFULLY_EXECUTED_DROP_STATEMENT = SUCCESSFULLY_EXECUTED + DROP_STATEMENT;
protected static final String FAILED_TO_EXECUTE_DROP_STATEMENT = FAILED_TO_EXECUTE + DROP_STATEMENT + FAILED_DUE_TO;
protected static final String SUCCESSFULLY_EXECUTED_QUERY = SUCCESSFULLY_EXECUTED + QUERY;
protected static final String FAILED_TO_EXECUTE_QUERY = FAILED_TO_EXECUTE + QUERY + FAILED_DUE_TO;
@Value("${spring.datasource.url}")
protected String dbUrl;
@Value("${spring.datasource.username}")
protected String dbUserName;
@Value("${spring.datasource.password}")
protected String dbPassword;
@Autowired
protected InstallScripts installScripts;
protected abstract void loadSql(Connection conn);
protected void loadFunctions(Path sqlFile, Connection conn) throws Exception {
String sql = new String(Files.readAllBytes(sqlFile), StandardCharsets.UTF_8);
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
}
protected boolean checkVersion(Connection conn) {
log.info("Check the current PostgreSQL version...");
boolean versionValid = false;
try {
CallableStatement callableStatement = conn.prepareCall("{? = " + CALL_CHECK_VERSION + " }");
callableStatement.registerOutParameter(1, Types.BOOLEAN);
callableStatement.execute();
versionValid = callableStatement.getBoolean(1);
callableStatement.close();
} catch (Exception e) {
log.info("Failed to check current PostgreSQL version due to: {}", e.getMessage());
}
return versionValid;
}
protected void executeFunction(Connection conn, String query) {
log.info("{} ... ", query);
try {
CallableStatement callableStatement = conn.prepareCall("{" + query + "}");
callableStatement.execute();
callableStatement.close();
log.info(SUCCESSFULLY_EXECUTED_FUNCTION, query.replace(CALL_REGEX, ""));
Thread.sleep(2000);
} catch (Exception e) {
log.info(FAILED_TO_EXECUTE_FUNCTION_DUE_TO, query, e.getMessage());
}
}
protected void executeDropStatement(Connection conn, String query) {
try {
conn.createStatement().execute(query); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
log.info(SUCCESSFULLY_EXECUTED_DROP_STATEMENT, query);
Thread.sleep(5000);
} catch (InterruptedException | SQLException e) {
log.info(FAILED_TO_EXECUTE_DROP_STATEMENT, query, e.getMessage());
}
}
protected void executeQuery(Connection conn, String query) {
try {
conn.createStatement().execute(query); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
log.info(SUCCESSFULLY_EXECUTED_QUERY, query);
Thread.sleep(5000);
} catch (InterruptedException | SQLException e) {
log.info(FAILED_TO_EXECUTE_QUERY, query, e.getMessage());
}
}
}

127
application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java

@ -16,54 +16,55 @@
package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.SqlTsDao;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.sql.CallableStatement;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.SQLException;
import java.sql.Types;
@Service
@Profile("install")
@Slf4j
@SqlTsDao
@PsqlDao
public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
private static final String CALL_REGEX = "call ";
private static final String LOAD_FUNCTIONS_SQL = "schema_update_psql_ts.sql";
private static final String CHECK_VERSION = CALL_REGEX + "check_version()";
private static final String CREATE_PARTITION_TS_KV_TABLE = CALL_REGEX + "create_partition_ts_kv_table()";
private static final String CREATE_NEW_TS_KV_LATEST_TABLE = CALL_REGEX + "create_new_ts_kv_latest_table()";
private static final String CREATE_PARTITIONS = CALL_REGEX + "create_partitions()";
private static final String CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + "create_ts_kv_dictionary_table()";
private static final String INSERT_INTO_DICTIONARY = CALL_REGEX + "insert_into_dictionary()";
private static final String INSERT_INTO_TS_KV = CALL_REGEX + "insert_into_ts_kv()";
private static final String INSERT_INTO_TS_KV_LATEST = CALL_REGEX + "insert_into_ts_kv_latest()";
private static final String DROP_TABLE_TS_KV_OLD = "DROP TABLE ts_kv_old;";
private static final String DROP_TABLE_TS_KV_LATEST_OLD = "DROP TABLE ts_kv_latest_old;";
@Value("${spring.datasource.url}")
private String dbUrl;
private static final String TS_KV_OLD = "ts_kv_old;";
private static final String TS_KV_LATEST_OLD = "ts_kv_latest_old;";
@Value("${spring.datasource.username}")
private String dbUserName;
private static final String CREATE_PARTITION_TS_KV_TABLE = "create_partition_ts_kv_table()";
private static final String CREATE_NEW_TS_KV_LATEST_TABLE = "create_new_ts_kv_latest_table()";
private static final String CREATE_PARTITIONS = "create_partitions()";
private static final String CREATE_TS_KV_DICTIONARY_TABLE = "create_ts_kv_dictionary_table()";
private static final String INSERT_INTO_DICTIONARY = "insert_into_dictionary()";
private static final String INSERT_INTO_TS_KV = "insert_into_ts_kv()";
private static final String INSERT_INTO_TS_KV_LATEST = "insert_into_ts_kv_latest()";
@Value("${spring.datasource.password}")
private String dbPassword;
private static final String CALL_CREATE_PARTITION_TS_KV_TABLE = CALL_REGEX + CREATE_PARTITION_TS_KV_TABLE;
private static final String CALL_CREATE_NEW_TS_KV_LATEST_TABLE = CALL_REGEX + CREATE_NEW_TS_KV_LATEST_TABLE;
private static final String CALL_CREATE_PARTITIONS = CALL_REGEX + CREATE_PARTITIONS;
private static final String CALL_CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + CREATE_TS_KV_DICTIONARY_TABLE;
private static final String CALL_INSERT_INTO_DICTIONARY = CALL_REGEX + INSERT_INTO_DICTIONARY;
private static final String CALL_INSERT_INTO_TS_KV = CALL_REGEX + INSERT_INTO_TS_KV;
private static final String CALL_INSERT_INTO_TS_KV_LATEST = CALL_REGEX + INSERT_INTO_TS_KV_LATEST;
@Autowired
private InstallScripts installScripts;
private static final String DROP_TABLE_TS_KV_OLD = DROP_TABLE + TS_KV_OLD;
private static final String DROP_TABLE_TS_KV_LATEST_OLD = DROP_TABLE + TS_KV_LATEST_OLD;
private static final String DROP_FUNCTION_CHECK_VERSION = DROP_FUNCTION_IF_EXISTS + CHECK_VERSION;
private static final String DROP_FUNCTION_CREATE_PARTITION_TS_KV_TABLE = DROP_FUNCTION_IF_EXISTS + CREATE_PARTITION_TS_KV_TABLE;
private static final String DROP_FUNCTION_CREATE_NEW_TS_KV_LATEST_TABLE = DROP_FUNCTION_IF_EXISTS + CREATE_NEW_TS_KV_LATEST_TABLE;
private static final String DROP_FUNCTION_CREATE_PARTITIONS = DROP_FUNCTION_IF_EXISTS + CREATE_PARTITIONS;
private static final String DROP_FUNCTION_CREATE_TS_KV_DICTIONARY_TABLE = DROP_FUNCTION_IF_EXISTS + CREATE_TS_KV_DICTIONARY_TABLE;
private static final String DROP_FUNCTION_INSERT_INTO_DICTIONARY = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_DICTIONARY;
private static final String DROP_FUNCTION_INSERT_INTO_TS_KV = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_TS_KV;
private static final String DROP_FUNCTION_INSERT_INTO_TS_KV_LATEST = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_TS_KV_LATEST;
@Override
public void upgradeDatabase(String fromVersion) throws Exception {
@ -80,15 +81,26 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
} else {
log.info("PostgreSQL version is valid!");
log.info("Updating schema ...");
executeFunction(conn, CREATE_PARTITION_TS_KV_TABLE);
executeFunction(conn, CREATE_PARTITIONS);
executeFunction(conn, CREATE_TS_KV_DICTIONARY_TABLE);
executeFunction(conn, INSERT_INTO_DICTIONARY);
executeFunction(conn, INSERT_INTO_TS_KV);
executeFunction(conn, CREATE_NEW_TS_KV_LATEST_TABLE);
executeFunction(conn, INSERT_INTO_TS_KV_LATEST);
dropOldTable(conn, DROP_TABLE_TS_KV_OLD);
dropOldTable(conn, DROP_TABLE_TS_KV_LATEST_OLD);
executeFunction(conn, CALL_CREATE_PARTITION_TS_KV_TABLE);
executeFunction(conn, CALL_CREATE_PARTITIONS);
executeFunction(conn, CALL_CREATE_TS_KV_DICTIONARY_TABLE);
executeFunction(conn, CALL_INSERT_INTO_DICTIONARY);
executeFunction(conn, CALL_INSERT_INTO_TS_KV);
executeFunction(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE);
executeFunction(conn, CALL_INSERT_INTO_TS_KV_LATEST);
executeDropStatement(conn, DROP_TABLE_TS_KV_OLD);
executeDropStatement(conn, DROP_TABLE_TS_KV_LATEST_OLD);
executeDropStatement(conn, DROP_FUNCTION_CHECK_VERSION);
executeDropStatement(conn, DROP_FUNCTION_CREATE_PARTITION_TS_KV_TABLE);
executeDropStatement(conn, DROP_FUNCTION_CREATE_PARTITIONS);
executeDropStatement(conn, DROP_FUNCTION_CREATE_TS_KV_DICTIONARY_TABLE);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_DICTIONARY);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_TS_KV);
executeDropStatement(conn, DROP_FUNCTION_CREATE_NEW_TS_KV_LATEST_TABLE);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_TS_KV_LATEST);
log.info("schema timeseries updated!");
}
}
@ -98,7 +110,7 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
}
}
private void loadSql(Connection conn) {
protected void loadSql(Connection conn) {
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.4.3", LOAD_FUNCTIONS_SQL);
try {
loadFunctions(schemaUpdateFile, conn);
@ -107,45 +119,4 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
log.info("Failed to load PostgreSQL upgrade functions due to: {}", e.getMessage());
}
}
private void loadFunctions(Path sqlFile, Connection conn) throws Exception {
String sql = new String(Files.readAllBytes(sqlFile), StandardCharsets.UTF_8);
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
}
private boolean checkVersion(Connection conn) {
log.info("Check the current PostgreSQL version...");
boolean versionValid = false;
try {
CallableStatement callableStatement = conn.prepareCall("{? = " + CHECK_VERSION + " }");
callableStatement.registerOutParameter(1, Types.BOOLEAN);
callableStatement.execute();
versionValid = callableStatement.getBoolean(1);
callableStatement.close();
} catch (Exception e) {
log.info("Failed to check current PostgreSQL version due to: {}", e.getMessage());
}
return versionValid;
}
private void executeFunction(Connection conn, String query) {
log.info("{} ... ", query);
try {
CallableStatement callableStatement = conn.prepareCall("{" + query + "}");
callableStatement.execute();
callableStatement.close();
log.info("Successfully executed: {}", query.replace(CALL_REGEX, ""));
} catch (Exception e) {
log.info("Failed to execute {} due to: {}", query, e.getMessage());
}
}
private void dropOldTable(Connection conn, String query) {
try {
conn.createStatement().execute(query); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
Thread.sleep(5000);
} catch (InterruptedException | SQLException e) {
log.info("Failed to drop table {} due to: {}", query.replace("DROP TABLE ", ""), e.getMessage());
}
}
}

6
application/src/main/java/org/thingsboard/server/service/install/SqlAbstractDatabaseSchemaService.java

@ -32,13 +32,13 @@ public abstract class SqlAbstractDatabaseSchemaService implements DatabaseSchema
private static final String SQL_DIR = "sql";
@Value("${spring.datasource.url}")
private String dbUrl;
protected String dbUrl;
@Value("${spring.datasource.username}")
private String dbUserName;
protected String dbUserName;
@Value("${spring.datasource.password}")
private String dbPassword;
protected String dbPassword;
@Autowired
private InstallScripts installScripts;

31
application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java

@ -1,31 +0,0 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.install;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.SqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
@Service
@TimescaleDBTsDao
@Profile("install")
public class SqlTimescaleDatabaseSchemaService extends SqlAbstractDatabaseSchemaService
implements TsDatabaseSchemaService {
public SqlTimescaleDatabaseSchemaService() {
super("schema-timescale.sql", "schema-timescale-idx.sql");
}
}

147
application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseUpgradeService.java

@ -1,147 +0,0 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.sql.CallableStatement;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.SQLException;
import java.sql.Types;
@Service
@Profile("install")
@Slf4j
@TimescaleDBTsDao
@PsqlDao
public class SqlTimescaleDatabaseUpgradeService implements DatabaseTsUpgradeService {
private static final String CALL_REGEX = "call ";
private static final String LOAD_FUNCTIONS_SQL = "schema_update_timescale_ts.sql";
private static final String CHECK_VERSION = CALL_REGEX + "check_version()";
private static final String CREATE_TS_KV_LATEST_TABLE = CALL_REGEX + "create_ts_kv_latest_table()";
private static final String CREATE_TENANT_TS_KV_TABLE_COPY = CALL_REGEX + "create_tenant_ts_kv_table_copy()";
private static final String CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + "create_ts_kv_dictionary_table()";
private static final String INSERT_INTO_DICTIONARY = CALL_REGEX + "insert_into_dictionary()";
private static final String INSERT_INTO_TS_KV = CALL_REGEX + "insert_into_tenant_ts_kv()";
private static final String INSERT_INTO_TS_KV_LATEST = CALL_REGEX + "insert_into_ts_kv_latest()";
private static final String DROP_OLD_TS_KV_TABLE = "DROP TABLE tenant_ts_kv_old;";
@Value("${spring.datasource.url}")
private String dbUrl;
@Value("${spring.datasource.username}")
private String dbUserName;
@Value("${spring.datasource.password}")
private String dbPassword;
@Autowired
private InstallScripts installScripts;
@Override
public void upgradeDatabase(String fromVersion) throws Exception {
switch (fromVersion) {
case "2.4.3":
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating timescale schema ...");
log.info("Load upgrade functions ...");
loadSql(conn);
boolean versionValid = checkVersion(conn);
if (!versionValid) {
log.info("PostgreSQL version should be at least more than 9.6!");
log.info("Please upgrade your PostgreSQL and restart the script!");
} else {
log.info("PostgreSQL version is valid!");
log.info("Updating schema ...");
executeFunction(conn, CREATE_TS_KV_LATEST_TABLE);
executeFunction(conn, CREATE_TENANT_TS_KV_TABLE_COPY);
executeFunction(conn, CREATE_TS_KV_DICTIONARY_TABLE);
executeFunction(conn, INSERT_INTO_DICTIONARY);
executeFunction(conn, INSERT_INTO_TS_KV);
executeFunction(conn, INSERT_INTO_TS_KV_LATEST);
executeQuery(conn, DROP_OLD_TS_KV_TABLE);
log.info("schema timeseries updated!");
}
}
break;
default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
}
}
private void loadSql(Connection conn) {
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.4.3", LOAD_FUNCTIONS_SQL);
try {
loadFunctions(schemaUpdateFile, conn);
log.info("Upgrade functions successfully loaded!");
} catch (Exception e) {
log.info("Failed to load Timescale upgrade functions due to: {}", e.getMessage());
}
}
private void loadFunctions(Path sqlFile, Connection conn) throws Exception {
String sql = new String(Files.readAllBytes(sqlFile), StandardCharsets.UTF_8);
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
}
private boolean checkVersion(Connection conn) {
log.info("Check the current PostgreSQL version...");
boolean versionValid = false;
try {
CallableStatement callableStatement = conn.prepareCall("{? = " + CHECK_VERSION + " }");
callableStatement.registerOutParameter(1, Types.BOOLEAN);
callableStatement.execute();
versionValid = callableStatement.getBoolean(1);
callableStatement.close();
} catch (Exception e) {
log.info("Failed to check current PostgreSQL version due to: {}", e.getMessage());
}
return versionValid;
}
private void executeFunction(Connection conn, String query) {
log.info("{} ... ", query);
try {
CallableStatement callableStatement = conn.prepareCall("{" + query + "}");
callableStatement.execute();
callableStatement.close();
log.info("Successfully executed: {}", query.replace(CALL_REGEX, ""));
} catch (Exception e) {
log.info("Failed to execute {} due to: {}", query, e.getMessage());
}
}
private void executeQuery(Connection conn, String query) {
try {
conn.createStatement().execute(query); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
Thread.sleep(5000);
} catch (InterruptedException | SQLException e) {
log.info("Failed to drop table {} due to: {}", query.replace("DROP TABLE ", ""), e.getMessage());
}
}
}

68
application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseSchemaService.java

@ -0,0 +1,68 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.SQLException;
@Service
@TimescaleDBTsDao
@PsqlDao
@Profile("install")
@Slf4j
public class TimescaleTsDatabaseSchemaService extends SqlAbstractDatabaseSchemaService implements TsDatabaseSchemaService {
private static final String QUERY = "query: {}";
private static final String SUCCESSFULLY_EXECUTED = "Successfully executed ";
private static final String FAILED_TO_EXECUTE = "Failed to execute ";
private static final String FAILED_DUE_TO = " due to: {}";
private static final String SUCCESSFULLY_EXECUTED_QUERY = SUCCESSFULLY_EXECUTED + QUERY;
private static final String FAILED_TO_EXECUTE_QUERY = FAILED_TO_EXECUTE + QUERY + FAILED_DUE_TO;
@Value("${sql.timescale.chunk_time_interval:86400000}")
private long chunkTimeInterval;
public TimescaleTsDatabaseSchemaService() {
super("schema-timescale.sql", "schema-timescale-idx.sql");
}
@Override
public void createDatabaseSchema() throws Exception {
super.createDatabaseSchema();
executeQuery("SELECT create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => " + chunkTimeInterval + ", if_not_exists => true);");
}
private void executeQuery(String query) {
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
conn.createStatement().execute(query); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
log.info(SUCCESSFULLY_EXECUTED_QUERY, query);
Thread.sleep(5000);
} catch (InterruptedException | SQLException e) {
log.info(FAILED_TO_EXECUTE_QUERY, query, e.getMessage());
}
}
}

125
application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java

@ -0,0 +1,125 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.sql.Connection;
import java.sql.DriverManager;
@Service
@Profile("install")
@Slf4j
@TimescaleDBTsDao
@PsqlDao
public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
@Value("${sql.timescale.chunk_time_interval:86400000}")
private long chunkTimeInterval;
private static final String LOAD_FUNCTIONS_SQL = "schema_update_timescale_ts.sql";
private static final String TENANT_TS_KV_OLD_TABLE = "tenant_ts_kv_old;";
private static final String CREATE_TS_KV_LATEST_TABLE = "create_ts_kv_latest_table()";
private static final String CREATE_NEW_TENANT_TS_KV_TABLE = "create_new_tenant_ts_kv_table()";
private static final String CREATE_TS_KV_DICTIONARY_TABLE = "create_ts_kv_dictionary_table()";
private static final String INSERT_INTO_DICTIONARY = "insert_into_dictionary()";
private static final String INSERT_INTO_TENANT_TS_KV = "insert_into_tenant_ts_kv()";
private static final String INSERT_INTO_TS_KV_LATEST = "insert_into_ts_kv_latest()";
private static final String CALL_CREATE_TS_KV_LATEST_TABLE = CALL_REGEX + CREATE_TS_KV_LATEST_TABLE;
private static final String CALL_CREATE_NEW_TENANT_TS_KV_TABLE = CALL_REGEX + CREATE_NEW_TENANT_TS_KV_TABLE;
private static final String CALL_CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + CREATE_TS_KV_DICTIONARY_TABLE;
private static final String CALL_INSERT_INTO_DICTIONARY = CALL_REGEX + INSERT_INTO_DICTIONARY;
private static final String CALL_INSERT_INTO_TS_KV = CALL_REGEX + INSERT_INTO_TENANT_TS_KV;
private static final String CALL_INSERT_INTO_TS_KV_LATEST = CALL_REGEX + INSERT_INTO_TS_KV_LATEST;
private static final String DROP_OLD_TENANT_TS_KV_TABLE = DROP_TABLE + TENANT_TS_KV_OLD_TABLE;
private static final String DROP_FUNCTION_CREATE_TS_KV_LATEST_TABLE = DROP_FUNCTION_IF_EXISTS + CREATE_TS_KV_LATEST_TABLE;
private static final String DROP_FUNCTION_CREATE_TENANT_TS_KV_TABLE_COPY = DROP_FUNCTION_IF_EXISTS + CREATE_NEW_TENANT_TS_KV_TABLE;
private static final String DROP_FUNCTION_CREATE_TS_KV_DICTIONARY_TABLE = DROP_FUNCTION_IF_EXISTS + CREATE_TS_KV_DICTIONARY_TABLE;
private static final String DROP_FUNCTION_INSERT_INTO_DICTIONARY = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_DICTIONARY;
private static final String DROP_FUNCTION_INSERT_INTO_TENANT_TS_KV = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_TENANT_TS_KV;
private static final String DROP_FUNCTION_INSERT_INTO_TS_KV_LATEST = DROP_FUNCTION_IF_EXISTS + INSERT_INTO_TS_KV_LATEST;
@Autowired
private InstallScripts installScripts;
@Override
public void upgradeDatabase(String fromVersion) throws Exception {
switch (fromVersion) {
case "2.4.3":
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating timescale schema ...");
log.info("Load upgrade functions ...");
loadSql(conn);
boolean versionValid = checkVersion(conn);
if (!versionValid) {
log.info("PostgreSQL version should be at least more than 9.6!");
log.info("Please upgrade your PostgreSQL and restart the script!");
} else {
log.info("PostgreSQL version is valid!");
log.info("Updating schema ...");
executeFunction(conn, CALL_CREATE_TS_KV_LATEST_TABLE);
executeFunction(conn, CALL_CREATE_NEW_TENANT_TS_KV_TABLE);
executeQuery(conn, "SELECT create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => " + chunkTimeInterval + ", if_not_exists => true);");
executeFunction(conn, CALL_CREATE_TS_KV_DICTIONARY_TABLE);
executeFunction(conn, CALL_INSERT_INTO_DICTIONARY);
executeFunction(conn, CALL_INSERT_INTO_TS_KV);
executeFunction(conn, CALL_INSERT_INTO_TS_KV_LATEST);
//executeQuery(conn, "SELECT set_chunk_time_interval('tenant_ts_kv', " + chunkTimeInterval +");");
executeDropStatement(conn, DROP_OLD_TENANT_TS_KV_TABLE);
executeDropStatement(conn, DROP_FUNCTION_CREATE_TS_KV_LATEST_TABLE);
executeDropStatement(conn, DROP_FUNCTION_CREATE_TENANT_TS_KV_TABLE_COPY);
executeDropStatement(conn, DROP_FUNCTION_CREATE_TS_KV_DICTIONARY_TABLE);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_DICTIONARY);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_TENANT_TS_KV);
executeDropStatement(conn, DROP_FUNCTION_INSERT_INTO_TS_KV_LATEST);
log.info("schema timeseries updated!");
}
}
break;
default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
}
}
protected void loadSql(Connection conn) {
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.4.3", LOAD_FUNCTIONS_SQL);
try {
loadFunctions(schemaUpdateFile, conn);
log.info("Upgrade functions successfully loaded!");
} catch (Exception e) {
log.info("Failed to load Timescale upgrade functions due to: {}", e.getMessage());
}
}
}

8
application/src/main/resources/thingsboard.yml

@ -210,8 +210,12 @@ sql:
stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}"
# Specify whether to remove null characters from strValue of attributes and timeseries before insert
remove_null_chars: "${SQL_REMOVE_NULL_CHARS:true}"
# Specify partitioning size for timestamp key-value storage. Example: DAYS, MONTHS, YEARS, INDEFINITE
ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}"
postgres:
# Specify partitioning size for timestamp key-value storage. Example: DAYS, MONTHS, YEARS, INDEFINITE.
ts_key_value_partitioning: "${SQL_POSTGRES_TS_KV_PARTITIONING:MONTHS}"
timescale:
# Specify Interval size for new data chunks storage.
chunk_time_interval: "${SQL_TIMESCALE_CHUNK_TIME_INTERVAL:604800000}"
# Actor system parameters
actors:

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractPsqlHsqlTimeseriesDao.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java

@ -35,7 +35,7 @@ import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;
@Slf4j
public abstract class AbstractPsqlHsqlTimeseriesDao<T extends AbstractTsKvEntity> extends AbstractSqlTimeseriesDao {
public abstract class AbstractChunkedAggregationTimeseriesDao<T extends AbstractTsKvEntity> extends AbstractSqlTimeseriesDao {
@Autowired
protected InsertTsRepository<T> insertRepository;
@ -65,7 +65,7 @@ public abstract class AbstractPsqlHsqlTimeseriesDao<T extends AbstractTsKvEntity
protected abstract ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation);
protected void switchAgregation(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List<CompletableFuture<T>> entitiesFutures) {
protected void switchAggregation(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List<CompletableFuture<T>> entitiesFutures) {
switch (aggregation) {
case AVG:
findAvg(tenantId, entityId, key, startTs, endTs, entitiesFutures);

5
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java

@ -34,7 +34,6 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestCompositeKey;
@ -76,7 +75,7 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
private SearchTsKvLatestRepository searchTsKvLatestRepository;
@Autowired
private InsertLatestRepository insertLatestRepository;
private InsertLatestTsRepository insertLatestTsRepository;
@Autowired
private TsKvDictionaryRepository dictionaryRepository;
@ -113,7 +112,7 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
.statsPrintIntervalMs(tsLatestStatsPrintIntervalMs)
.build();
tsLatestQueue = new TbSqlBlockingQueue<>(tsLatestParams);
tsLatestQueue.init(logExecutor, v -> insertLatestRepository.saveOrUpdate(v));
tsLatestQueue.init(logExecutor, v -> insertLatestTsRepository.saveOrUpdate(v));
}
@PreDestroy

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/InsertLatestRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/InsertLatestTsRepository.java

@ -19,7 +19,7 @@ import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import java.util.List;
public interface InsertLatestRepository {
public interface InsertLatestTsRepository {
void saveOrUpdate(List<TsKvLatestEntity> entities);

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlTimeseriesInsertRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlInsertTsRepository.java

@ -34,7 +34,7 @@ import java.util.List;
@HsqlDao
@Repository
@Transactional
public class HsqlTimeseriesInsertRepository extends AbstractInsertRepository implements InsertTsRepository<TsKvEntity> {
public class HsqlInsertTsRepository extends AbstractInsertRepository implements InsertTsRepository<TsKvEntity> {
private static final String INSERT_OR_UPDATE =
"MERGE INTO ts_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?) " +

6
dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java

@ -30,7 +30,7 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sqlts.hsql.TsKvEntity;
import org.thingsboard.server.dao.sqlts.AbstractPsqlHsqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.AbstractChunkedAggregationTimeseriesDao;
import org.thingsboard.server.dao.sqlts.EntityContainer;
import org.thingsboard.server.dao.timeseries.TimeseriesDao;
import org.thingsboard.server.dao.util.HsqlDao;
@ -46,7 +46,7 @@ import java.util.concurrent.CompletableFuture;
@Slf4j
@SqlTsDao
@HsqlDao
public class JpaHsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
public class JpaHsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDao<TsKvEntity> implements TimeseriesDao {
@Autowired
private TsKvHsqlRepository tsKvRepository;
@ -150,7 +150,7 @@ public class JpaHsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEnti
@Override
protected ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) {
List<CompletableFuture<TsKvEntity>> entitiesFutures = new ArrayList<>();
switchAgregation(tenantId, entityId, key, startTs, endTs, aggregation, entitiesFutures);
switchAggregation(tenantId, entityId, key, startTs, endTs, aggregation, entitiesFutures);
return Futures.transform(setFutures(entitiesFutures), entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertTsRepository.java

@ -20,7 +20,7 @@ import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.sqlts.AbstractInsertRepository;
import org.thingsboard.server.dao.sqlts.InsertLatestRepository;
import org.thingsboard.server.dao.sqlts.InsertLatestTsRepository;
import org.thingsboard.server.dao.util.HsqlDao;
import org.thingsboard.server.dao.util.SqlTsDao;
@ -33,7 +33,7 @@ import java.util.List;
@HsqlDao
@Repository
@Transactional
public class HsqlLatestInsertRepository extends AbstractInsertRepository implements InsertLatestRepository {
public class HsqlLatestInsertTsRepository extends AbstractInsertRepository implements InsertLatestTsRepository {
private static final String INSERT_OR_UPDATE =
"MERGE INTO ts_kv_latest USING(VALUES ?, ?, ?, ?, ?, ?, ?) " +

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertTsRepository.java

@ -22,7 +22,7 @@ import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionCallbackWithoutResult;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.sqlts.AbstractInsertRepository;
import org.thingsboard.server.dao.sqlts.InsertLatestRepository;
import org.thingsboard.server.dao.sqlts.InsertLatestTsRepository;
import org.thingsboard.server.dao.util.PsqlTsAnyDao;
import java.sql.PreparedStatement;
@ -35,7 +35,7 @@ import java.util.List;
@PsqlTsAnyDao
@Repository
@Transactional
public class PsqlLatestInsertRepository extends AbstractInsertRepository implements InsertLatestRepository {
public class PsqlLatestInsertTsRepository extends AbstractInsertRepository implements InsertLatestTsRepository {
private static final String BATCH_UPDATE =
"UPDATE ts_kv_latest SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ? WHERE entity_id = ? and key = ?";

8
dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java

@ -31,7 +31,7 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sqlts.psql.TsKvEntity;
import org.thingsboard.server.dao.sqlts.AbstractPsqlHsqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.AbstractChunkedAggregationTimeseriesDao;
import org.thingsboard.server.dao.sqlts.EntityContainer;
import org.thingsboard.server.dao.timeseries.PsqlPartition;
import org.thingsboard.server.dao.timeseries.SqlTsPartitionDate;
@ -59,7 +59,7 @@ import static org.thingsboard.server.dao.timeseries.SqlTsPartitionDate.EPOCH_STA
@Slf4j
@SqlTsDao
@PsqlDao
public class JpaPsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
public class JpaPsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDao<TsKvEntity> implements TimeseriesDao {
private final Map<Long, PsqlPartition> partitions = new ConcurrentHashMap<>();
private static final ReentrantLock partitionCreationLock = new ReentrantLock();
@ -73,7 +73,7 @@ public class JpaPsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEnti
private SqlTsPartitionDate tsFormat;
private PsqlPartition indefinitePartition;
@Value("${sql.ts_key_value_partitioning}")
@Value("${sql.postgres.ts_key_value_partitioning:MONTHS}")
private String partitioning;
@Override
@ -188,7 +188,7 @@ public class JpaPsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEnti
@Override
protected ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) {
List<CompletableFuture<TsKvEntity>> entitiesFutures = new ArrayList<>();
switchAgregation(tenantId, entityId, key, startTs, endTs, aggregation, entitiesFutures);
switchAggregation(tenantId, entityId, key, startTs, endTs, aggregation, entitiesFutures);
return Futures.transform(setFutures(entitiesFutures), entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlTimeseriesInsertRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/PsqlInsertTsRepository.java

@ -37,7 +37,7 @@ import java.util.Map;
@PsqlDao
@Repository
@Transactional
public class PsqlTimeseriesInsertRepository extends AbstractInsertRepository implements InsertTsRepository<TsKvEntity> {
public class PsqlInsertTsRepository extends AbstractInsertRepository implements InsertTsRepository<TsKvEntity> {
private static final String INSERT_INTO_TS_KV = "INSERT INTO ts_kv_";

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertRepository.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertTsRepository.java

@ -34,7 +34,7 @@ import java.util.List;
@PsqlDao
@Repository
@Transactional
public class TimescaleInsertRepository extends AbstractInsertRepository implements InsertTsRepository<TimescaleTsKvEntity> {
public class TimescaleInsertTsRepository extends AbstractInsertRepository implements InsertTsRepository<TimescaleTsKvEntity> {
private static final String INSERT_OR_UPDATE =
"INSERT INTO tenant_ts_kv (tenant_id, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES(?, ?, ?, ?, ?, ?, ?, ?) " +

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java

@ -117,7 +117,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
}
private ListenableFuture<List<Optional<TsKvEntry>>> findAllAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) {
CompletableFuture<List<TimescaleTsKvEntity>> listCompletableFuture = switchAgregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId(), tenantId.getId());
CompletableFuture<List<TimescaleTsKvEntity>> listCompletableFuture = switchAggregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId(), tenantId.getId());
SettableFuture<List<TimescaleTsKvEntity>> listenableFuture = SettableFuture.create();
listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> {
if (throwable != null) {
@ -213,7 +213,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
return service.submit(() -> null);
}
private CompletableFuture<List<TimescaleTsKvEntity>> switchAgregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId, UUID tenantId) {
private CompletableFuture<List<TimescaleTsKvEntity>> switchAggregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId, UUID tenantId) {
switch (aggregation) {
case AVG:
return findAvg(key, startTs, endTs, timeBucket, entityId, tenantId);

4
dao/src/main/resources/sql/schema-timescale.sql

@ -43,6 +43,4 @@ CREATE TABLE IF NOT EXISTS ts_kv_latest (
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key)
);
SELECT create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => 86400000, if_not_exists => true);
);
Loading…
Cancel
Save