Browse Source

bug fixes & improvements / sql-timeseries (#2382)

* fixed the partion date extracting

* fix imports

* ts-keys dictionary for latest, hsqldb

* removed AbstractSimpleSqlTimeseriesDao class & fix beanCreationException in ThingsboardInstallService

* timescale-db upgrade added

* added postgreSQL upgrade

* fix logging

* refactoring timeseries-dao implementation
pull/2412/head
ShvaykaD 7 years ago
committed by Andrew Shvayka
parent
commit
3955600a9c
  1. 86
      application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql
  2. 213
      application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql
  3. 2
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  4. 16
      application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java
  5. 147
      application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseUpgradeService.java
  6. 4
      application/src/main/resources/thingsboard.yml
  7. 22
      common/dao-api/src/main/java/org/thingsboard/server/dao/util/SqlTsAnyDao.java
  8. 0
      common/dao-api/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java
  9. 4
      dao/src/main/java/org/thingsboard/server/dao/HsqlTsDaoConfig.java
  10. 1
      dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java
  11. 2
      dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java
  12. 37
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java
  13. 2
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/TsKvDictionary.java
  14. 6
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvCompositeKey.java
  15. 29
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java
  16. 6
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestCompositeKey.java
  17. 84
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestEntity.java
  18. 27
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/psql/TsKvEntity.java
  19. 26
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java
  20. 96
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractPsqlHsqlTimeseriesDao.java
  21. 96
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  22. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/TsKvDictionaryRepository.java
  23. 32
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlTimeseriesInsertRepository.java
  24. 131
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java
  25. 76
      dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java
  26. 32
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertRepository.java
  27. 48
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertRepository.java
  28. 45
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java
  29. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java
  30. 122
      dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java
  31. 198
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  32. 9
      dao/src/main/resources/sql/schema-timescale.sql
  33. 22
      dao/src/main/resources/sql/schema-ts-hsql.sql
  34. 19
      dao/src/main/resources/sql/schema-ts-psql.sql

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

@ -14,7 +14,7 @@
-- limitations under the License.
--
-- load function check_version()
-- select check_version();
CREATE OR REPLACE FUNCTION check_version() RETURNS boolean AS $$
DECLARE
@ -38,9 +38,9 @@ BEGIN
END;
$$ LANGUAGE 'plpgsql';
-- load function create_partition_table()
-- select create_partition_ts_kv_table();
CREATE OR REPLACE FUNCTION create_partition_table() RETURNS VOID AS $$
CREATE OR REPLACE FUNCTION create_partition_ts_kv_table() RETURNS VOID AS $$
BEGIN
ALTER TABLE ts_kv
@ -59,8 +59,32 @@ BEGIN
END;
$$ LANGUAGE 'plpgsql';
-- select create_new_ts_kv_latest_table();
-- load function create_partitions()
CREATE OR REPLACE FUNCTION create_new_ts_kv_latest_table() RETURNS VOID AS $$
BEGIN
ALTER TABLE ts_kv_latest
RENAME TO ts_kv_latest_old;
ALTER TABLE ts_kv_latest_old
RENAME CONSTRAINT ts_kv_latest_pkey TO ts_kv_latest_pkey_old;
CREATE TABLE IF NOT EXISTS ts_kv_latest
(
LIKE ts_kv_latest_old
);
ALTER TABLE ts_kv_latest
DROP COLUMN entity_type;
ALTER TABLE ts_kv_latest
ALTER COLUMN entity_id TYPE uuid USING entity_id::uuid;
ALTER TABLE ts_kv_latest
ALTER COLUMN key TYPE integer USING key::integer;
ALTER TABLE ts_kv_latest
ADD CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key);
END;
$$ LANGUAGE 'plpgsql';
-- select create_partitions();
CREATE OR REPLACE FUNCTION create_partitions() RETURNS VOID AS
$$
@ -89,7 +113,7 @@ BEGIN
END;
$$ language 'plpgsql';
-- load function create_ts_kv_dictionary_table()
-- select create_ts_kv_dictionary_table();
CREATE OR REPLACE FUNCTION create_ts_kv_dictionary_table() RETURNS VOID AS $$
@ -103,7 +127,7 @@ BEGIN
END;
$$ LANGUAGE 'plpgsql';
-- load function insert_into_dictionary()
-- select insert_into_dictionary();
CREATE OR REPLACE FUNCTION insert_into_dictionary() RETURNS VOID AS
$$
@ -128,7 +152,7 @@ BEGIN
END;
$$ language 'plpgsql';
-- load function insert_into_ts_kv()
-- select insert_into_ts_kv();
CREATE OR REPLACE FUNCTION insert_into_ts_kv() RETURNS void AS
$$
@ -176,4 +200,52 @@ BEGIN
END;
$$ LANGUAGE 'plpgsql';
-- select insert_into_ts_kv_latest();
CREATE OR REPLACE FUNCTION insert_into_ts_kv_latest() RETURNS void AS
$$
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,
key_id AS key,
ts,
bool_v,
str_v,
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;
BEGIN
OPEN insert_cursor;
LOOP
insert_counter := insert_counter + 1;
FETCH insert_cursor INTO insert_record;
IF NOT FOUND THEN
RAISE NOTICE '% records have been inserted into the ts_kv_latest!',insert_counter - 1;
EXIT;
END IF;
INSERT INTO ts_kv_latest(entity_id, key, ts, bool_v, str_v, long_v, dbl_v)
VALUES (insert_record.entity_id, insert_record.key, insert_record.ts, insert_record.bool_v, insert_record.str_v,
insert_record.long_v, insert_record.dbl_v);
IF MOD(insert_counter, insert_size) = 0 THEN
RAISE NOTICE '% records have been inserted into the ts_kv_latest!',insert_counter;
END IF;
END LOOP;
CLOSE insert_cursor;
END;
$$ LANGUAGE 'plpgsql';

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

@ -0,0 +1,213 @@
--
-- 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.
--
-- select check_version();
CREATE OR REPLACE FUNCTION check_version() RETURNS boolean AS $$
DECLARE
current_version integer;
valid_version boolean;
BEGIN
RAISE NOTICE 'Check the current installed PostgreSQL version...';
SELECT current_setting('server_version_num') INTO current_version;
IF current_version < 90600 THEN
valid_version := FALSE;
ELSE
valid_version := TRUE;
END IF;
IF valid_version = FALSE THEN
RAISE NOTICE 'Postgres version should be at least more than 9.6!';
ELSE
RAISE NOTICE 'PostgreSQL version is valid!';
RAISE NOTICE 'Schema update started...';
END IF;
RETURN valid_version;
END;
$$ LANGUAGE 'plpgsql';
-- select create_tenant_ts_kv_table_copy();
CREATE OR REPLACE FUNCTION create_tenant_ts_kv_table_copy() RETURNS VOID AS $$
BEGIN
ALTER TABLE tenant_ts_kv
RENAME TO tenant_ts_kv_old;
CREATE TABLE IF NOT EXISTS tenant_ts_kv
(
LIKE tenant_ts_kv_old
);
ALTER TABLE tenant_ts_kv
ALTER COLUMN tenant_id TYPE uuid USING tenant_id::uuid;
ALTER TABLE tenant_ts_kv
ALTER COLUMN entity_id TYPE uuid USING entity_id::uuid;
ALTER TABLE tenant_ts_kv
ALTER COLUMN key TYPE integer USING key::integer;
ALTER TABLE tenant_ts_kv
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);
CREATE INDEX IF NOT EXISTS idx_tenant_ts_kv ON tenant_ts_kv(tenant_id, entity_id, key, ts);
END;
$$ LANGUAGE 'plpgsql';
-- select create_ts_kv_latest_table();
CREATE OR REPLACE FUNCTION create_ts_kv_latest_table() RETURNS VOID AS $$
BEGIN
CREATE TABLE IF NOT EXISTS ts_kv_latest
(
entity_id uuid NOT NULL,
key int NOT NULL,
ts bigint NOT NULL,
bool_v boolean,
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key)
);
END;
$$ LANGUAGE 'plpgsql';
-- select create_ts_kv_dictionary_table();
CREATE OR REPLACE FUNCTION create_ts_kv_dictionary_table() RETURNS VOID AS $$
BEGIN
CREATE TABLE IF NOT EXISTS ts_kv_dictionary
(
key varchar(255) NOT NULL,
key_id serial UNIQUE,
CONSTRAINT ts_key_id_pkey PRIMARY KEY (key)
);
END;
$$ LANGUAGE 'plpgsql';
-- select insert_into_dictionary();
CREATE OR REPLACE FUNCTION insert_into_dictionary() RETURNS VOID AS
$$
DECLARE
insert_record RECORD;
key_cursor CURSOR FOR SELECT DISTINCT key
FROM tenant_ts_kv_old
ORDER BY key;
BEGIN
OPEN key_cursor;
LOOP
FETCH key_cursor INTO insert_record;
EXIT WHEN NOT FOUND;
IF NOT EXISTS(SELECT key FROM ts_kv_dictionary WHERE key = insert_record.key) THEN
INSERT INTO ts_kv_dictionary(key) VALUES (insert_record.key);
RAISE NOTICE 'Key: % has been inserted into the dictionary!',insert_record.key;
ELSE
RAISE NOTICE 'Key: % already exists in the dictionary!',insert_record.key;
END IF;
END LOOP;
CLOSE key_cursor;
END;
$$ language 'plpgsql';
-- select insert_into_tenant_ts_kv();
CREATE OR REPLACE FUNCTION insert_into_tenant_ts_kv() RETURNS void AS
$$
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,
key_id AS key,
ts,
bool_v,
str_v,
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;
BEGIN
OPEN insert_cursor;
LOOP
insert_counter := insert_counter + 1;
FETCH insert_cursor INTO insert_record;
IF NOT FOUND THEN
RAISE NOTICE '% records have been inserted into the new tenant_ts_kv table!',insert_counter - 1;
EXIT;
END IF;
INSERT INTO tenant_ts_kv(tenant_id, entity_id, key, ts, bool_v, str_v, long_v, dbl_v)
VALUES (insert_record.tenant_id, insert_record.entity_id, insert_record.key, insert_record.ts, insert_record.bool_v, insert_record.str_v,
insert_record.long_v, insert_record.dbl_v);
IF MOD(insert_counter, insert_size) = 0 THEN
RAISE NOTICE '% records have been inserted into the new tenant_ts_kv table!',insert_counter;
END IF;
END LOOP;
CLOSE insert_cursor;
END;
$$ LANGUAGE 'plpgsql';
-- select insert_into_ts_kv_latest();
CREATE OR REPLACE FUNCTION insert_into_ts_kv_latest() RETURNS void AS
$$
DECLARE
insert_size CONSTANT integer := 10000;
insert_counter integer DEFAULT 0;
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;
BEGIN
OPEN insert_cursor;
LOOP
insert_counter := insert_counter + 1;
FETCH insert_cursor INTO latest_record;
IF NOT FOUND THEN
RAISE NOTICE '% records have been inserted into the ts_kv_latest table!',insert_counter - 1;
EXIT;
END IF;
SELECT entity_id AS entity_id, key AS key, ts AS ts, bool_v AS bool_v, str_v AS str_v, long_v AS long_v, dbl_v AS dbl_v INTO insert_record FROM tenant_ts_kv WHERE entity_id = latest_record.entity_id AND key = latest_record.key AND ts = latest_record.ts;
INSERT INTO ts_kv_latest(entity_id, key, ts, bool_v, str_v, long_v, dbl_v)
VALUES (insert_record.entity_id, insert_record.key, insert_record.ts, insert_record.bool_v, insert_record.str_v, insert_record.long_v, insert_record.dbl_v);
IF MOD(insert_counter, insert_size) = 0 THEN
RAISE NOTICE '% records have been inserted into the ts_kv_latest table!',insert_counter;
END IF;
END LOOP;
CLOSE insert_cursor;
END;
$$ LANGUAGE 'plpgsql';

2
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -53,7 +53,7 @@ public class ThingsboardInstallService {
@Autowired
private DatabaseEntitiesUpgradeService databaseEntitiesUpgradeService;
@Autowired
@Autowired(required = false)
private DatabaseTsUpgradeService databaseTsUpgradeService;
@Autowired

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

@ -43,12 +43,15 @@ public class PsqlTsDatabaseUpgradeService 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_TABLE = CALL_REGEX + "create_partition_table()";
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 DROP_OLD_TABLE = "DROP TABLE ts_kv_old;";
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;
@ -70,7 +73,6 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
log.info("Updating timeseries schema ...");
log.info("Load upgrade functions ...");
loadSql(conn);
log.info("Upgrade functions successfully loaded!");
boolean versionValid = checkVersion(conn);
if (!versionValid) {
log.info("PostgreSQL version should be at least more than 10!");
@ -78,12 +80,15 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
} else {
log.info("PostgreSQL version is valid!");
log.info("Updating schema ...");
executeFunction(conn, CREATE_PARTITION_TABLE);
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);
dropOldTable(conn, DROP_OLD_TABLE);
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);
log.info("schema timeseries updated!");
}
}
@ -97,6 +102,7 @@ public class PsqlTsDatabaseUpgradeService implements DatabaseTsUpgradeService {
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 PostgreSQL upgrade functions due to: {}", e.getMessage());
}

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

@ -0,0 +1,147 @@
/**
* 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());
}
}
}

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

@ -208,10 +208,6 @@ sql:
batch_size: "${SQL_TS_LATEST_BATCH_SIZE:10000}"
batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:100}"
stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}"
ts_timescale:
batch_size: "${SQL_TS_TIMESCALE_BATCH_SIZE:10000}"
batch_max_delay: "${SQL_TS_TIMESCALE_BATCH_MAX_DELAY_MS:100}"
stats_print_interval_ms: "${SQL_TS_TIMESCALE_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

22
common/dao-api/src/main/java/org/thingsboard/server/dao/util/SqlTsAnyDao.java

@ -0,0 +1,22 @@
/**
* 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.dao.util;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
@ConditionalOnExpression("'${database.ts.type}'=='sql' || '${database.ts.type}'=='timescale'")
public @interface SqlTsAnyDao {
}

0
dao/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java → common/dao-api/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java

4
dao/src/main/java/org/thingsboard/server/dao/HsqlTsDaoConfig.java

@ -27,8 +27,8 @@ import org.thingsboard.server.dao.util.SqlTsDao;
@Configuration
@EnableAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sqlts.hsql", "org.thingsboard.server.dao.sqlts.latest"})
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.hsql", "org.thingsboard.server.dao.sqlts.latest"})
@EntityScan({"org.thingsboard.server.dao.model.sqlts.hsql", "org.thingsboard.server.dao.model.sqlts.latest"})
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.hsql", "org.thingsboard.server.dao.sqlts.latest", "org.thingsboard.server.dao.sqlts.dictionary"})
@EntityScan({"org.thingsboard.server.dao.model.sqlts.hsql", "org.thingsboard.server.dao.model.sqlts.latest", "org.thingsboard.server.dao.model.sqlts.dictionary"})
@EnableTransactionManagement
@SqlTsDao
@HsqlDao

1
dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java

@ -22,7 +22,6 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.thingsboard.server.dao.util.SqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
/**
* @author Valerii Sosliuk

2
dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java

@ -21,6 +21,7 @@ import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
@Configuration
@ -30,6 +31,7 @@ import org.thingsboard.server.dao.util.TimescaleDBTsDao;
@EntityScan({"org.thingsboard.server.dao.model.sqlts.timescale", "org.thingsboard.server.dao.model.sqlts.dictionary", "org.thingsboard.server.dao.model.sqlts.latest"})
@EnableTransactionManagement
@TimescaleDBTsDao
@PsqlDao
public class TimescaleDaoConfig {
}

37
dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java

@ -16,26 +16,42 @@
package org.thingsboard.server.dao.model.sql;
import lombok.Data;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.model.ToData;
import javax.persistence.Column;
import javax.persistence.Id;
import javax.persistence.MappedSuperclass;
import javax.persistence.Transient;
import java.util.UUID;
import static org.thingsboard.server.dao.model.ModelConstants.BOOLEAN_VALUE_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.DOUBLE_VALUE_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.LONG_VALUE_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.STRING_VALUE_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.TS_COLUMN;
@Data
@MappedSuperclass
public abstract class AbstractTsKvEntity {
public abstract class AbstractTsKvEntity implements ToData<TsKvEntry> {
protected static final String SUM = "SUM";
protected static final String AVG = "AVG";
protected static final String MIN = "MIN";
protected static final String MAX = "MAX";
@Id
@Column(name = ENTITY_ID_COLUMN, columnDefinition = "uuid")
protected UUID entityId;
@Id
@Column(name = TS_COLUMN)
protected Long ts;
@ -52,6 +68,9 @@ public abstract class AbstractTsKvEntity {
@Column(name = DOUBLE_VALUE_COLUMN)
protected Double doubleValue;
@Transient
protected String strKey;
public abstract boolean isNotEmpty();
protected static boolean isAllNull(Object... args) {
@ -62,4 +81,20 @@ public abstract class AbstractTsKvEntity {
}
return true;
}
@Override
public TsKvEntry toData() {
KvEntry kvEntry = null;
if (strValue != null) {
kvEntry = new StringDataEntry(strKey, strValue);
} else if (longValue != null) {
kvEntry = new LongDataEntry(strKey, longValue);
} else if (doubleValue != null) {
kvEntry = new DoubleDataEntry(strKey, doubleValue);
} else if (booleanValue != null) {
kvEntry = new BooleanDataEntry(strKey, booleanValue);
}
return new BasicTsKvEntry(ts, kvEntry);
}
}

2
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/TsKvDictionary.java

@ -38,7 +38,7 @@ public final class TsKvDictionary {
@Column(name = KEY_COLUMN)
private String key;
@Column(name = KEY_ID_COLUMN, unique = true, columnDefinition="serial")
@Column(name = KEY_ID_COLUMN, unique = true, columnDefinition="int")
@Generated(GenerationTime.INSERT)
private int keyId;

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

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.EntityType;
import javax.persistence.Transient;
import java.io.Serializable;
import java.util.UUID;
@Data
@AllArgsConstructor
@ -31,9 +32,8 @@ public class TsKvCompositeKey implements Serializable {
@Transient
private static final long serialVersionUID = -4089175869616037523L;
private EntityType entityType;
private String entityId;
private String key;
private UUID entityId;
private int key;
private long ts;
}

29
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/hsql/TsKvEntity.java

@ -34,6 +34,9 @@ import javax.persistence.Enumerated;
import javax.persistence.Id;
import javax.persistence.IdClass;
import javax.persistence.Table;
import javax.persistence.Transient;
import java.util.UUID;
import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_TYPE_COLUMN;
@ -45,18 +48,9 @@ import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN;
@IdClass(TsKvCompositeKey.class)
public final class TsKvEntity extends AbstractTsKvEntity implements ToData<TsKvEntry> {
@Id
@Enumerated(EnumType.STRING)
@Column(name = ENTITY_TYPE_COLUMN)
private EntityType entityType;
@Id
@Column(name = ENTITY_ID_COLUMN)
private String entityId;
@Id
@Column(name = KEY_COLUMN)
private String key;
private int key;
public TsKvEntity() {
}
@ -120,19 +114,4 @@ public final class TsKvEntity extends AbstractTsKvEntity implements ToData<TsKvE
public boolean isNotEmpty() {
return strValue != null || longValue != null || doubleValue != null || booleanValue != null;
}
@Override
public TsKvEntry toData() {
KvEntry kvEntry = null;
if (strValue != null) {
kvEntry = new StringDataEntry(key, strValue);
} else if (longValue != null) {
kvEntry = new LongDataEntry(key, longValue);
} else if (doubleValue != null) {
kvEntry = new DoubleDataEntry(key, doubleValue);
} else if (booleanValue != null) {
kvEntry = new BooleanDataEntry(key, booleanValue);
}
return new BasicTsKvEntry(ts, kvEntry);
}
}

6
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestCompositeKey.java

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.EntityType;
import javax.persistence.Transient;
import java.io.Serializable;
import java.util.UUID;
@Data
@NoArgsConstructor
@ -31,7 +32,6 @@ public class TsKvLatestCompositeKey implements Serializable{
@Transient
private static final long serialVersionUID = -4089175869616037523L;
private EntityType entityType;
private String entityId;
private String key;
private UUID entityId;
private int key;
}

84
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/latest/TsKvLatestEntity.java

@ -16,66 +16,78 @@
package org.thingsboard.server.dao.model.sqlts.latest;
import lombok.Data;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.model.ToData;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository;
import javax.persistence.Column;
import javax.persistence.ColumnResult;
import javax.persistence.ConstructorResult;
import javax.persistence.Entity;
import javax.persistence.EnumType;
import javax.persistence.Enumerated;
import javax.persistence.Id;
import javax.persistence.IdClass;
import javax.persistence.NamedNativeQueries;
import javax.persistence.NamedNativeQuery;
import javax.persistence.SqlResultSetMapping;
import javax.persistence.SqlResultSetMappings;
import javax.persistence.Table;
import java.util.UUID;
import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_TYPE_COLUMN;
import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN;
@Data
@Entity
@Table(name = "ts_kv_latest")
@IdClass(TsKvLatestCompositeKey.class)
public final class TsKvLatestEntity extends AbstractTsKvEntity implements ToData<TsKvEntry> {
@SqlResultSetMappings({
@SqlResultSetMapping(
name = "tsKvLatestFindMapping",
classes = {
@ConstructorResult(
targetClass = TsKvLatestEntity.class,
columns = {
@ColumnResult(name = "entityId", type = UUID.class),
@ColumnResult(name = "key", type = Integer.class),
@ColumnResult(name = "strKey", type = String.class),
@ColumnResult(name = "strValue", type = String.class),
@ColumnResult(name = "boolValue", type = Boolean.class),
@ColumnResult(name = "longValue", type = Long.class),
@ColumnResult(name = "doubleValue", type = Double.class),
@ColumnResult(name = "ts", type = Long.class),
@Id
@Enumerated(EnumType.STRING)
@Column(name = ENTITY_TYPE_COLUMN)
private EntityType entityType;
@Id
@Column(name = ENTITY_ID_COLUMN)
private String entityId;
}
),
})
})
@NamedNativeQueries({
@NamedNativeQuery(
name = SearchTsKvLatestRepository.FIND_ALL_BY_ENTITY_ID,
query = SearchTsKvLatestRepository.FIND_ALL_BY_ENTITY_ID_QUERY,
resultSetMapping = "tsKvLatestFindMapping",
resultClass = TsKvLatestEntity.class
)
})
public final class TsKvLatestEntity extends AbstractTsKvEntity {
@Id
@Column(name = KEY_COLUMN)
private String key;
private int key;
@Override
public boolean isNotEmpty() {
return strValue != null || longValue != null || doubleValue != null || booleanValue != null;
}
@Override
public TsKvEntry toData() {
KvEntry kvEntry = null;
if (strValue != null) {
kvEntry = new StringDataEntry(key, strValue);
} else if (longValue != null) {
kvEntry = new LongDataEntry(key, longValue);
} else if (doubleValue != null) {
kvEntry = new DoubleDataEntry(key, doubleValue);
} else if (booleanValue != null) {
kvEntry = new BooleanDataEntry(key, booleanValue);
}
return new BasicTsKvEntry(ts, kvEntry);
public TsKvLatestEntity() {
}
public TsKvLatestEntity(UUID entityId, Integer key, String strKey, String strValue, Boolean boolValue, Long longValue, Double doubleValue, Long ts) {
this.entityId = entityId;
this.key = key;
this.ts = ts;
this.longValue = longValue;
this.doubleValue = doubleValue;
this.strValue = strValue;
this.booleanValue = boolValue;
this.strKey = strKey;
}
}

27
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/psql/TsKvEntity.java

@ -41,18 +41,11 @@ import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN;
@Entity
@Table(name = "ts_kv")
@IdClass(TsKvCompositeKey.class)
public final class TsKvEntity extends AbstractTsKvEntity implements ToData<TsKvEntry> {
@Id
@Column(name = ENTITY_ID_COLUMN, columnDefinition = "uuid")
protected UUID entityId;
public final class TsKvEntity extends AbstractTsKvEntity {
@Id
@Column(name = KEY_COLUMN)
protected int key;
@Transient
protected String strKey;
private int key;
public TsKvEntity() {
}
@ -116,20 +109,4 @@ public final class TsKvEntity extends AbstractTsKvEntity implements ToData<TsKvE
public boolean isNotEmpty() {
return strValue != null || longValue != null || doubleValue != null || booleanValue != null;
}
@Override
public TsKvEntry toData() {
KvEntry kvEntry = null;
if (strValue != null) {
kvEntry = new StringDataEntry(strKey, strValue);
} else if (longValue != null) {
kvEntry = new LongDataEntry(strKey, longValue);
} else if (doubleValue != null) {
kvEntry = new DoubleDataEntry(strKey, doubleValue);
} else if (booleanValue != null) {
kvEntry = new BooleanDataEntry(strKey, booleanValue);
}
return new BasicTsKvEntry(ts, kvEntry);
}
}

26
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java

@ -129,17 +129,9 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity implements ToD
@Column(name = TENANT_ID_COLUMN, columnDefinition = "uuid")
private UUID tenantId;
@Id
@Column(name = ENTITY_ID_COLUMN, columnDefinition = "uuid")
protected UUID entityId;
@Id
@Column(name = KEY_COLUMN)
protected int key;
@Transient
protected String strKey;
private int key;
public TimescaleTsKvEntity() {
}
@ -204,20 +196,4 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity implements ToD
public boolean isNotEmpty() {
return ts != null && (strValue != null || longValue != null || doubleValue != null || booleanValue != null);
}
@Override
public TsKvEntry toData() {
KvEntry kvEntry = null;
if (strValue != null) {
kvEntry = new StringDataEntry(strKey, strValue);
} else if (longValue != null) {
kvEntry = new LongDataEntry(strKey, longValue);
} else if (doubleValue != null) {
kvEntry = new DoubleDataEntry(strKey, doubleValue);
} else if (booleanValue != null) {
kvEntry = new BooleanDataEntry(strKey, booleanValue);
}
return new BasicTsKvEntry(ts, kvEntry);
}
}

96
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSimpleSqlTimeseriesDao.java → dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractPsqlHsqlTimeseriesDao.java

@ -15,16 +15,13 @@
*/
package org.thingsboard.server.dao.sqlts;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueue;
@ -32,26 +29,16 @@ import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;
@Slf4j
public abstract class AbstractSimpleSqlTimeseriesDao<T extends AbstractTsKvEntity> extends AbstractSqlTimeseriesDao {
public abstract class AbstractPsqlHsqlTimeseriesDao<T extends AbstractTsKvEntity> extends AbstractSqlTimeseriesDao {
@Autowired
private InsertTsRepository<T> insertRepository;
@Value("${sql.ts.batch_size:1000}")
private int tsBatchSize;
@Value("${sql.ts.batch_max_delay:100}")
private long tsMaxDelay;
@Value("${sql.ts.stats_print_interval_ms:1000}")
private long tsStatsPrintIntervalMs;
protected InsertTsRepository<T> insertRepository;
protected TbSqlBlockingQueue<EntityContainer<T>> tsQueue;
@ -76,26 +63,39 @@ public abstract class AbstractSimpleSqlTimeseriesDao<T extends AbstractTsKvEntit
}
}
protected ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(entityId, query);
} else {
long stepTs = query.getStartTs();
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>();
while (stepTs < query.getEndTs()) {
long startTs = stepTs;
long endTs = stepTs + query.getInterval();
long ts = startTs + (endTs - startTs) / 2;
futures.add(findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()));
stepTs = endTs;
}
return getTskvEntriesFuture(Futures.allAsList(futures));
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) {
switch (aggregation) {
case AVG:
findAvg(tenantId, entityId, key, startTs, endTs, entitiesFutures);
break;
case MAX:
findMax(tenantId, entityId, key, startTs, endTs, entitiesFutures);
break;
case MIN:
findMin(tenantId, entityId, key, startTs, endTs, entitiesFutures);
break;
case SUM:
findSum(tenantId, entityId, key, startTs, endTs, entitiesFutures);
break;
case COUNT:
findCount(tenantId, entityId, key, startTs, endTs, entitiesFutures);
break;
default:
throw new IllegalArgumentException("Not supported aggregation type: " + aggregation);
}
}
protected abstract ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation);
protected abstract void findCount(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findSum(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findMin(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findMax(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query);
protected abstract void findAvg(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected SettableFuture<T> setFutures(List<CompletableFuture<T>> entitiesFutures) {
SettableFuture<T> listenableFuture = SettableFuture.create();
@ -121,36 +121,4 @@ public abstract class AbstractSimpleSqlTimeseriesDao<T extends AbstractTsKvEntit
});
return listenableFuture;
}
protected void switchAgregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List<CompletableFuture<T>> entitiesFutures) {
switch (aggregation) {
case AVG:
findAvg(entityId, key, startTs, endTs, entitiesFutures);
break;
case MAX:
findMax(entityId, key, startTs, endTs, entitiesFutures);
break;
case MIN:
findMin(entityId, key, startTs, endTs, entitiesFutures);
break;
case SUM:
findSum(entityId, key, startTs, endTs, entitiesFutures);
break;
case COUNT:
findCount(entityId, key, startTs, endTs, entitiesFutures);
break;
default:
throw new IllegalArgumentException("Not supported aggregation type: " + aggregation);
}
}
protected abstract void findCount(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findSum(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findMin(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findMax(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
protected abstract void findAvg(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<T>> entitiesFutures);
}
}

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

@ -21,9 +21,9 @@ import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.thingsboard.server.common.data.UUIDConverter;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.Aggregation;
@ -34,12 +34,17 @@ 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;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueue;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository;
import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository;
import org.thingsboard.server.dao.sqlts.latest.TsKvLatestRepository;
import org.thingsboard.server.dao.timeseries.SimpleListenableFuture;
@ -49,24 +54,34 @@ import javax.annotation.PreDestroy;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID;
@Slf4j
public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService {
private static final String DESC_ORDER = "DESC";
private final ConcurrentMap<String, Integer> tsKvDictionaryMap = new ConcurrentHashMap<>();
private static final ReentrantLock tsCreationLock = new ReentrantLock();
@Autowired
private TsKvLatestRepository tsKvLatestRepository;
@Autowired
private SearchTsKvLatestRepository searchTsKvLatestRepository;
@Autowired
private InsertLatestRepository insertLatestRepository;
@Autowired
protected ScheduledLogExecutorComponent logExecutor;
private TsKvDictionaryRepository dictionaryRepository;
private TbSqlBlockingQueue<TsKvLatestEntity> tsLatestQueue;
@Value("${sql.ts_latest.batch_size:1000}")
private int tsLatestBatchSize;
@ -77,7 +92,17 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
@Value("${sql.ts_latest.stats_print_interval_ms:1000}")
private long tsLatestStatsPrintIntervalMs;
private TbSqlBlockingQueue<TsKvLatestEntity> tsLatestQueue;
@Autowired
protected ScheduledLogExecutorComponent logExecutor;
@Value("${sql.ts.batch_size:1000}")
protected int tsBatchSize;
@Value("${sql.ts.batch_max_delay:100}")
protected long tsMaxDelay;
@Value("${sql.ts.stats_print_interval_ms:1000}")
protected long tsStatsPrintIntervalMs;
@PostConstruct
protected void init() {
@ -120,6 +145,8 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
protected abstract ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query);
protected abstract ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query);
protected ListenableFuture<List<TsKvEntry>> getTskvEntriesFuture(ListenableFuture<List<Optional<TsKvEntry>>> future) {
return Futures.transform(future, new Function<List<Optional<TsKvEntry>>, List<TsKvEntry>>() {
@Nullable
@ -147,13 +174,14 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
protected ListenableFuture<TsKvEntry> getFindLatestFuture(EntityId entityId, String key) {
TsKvLatestCompositeKey compositeKey =
new TsKvLatestCompositeKey(
entityId.getEntityType(),
fromTimeUUID(entityId.getId()),
key);
entityId.getId(),
getOrSaveKeyId(key));
Optional<TsKvLatestEntity> entry = tsKvLatestRepository.findById(compositeKey);
TsKvEntry result;
if (entry.isPresent()) {
result = DaoUtil.getData(entry.get());
TsKvLatestEntity tsKvLatestEntity = entry.get();
tsKvLatestEntity.setStrKey(key);
result = DaoUtil.getData(tsKvLatestEntity);
} else {
result = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(key, null));
}
@ -171,9 +199,8 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
ListenableFuture<Void> removedLatestFuture = Futures.transformAsync(booleanFuture, isRemove -> {
if (isRemove) {
TsKvLatestEntity latestEntity = new TsKvLatestEntity();
latestEntity.setEntityType(entityId.getEntityType());
latestEntity.setEntityId(fromTimeUUID(entityId.getId()));
latestEntity.setKey(query.getKey());
latestEntity.setEntityId(entityId.getId());
latestEntity.setKey(getOrSaveKeyId(query.getKey()));
return service.submit(() -> {
tsKvLatestRepository.delete(latestEntity);
return null;
@ -215,17 +242,14 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
protected ListenableFuture<List<TsKvEntry>> getFindAllLatestFuture(EntityId entityId) {
return Futures.immediateFuture(
DaoUtil.convertDataList(Lists.newArrayList(
tsKvLatestRepository.findAllByEntityTypeAndEntityId(
entityId.getEntityType(),
UUIDConverter.fromTimeUUID(entityId.getId())))));
searchTsKvLatestRepository.findAllByEntityId(entityId.getId()))));
}
protected ListenableFuture<Void> getSaveLatestFuture(EntityId entityId, TsKvEntry tsKvEntry) {
TsKvLatestEntity latestEntity = new TsKvLatestEntity();
latestEntity.setEntityType(entityId.getEntityType());
latestEntity.setEntityId(fromTimeUUID(entityId.getId()));
latestEntity.setEntityId(entityId.getId());
latestEntity.setTs(tsKvEntry.getTs());
latestEntity.setKey(tsKvEntry.getKey());
latestEntity.setKey(getOrSaveKeyId(tsKvEntry.getKey()));
latestEntity.setStrValue(tsKvEntry.getStrValue().orElse(null));
latestEntity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null));
latestEntity.setLongValue(tsKvEntry.getLongValue().orElse(null));
@ -233,6 +257,42 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx
return tsLatestQueue.add(latestEntity);
}
protected Integer getOrSaveKeyId(String strKey) {
Integer keyId = tsKvDictionaryMap.get(strKey);
if (keyId == null) {
Optional<TsKvDictionary> tsKvDictionaryOptional;
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
tsCreationLock.lock();
try {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
TsKvDictionary tsKvDictionary = new TsKvDictionary();
tsKvDictionary.setKey(strKey);
try {
TsKvDictionary saved = dictionaryRepository.save(tsKvDictionary);
tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId());
keyId = saved.getKeyId();
} catch (ConstraintViolationException e) {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
TsKvDictionary dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!"));
tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId());
keyId = dictionary.getKeyId();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
}
} finally {
tsCreationLock.unlock();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
tsKvDictionaryMap.put(strKey, keyId);
}
}
return keyId;
}
private ListenableFuture<Void> getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) {
ListenableFuture<List<TsKvEntry>> future = findNewLatestEntryFuture(tenantId, entityId, query);
return Futures.transformAsync(future, entryList -> {

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/TsKvDictionaryRepository.java

@ -18,11 +18,11 @@ package org.thingsboard.server.dao.sqlts.dictionary;
import org.springframework.data.repository.CrudRepository;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey;
import org.thingsboard.server.dao.util.PsqlDao;
import org.thingsboard.server.dao.util.SqlTsAnyDao;
import java.util.Optional;
@PsqlDao
@SqlTsAnyDao
public interface TsKvDictionaryRepository extends CrudRepository<TsKvDictionary, TsKvDictionaryCompositeKey> {
Optional<TsKvDictionary> findByKeyId(int keyId);

32
dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/HsqlTimeseriesInsertRepository.java

@ -37,15 +37,14 @@ import java.util.List;
public class HsqlTimeseriesInsertRepository extends AbstractInsertRepository implements InsertTsRepository<TsKvEntity> {
private static final String INSERT_OR_UPDATE =
"MERGE INTO ts_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?) " +
"T (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"ON (ts_kv.entity_type=T.entity_type " +
"AND ts_kv.entity_id=T.entity_id " +
"MERGE INTO ts_kv USING(VALUES ?, ?, ?, ?, ?, ?, ?) " +
"T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"ON (ts_kv.entity_id=T.entity_id " +
"AND ts_kv.key=T.key " +
"AND ts_kv.ts=T.ts) " +
"WHEN MATCHED THEN UPDATE SET ts_kv.bool_v = T.bool_v, ts_kv.str_v = T.str_v, ts_kv.long_v = T.long_v, ts_kv.dbl_v = T.dbl_v " +
"WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"VALUES (T.entity_type, T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v);";
"WHEN NOT MATCHED THEN INSERT (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"VALUES (T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v);";
@Override
public void saveOrUpdate(List<EntityContainer<TsKvEntity>> entities) {
@ -54,29 +53,28 @@ public class HsqlTimeseriesInsertRepository extends AbstractInsertRepository imp
public void setValues(PreparedStatement ps, int i) throws SQLException {
EntityContainer<TsKvEntity> tsKvEntityEntityContainer = entities.get(i);
TsKvEntity tsKvEntity = tsKvEntityEntityContainer.getEntity();
ps.setString(1, tsKvEntity.getEntityType().name());
ps.setString(2, tsKvEntity.getEntityId());
ps.setString(3, tsKvEntity.getKey());
ps.setLong(4, tsKvEntity.getTs());
ps.setObject(1, tsKvEntity.getEntityId());
ps.setInt(2, tsKvEntity.getKey());
ps.setLong(3, tsKvEntity.getTs());
if (tsKvEntity.getBooleanValue() != null) {
ps.setBoolean(5, tsKvEntity.getBooleanValue());
ps.setBoolean(4, tsKvEntity.getBooleanValue());
} else {
ps.setNull(5, Types.BOOLEAN);
ps.setNull(4, Types.BOOLEAN);
}
ps.setString(6, tsKvEntity.getStrValue());
ps.setString(5, tsKvEntity.getStrValue());
if (tsKvEntity.getLongValue() != null) {
ps.setLong(7, tsKvEntity.getLongValue());
ps.setLong(6, tsKvEntity.getLongValue());
} else {
ps.setNull(7, Types.BIGINT);
ps.setNull(6, Types.BIGINT);
}
if (tsKvEntity.getDoubleValue() != null) {
ps.setDouble(8, tsKvEntity.getDoubleValue());
ps.setDouble(7, tsKvEntity.getDoubleValue());
} else {
ps.setNull(8, Types.DOUBLE);
ps.setNull(7, Types.DOUBLE);
}
}

131
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.AbstractSimpleSqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.AbstractPsqlHsqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.EntityContainer;
import org.thingsboard.server.dao.timeseries.TimeseriesDao;
import org.thingsboard.server.dao.util.HsqlDao;
@ -41,14 +41,12 @@ import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID;
@Component
@Slf4j
@SqlTsDao
@HsqlDao
public class JpaHsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
public class JpaHsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
@Autowired
private TsKvHsqlRepository tsKvRepository;
@ -60,11 +58,12 @@ public class JpaHsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
@Override
public ListenableFuture<Void> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl) {
String strKey = tsKvEntry.getKey();
Integer keyId = getOrSaveKeyId(strKey);
TsKvEntity entity = new TsKvEntity();
entity.setEntityType(entityId.getEntityType());
entity.setEntityId(fromTimeUUID(entityId.getId()));
entity.setEntityId(entityId.getId());
entity.setTs(tsKvEntry.getTs());
entity.setKey(tsKvEntry.getKey());
entity.setKey(keyId);
entity.setStrValue(tsKvEntry.getStrValue().orElse(null));
entity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null));
entity.setLongValue(tsKvEntry.getLongValue().orElse(null));
@ -77,9 +76,8 @@ public class JpaHsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
public ListenableFuture<Void> remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) {
return service.submit(() -> {
tsKvRepository.delete(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
query.getKey(),
entityId.getId(),
getOrSaveKeyId(query.getKey()),
query.getStartTs(),
query.getEndTs());
return null;
@ -116,14 +114,47 @@ public class JpaHsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
return Futures.immediateFuture(null);
}
protected ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) {
protected ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(tenantId, entityId, query);
} else {
long stepTs = query.getStartTs();
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>();
while (stepTs < query.getEndTs()) {
long startTs = stepTs;
long endTs = stepTs + query.getInterval();
long ts = startTs + (endTs - startTs) / 2;
futures.add(findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()));
stepTs = endTs;
}
return getTskvEntriesFuture(Futures.allAsList(futures));
}
}
@Override
protected ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
List<TsKvEntity> tsKvEntities = tsKvRepository.findAllWithLimit(
entityId.getId(),
getOrSaveKeyId(query.getKey()),
query.getStartTs(),
query.getEndTs(),
new PageRequest(0, query.getLimit(),
new Sort(Sort.Direction.fromString(
query.getOrderBy()), "ts")));
tsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(query.getKey()));
return Futures.immediateFuture(
DaoUtil.convertDataList(
tsKvEntities));
}
@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(entityId, key, startTs, endTs, aggregation, entitiesFutures);
switchAgregation(tenantId, entityId, key, startTs, endTs, aggregation, entitiesFutures);
return Futures.transform(setFutures(entitiesFutures), entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(fromTimeUUID(entityId.getId()));
entity.setEntityType(entityId.getEntityType());
entity.setKey(key);
entity.setEntityId(entityId.getId());
entity.setKey(getOrSaveKeyId(key));
entity.setTs(ts);
return Optional.of(DaoUtil.getData(entity));
} else {
@ -132,75 +163,63 @@ public class JpaHsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
});
}
protected ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) {
return Futures.immediateFuture(
DaoUtil.convertDataList(
tsKvRepository.findAllWithLimit(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
query.getKey(),
query.getStartTs(),
query.getEndTs(),
new PageRequest(0, query.getLimit(),
new Sort(Sort.Direction.fromString(
query.getOrderBy()), "ts")))));
}
protected void findCount(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findCount(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findCount(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findSum(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findSum(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findSum(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findMin(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findMin(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMin(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
entitiesFutures.add(tsKvRepository.findNumericMin(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findMax(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findMax(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMax(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
entitiesFutures.add(tsKvRepository.findNumericMax(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findAvg(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findAvg(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findAvg(
fromTimeUUID(entityId.getId()),
entityId.getEntityType(),
key,
entityId.getId(),
keyId,
startTs,
endTs));
}
}

76
dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/TsKvHsqlRepository.java

@ -22,23 +22,21 @@ import org.springframework.data.repository.CrudRepository;
import org.springframework.data.repository.query.Param;
import org.springframework.scheduling.annotation.Async;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.dao.model.sqlts.hsql.TsKvCompositeKey;
import org.thingsboard.server.dao.model.sqlts.hsql.TsKvEntity;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
@SqlDao
public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompositeKey> {
@Query("SELECT tskv FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.entityType = :entityType AND tskv.key = :entityKey " +
"AND tskv.ts > :startTs AND tskv.ts <= :endTs")
List<TsKvEntity> findAllWithLimit(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String key,
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
List<TsKvEntity> findAllWithLimit(@Param("entityId") UUID entityId,
@Param("entityKey") int key,
@Param("startTs") long startTs,
@Param("endTs") long endTs,
Pageable pageable);
@ -46,22 +44,18 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
@Transactional
@Modifying
@Query("DELETE FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.entityType = :entityType AND tskv.key = :entityKey " +
"AND tskv.ts > :startTs AND tskv.ts <= :endTs")
void delete(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String key,
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
void delete(@Param("entityId") UUID entityId,
@Param("entityKey") int key,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(MAX(tskv.strValue)) FROM TsKvEntity tskv " +
"WHERE tskv.strValue IS NOT NULL " +
"AND tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findStringMax(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
"WHERE tskv.strValue IS NOT NULL AND tskv.entityId = :entityId AND tskv.key = :entityKey" +
" AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findStringMax(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@ -70,24 +64,20 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
"MAX(COALESCE(tskv.doubleValue, -1.79769E+308)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'MAX') FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"'MAX') FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findNumericMax(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
CompletableFuture<TsKvEntity> findNumericMax(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(MIN(tskv.strValue)) FROM TsKvEntity tskv " +
"WHERE tskv.strValue IS NOT NULL " +
"AND tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"WHERE tskv.strValue IS NOT NULL AND tskv.entityId = :entityId " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findStringMin(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
CompletableFuture<TsKvEntity> findStringMin(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@ -96,12 +86,10 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
"MIN(COALESCE(tskv.doubleValue, 1.79769E+308)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'MIN') FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"'MIN') FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findNumericMin(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
CompletableFuture<TsKvEntity> findNumericMin(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@ -110,11 +98,9 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
"SUM(CASE WHEN tskv.strValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findCount(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findCount(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@ -123,12 +109,10 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
"SUM(COALESCE(tskv.doubleValue, 0.0)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'AVG') FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"'AVG') FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findAvg(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
CompletableFuture<TsKvEntity> findAvg(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);
@ -137,12 +121,10 @@ public interface TsKvHsqlRepository extends CrudRepository<TsKvEntity, TsKvCompo
"SUM(COALESCE(tskv.doubleValue, 0.0)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'SUM') FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " +
"'SUM') FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " +
"AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts <= :endTs")
CompletableFuture<TsKvEntity> findSum(@Param("entityId") String entityId,
@Param("entityType") EntityType entityType,
@Param("entityKey") String entityKey,
CompletableFuture<TsKvEntity> findSum(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey,
@Param("startTs") long startTs,
@Param("endTs") long endTs);

32
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/HsqlLatestInsertRepository.java

@ -36,43 +36,41 @@ import java.util.List;
public class HsqlLatestInsertRepository extends AbstractInsertRepository implements InsertLatestRepository {
private static final String INSERT_OR_UPDATE =
"MERGE INTO ts_kv_latest USING(VALUES ?, ?, ?, ?, ?, ?, ?, ?) " +
"T (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"ON (ts_kv_latest.entity_type=T.entity_type " +
"AND ts_kv_latest.entity_id=T.entity_id " +
"MERGE INTO ts_kv_latest USING(VALUES ?, ?, ?, ?, ?, ?, ?) " +
"T (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"ON (ts_kv_latest.entity_id=T.entity_id " +
"AND ts_kv_latest.key=T.key) " +
"WHEN MATCHED THEN UPDATE SET ts_kv_latest.ts = T.ts, ts_kv_latest.bool_v = T.bool_v, ts_kv_latest.str_v = T.str_v, ts_kv_latest.long_v = T.long_v, ts_kv_latest.dbl_v = T.dbl_v " +
"WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"VALUES (T.entity_type, T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v);";
"WHEN NOT MATCHED THEN INSERT (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) " +
"VALUES (T.entity_id, T.key, T.ts, T.bool_v, T.str_v, T.long_v, T.dbl_v);";
@Override
public void saveOrUpdate(List<TsKvLatestEntity> entities) {
jdbcTemplate.batchUpdate(INSERT_OR_UPDATE, new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
ps.setString(1, entities.get(i).getEntityType().name());
ps.setString(2, entities.get(i).getEntityId());
ps.setString(3, entities.get(i).getKey());
ps.setLong(4, entities.get(i).getTs());
ps.setObject(1, entities.get(i).getEntityId());
ps.setInt(2, entities.get(i).getKey());
ps.setLong(3, entities.get(i).getTs());
if (entities.get(i).getBooleanValue() != null) {
ps.setBoolean(5, entities.get(i).getBooleanValue());
ps.setBoolean(4, entities.get(i).getBooleanValue());
} else {
ps.setNull(5, Types.BOOLEAN);
ps.setNull(4, Types.BOOLEAN);
}
ps.setString(6, entities.get(i).getStrValue());
ps.setString(5, entities.get(i).getStrValue());
if (entities.get(i).getLongValue() != null) {
ps.setLong(7, entities.get(i).getLongValue());
ps.setLong(6, entities.get(i).getLongValue());
} else {
ps.setNull(7, Types.BIGINT);
ps.setNull(6, Types.BIGINT);
}
if (entities.get(i).getDoubleValue() != null) {
ps.setDouble(8, entities.get(i).getDoubleValue());
ps.setDouble(7, entities.get(i).getDoubleValue());
} else {
ps.setNull(8, Types.DOUBLE);
ps.setNull(7, Types.DOUBLE);
}
}

48
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/PsqlLatestInsertRepository.java

@ -38,12 +38,12 @@ import java.util.List;
public class PsqlLatestInsertRepository extends AbstractInsertRepository implements InsertLatestRepository {
private static final String BATCH_UPDATE =
"UPDATE ts_kv_latest SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ? WHERE entity_type = ? AND entity_id = ? and key = ?";
"UPDATE ts_kv_latest SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ? WHERE entity_id = ? and key = ?";
private static final String INSERT_OR_UPDATE =
"INSERT INTO ts_kv_latest (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES(?, ?, ?, ?, ?, ?, ?, ?) " +
"ON CONFLICT (entity_type, entity_id, key) DO UPDATE SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?;";
"INSERT INTO ts_kv_latest (entity_id, key, ts, bool_v, str_v, long_v, dbl_v) VALUES(?, ?, ?, ?, ?, ?, ?) " +
"ON CONFLICT (entity_id, key) DO UPDATE SET ts = ?, bool_v = ?, str_v = ?, long_v = ?, dbl_v = ?;";
@Override
public void saveOrUpdate(List<TsKvLatestEntity> entities) {
@ -76,9 +76,8 @@ public class PsqlLatestInsertRepository extends AbstractInsertRepository impleme
ps.setNull(5, Types.DOUBLE);
}
ps.setString(6, tsKvLatestEntity.getEntityType().name());
ps.setString(7, tsKvLatestEntity.getEntityId());
ps.setString(8, tsKvLatestEntity.getKey());
ps.setObject(6, tsKvLatestEntity.getEntityId());
ps.setInt(7, tsKvLatestEntity.getKey());
}
@Override
@ -105,38 +104,37 @@ public class PsqlLatestInsertRepository extends AbstractInsertRepository impleme
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
TsKvLatestEntity tsKvLatestEntity = insertEntities.get(i);
ps.setString(1, tsKvLatestEntity.getEntityType().name());
ps.setString(2, tsKvLatestEntity.getEntityId());
ps.setString(3, tsKvLatestEntity.getKey());
ps.setLong(4, tsKvLatestEntity.getTs());
ps.setLong(9, tsKvLatestEntity.getTs());
ps.setObject(1, tsKvLatestEntity.getEntityId());
ps.setInt(2, tsKvLatestEntity.getKey());
ps.setLong(3, tsKvLatestEntity.getTs());
ps.setLong(8, tsKvLatestEntity.getTs());
if (tsKvLatestEntity.getBooleanValue() != null) {
ps.setBoolean(5, tsKvLatestEntity.getBooleanValue());
ps.setBoolean(10, tsKvLatestEntity.getBooleanValue());
ps.setBoolean(4, tsKvLatestEntity.getBooleanValue());
ps.setBoolean(9, tsKvLatestEntity.getBooleanValue());
} else {
ps.setNull(5, Types.BOOLEAN);
ps.setNull(10, Types.BOOLEAN);
ps.setNull(4, Types.BOOLEAN);
ps.setNull(9, Types.BOOLEAN);
}
ps.setString(6, replaceNullChars(tsKvLatestEntity.getStrValue()));
ps.setString(11, replaceNullChars(tsKvLatestEntity.getStrValue()));
ps.setString(5, replaceNullChars(tsKvLatestEntity.getStrValue()));
ps.setString(10, replaceNullChars(tsKvLatestEntity.getStrValue()));
if (tsKvLatestEntity.getLongValue() != null) {
ps.setLong(7, tsKvLatestEntity.getLongValue());
ps.setLong(12, tsKvLatestEntity.getLongValue());
ps.setLong(6, tsKvLatestEntity.getLongValue());
ps.setLong(11, tsKvLatestEntity.getLongValue());
} else {
ps.setNull(7, Types.BIGINT);
ps.setNull(12, Types.BIGINT);
ps.setNull(6, Types.BIGINT);
ps.setNull(11, Types.BIGINT);
}
if (tsKvLatestEntity.getDoubleValue() != null) {
ps.setDouble(8, tsKvLatestEntity.getDoubleValue());
ps.setDouble(13, tsKvLatestEntity.getDoubleValue());
ps.setDouble(7, tsKvLatestEntity.getDoubleValue());
ps.setDouble(12, tsKvLatestEntity.getDoubleValue());
} else {
ps.setNull(8, Types.DOUBLE);
ps.setNull(13, Types.DOUBLE);
ps.setNull(7, Types.DOUBLE);
ps.setNull(12, Types.DOUBLE);
}
}

45
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/SearchTsKvLatestRepository.java

@ -0,0 +1,45 @@
/**
* 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.dao.sqlts.latest;
import org.springframework.stereotype.Repository;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.util.SqlTsAnyDao;
import javax.persistence.EntityManager;
import javax.persistence.PersistenceContext;
import java.util.List;
import java.util.UUID;
@SqlTsAnyDao
@Repository
public class SearchTsKvLatestRepository {
public static final String FIND_ALL_BY_ENTITY_ID = "findAllByEntityId";
public static final String FIND_ALL_BY_ENTITY_ID_QUERY = "SELECT ts_kv_latest.entity_id AS entityId, ts_kv_latest.key AS key, ts_kv_dictionary.key AS strKey, ts_kv_latest.str_v AS strValue," +
" ts_kv_latest.bool_v AS boolValue, ts_kv_latest.long_v AS longValue, ts_kv_latest.dbl_v AS doubleValue, ts_kv_latest.ts AS ts FROM ts_kv_latest " +
"INNER JOIN ts_kv_dictionary ON ts_kv_latest.key = ts_kv_dictionary.key_id WHERE ts_kv_latest.entity_id = cast(:id AS uuid)";
@PersistenceContext
private EntityManager entityManager;
public List<TsKvLatestEntity> findAllByEntityId(UUID entityId) {
return entityManager.createNamedQuery(FIND_ALL_BY_ENTITY_ID, TsKvLatestEntity.class)
.setParameter("id", entityId)
.getResultList();
}
}

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java

@ -16,15 +16,15 @@
package org.thingsboard.server.dao.sqlts.latest;
import org.springframework.data.repository.CrudRepository;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestCompositeKey;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.List;
import java.util.UUID;
@SqlDao
public interface TsKvLatestRepository extends CrudRepository<TsKvLatestEntity, TsKvLatestCompositeKey> {
List<TsKvLatestEntity> findAllByEntityTypeAndEntityId(EntityType entityType, String entityId);
List<TsKvLatestEntity> findAllByEntityId(UUID entityId);
}

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

@ -18,7 +18,6 @@ package org.thingsboard.server.dao.sqlts.psql;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.domain.PageRequest;
@ -31,12 +30,9 @@ import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
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.dictionary.TsKvDictionary;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey;
import org.thingsboard.server.dao.model.sqlts.psql.TsKvEntity;
import org.thingsboard.server.dao.sqlts.AbstractSimpleSqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.AbstractPsqlHsqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.EntityContainer;
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository;
import org.thingsboard.server.dao.timeseries.PsqlPartition;
import org.thingsboard.server.dao.timeseries.SqlTsPartitionDate;
import org.thingsboard.server.dao.timeseries.TimeseriesDao;
@ -48,10 +44,12 @@ import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.time.ZonedDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.ReentrantLock;
import static org.thingsboard.server.dao.timeseries.SqlTsPartitionDate.EPOCH_START;
@ -61,17 +59,11 @@ import static org.thingsboard.server.dao.timeseries.SqlTsPartitionDate.EPOCH_STA
@Slf4j
@SqlTsDao
@PsqlDao
public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
public class JpaPsqlTimeseriesDao extends AbstractPsqlHsqlTimeseriesDao<TsKvEntity> implements TimeseriesDao {
private final ConcurrentMap<String, Integer> tsKvDictionaryMap = new ConcurrentHashMap<>();
private final Map<Long, PsqlPartition> partitions = new ConcurrentHashMap<>();
private static final ReentrantLock tsCreationLock = new ReentrantLock();
private static final ReentrantLock partitionCreationLock = new ReentrantLock();
@Autowired
private TsKvDictionaryRepository dictionaryRepository;
@Autowired
private TsKvPsqlRepository tsKvRepository;
@ -100,11 +92,6 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
}
}
@Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return processFindAllAsync(tenantId, entityId, queries);
}
@Override
public ListenableFuture<Void> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl) {
String strKey = tsKvEntry.getKey();
@ -166,22 +153,25 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
return Futures.immediateFuture(null);
}
protected ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) {
List<CompletableFuture<TsKvEntity>> entitiesFutures = new ArrayList<>();
switchAgregation(entityId, key, startTs, endTs, aggregation, entitiesFutures);
return Futures.transform(setFutures(entitiesFutures), entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());
entity.setStrKey(key);
entity.setTs(ts);
return Optional.of(DaoUtil.getData(entity));
} else {
return Optional.empty();
protected ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(tenantId, entityId, query);
} else {
long stepTs = query.getStartTs();
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>();
while (stepTs < query.getEndTs()) {
long startTs = stepTs;
long endTs = stepTs + query.getInterval();
long ts = startTs + (endTs - startTs) / 2;
futures.add(findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()));
stepTs = endTs;
}
});
return getTskvEntriesFuture(Futures.allAsList(futures));
}
}
protected ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) {
@Override
protected ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
Integer keyId = getOrSaveKeyId(query.getKey());
List<TsKvEntity> tsKvEntities = tsKvRepository.findAllWithLimit(
entityId.getId(),
@ -195,7 +185,24 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
return Futures.immediateFuture(DaoUtil.convertDataList(tsKvEntities));
}
protected void findCount(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@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);
return Futures.transform(setFutures(entitiesFutures), entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());
entity.setStrKey(key);
entity.setTs(ts);
return Optional.of(DaoUtil.getData(entity));
} else {
return Optional.empty();
}
});
}
@Override
protected void findCount(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findCount(
entityId.getId(),
@ -204,7 +211,8 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
endTs));
}
protected void findSum(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findSum(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findSum(
entityId.getId(),
@ -213,7 +221,8 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
endTs));
}
protected void findMin(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findMin(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMin(
entityId.getId(),
@ -227,7 +236,8 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
endTs));
}
protected void findMax(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findMax(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMax(
entityId.getId(),
@ -241,7 +251,8 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
endTs));
}
protected void findAvg(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
@Override
protected void findAvg(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findAvg(
entityId.getId(),
@ -250,40 +261,9 @@ public class JpaPsqlTimeseriesDao extends AbstractSimpleSqlTimeseriesDao<TsKvEnt
endTs));
}
private Integer getOrSaveKeyId(String strKey) {
Integer keyId = tsKvDictionaryMap.get(strKey);
if (keyId == null) {
Optional<TsKvDictionary> tsKvDictionaryOptional;
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
tsCreationLock.lock();
try {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
TsKvDictionary tsKvDictionary = new TsKvDictionary();
tsKvDictionary.setKey(strKey);
try {
TsKvDictionary saved = dictionaryRepository.save(tsKvDictionary);
tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId());
keyId = saved.getKeyId();
} catch (ConstraintViolationException e) {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
TsKvDictionary dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!"));
tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId());
keyId = dictionary.getKeyId();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
}
} finally {
tsCreationLock.unlock();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
tsKvDictionaryMap.put(strKey, keyId);
}
}
return keyId;
@Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return processFindAllAsync(tenantId, entityId, queries);
}
private void savePartition(PsqlPartition psqlPartition) {

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

@ -19,9 +19,7 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.domain.PageRequest;
import org.springframework.data.domain.Sort;
import org.springframework.stereotype.Component;
@ -33,15 +31,12 @@ import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
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.dictionary.TsKvDictionary;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey;
import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvEntity;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueue;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao;
import org.thingsboard.server.dao.sqlts.EntityContainer;
import org.thingsboard.server.dao.sqlts.InsertTsRepository;
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository;
import org.thingsboard.server.dao.timeseries.TimeseriesDao;
import org.thingsboard.server.dao.util.TimescaleDBTsDao;
@ -53,25 +48,12 @@ import java.util.List;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.ReentrantLock;
@Component
@Slf4j
@TimescaleDBTsDao
public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao {
private static final String TS = "ts";
private final ConcurrentMap<String, Integer> tsKvDictionaryMap = new ConcurrentHashMap<>();
private static final ReentrantLock tsCreationLock = new ReentrantLock();
@Autowired
private TsKvDictionaryRepository dictionaryRepository;
@Autowired
private TsKvTimescaleRepository tsKvRepository;
@ -79,40 +61,32 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
private AggregationRepository aggregationRepository;
@Autowired
private InsertTsRepository<TimescaleTsKvEntity> insertRepository;
@Value("${sql.ts_timescale.batch_size:1000}")
private int batchSize;
@Value("${sql.ts_timescale.batch_max_delay:100}")
private long maxDelay;
protected InsertTsRepository<TimescaleTsKvEntity> insertRepository;
@Value("${sql.ts_timescale.stats_print_interval_ms:1000}")
private long statsPrintIntervalMs;
private TbSqlBlockingQueue<EntityContainer<TimescaleTsKvEntity>> queue;
protected TbSqlBlockingQueue<EntityContainer<TimescaleTsKvEntity>> tsQueue;
@PostConstruct
protected void init() {
super.init();
TbSqlBlockingQueueParams params = TbSqlBlockingQueueParams.builder()
TbSqlBlockingQueueParams tsParams = TbSqlBlockingQueueParams.builder()
.logName("TS Timescale")
.batchSize(batchSize)
.maxDelay(maxDelay)
.statsPrintIntervalMs(statsPrintIntervalMs)
.batchSize(tsBatchSize)
.maxDelay(tsMaxDelay)
.statsPrintIntervalMs(tsStatsPrintIntervalMs)
.build();
queue = new TbSqlBlockingQueue<>(params);
queue.init(logExecutor, v -> insertRepository.saveOrUpdate(v));
tsQueue = new TbSqlBlockingQueue<>(tsParams);
tsQueue.init(logExecutor, v -> insertRepository.saveOrUpdate(v));
}
@PreDestroy
protected void destroy() {
super.destroy();
if (queue != null) {
queue.destroy();
if (tsQueue != null) {
tsQueue.destroy();
}
}
@Override
protected ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(tenantId, entityId, query);
@ -120,11 +94,58 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
long startTs = query.getStartTs();
long endTs = query.getEndTs();
long timeBucket = query.getInterval();
ListenableFuture<List<Optional<TsKvEntry>>> future = findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation());
ListenableFuture<List<Optional<TsKvEntry>>> future = findAllAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation());
return getTskvEntriesFuture(future);
}
}
@Override
protected ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
String strKey = query.getKey();
Integer keyId = getOrSaveKeyId(strKey);
List<TimescaleTsKvEntity> timescaleTsKvEntities = tsKvRepository.findAllWithLimit(
tenantId.getId(),
entityId.getId(),
keyId,
query.getStartTs(),
query.getEndTs(),
new PageRequest(0, query.getLimit(),
new Sort(Sort.Direction.fromString(
query.getOrderBy()), "ts")));
timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey));
return Futures.immediateFuture(DaoUtil.convertDataList(timescaleTsKvEntities));
}
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());
SettableFuture<List<TimescaleTsKvEntity>> listenableFuture = SettableFuture.create();
listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> {
if (throwable != null) {
listenableFuture.setException(throwable);
} else {
listenableFuture.set(timescaleTsKvEntities);
}
});
return Futures.transform(listenableFuture, timescaleTsKvEntities -> {
if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) {
List<Optional<TsKvEntry>> result = new ArrayList<>();
timescaleTsKvEntities.forEach(entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());
entity.setTenantId(tenantId.getId());
entity.setStrKey(key);
result.add(Optional.of(DaoUtil.getData(entity)));
} else {
result.add(Optional.empty());
}
});
return result;
} else {
return Collections.emptyList();
}
});
}
@Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return processFindAllAsync(tenantId, entityId, queries);
@ -154,7 +175,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
entity.setLongValue(tsKvEntry.getLongValue().orElse(null));
entity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null));
log.trace("Saving entity to timescale db: {}", entity);
return queue.add(new EntityContainer(entity, null));
return tsQueue.add(new EntityContainer(entity, null));
}
@Override
@ -192,88 +213,6 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
return service.submit(() -> null);
}
private ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
String strKey = query.getKey();
Integer keyId = getOrSaveKeyId(strKey);
List<TimescaleTsKvEntity> timescaleTsKvEntities = tsKvRepository.findAllWithLimit(
tenantId.getId(),
entityId.getId(),
keyId,
query.getStartTs(),
query.getEndTs(),
new PageRequest(0, query.getLimit(),
new Sort(Sort.Direction.fromString(
query.getOrderBy()), TS)));
timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey));
return Futures.immediateFuture(DaoUtil.convertDataList(timescaleTsKvEntities));
}
private Integer getOrSaveKeyId(String strKey) {
Integer keyId = tsKvDictionaryMap.get(strKey);
if (keyId == null) {
Optional<TsKvDictionary> tsKvDictionaryOptional;
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
tsCreationLock.lock();
try {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
if (!tsKvDictionaryOptional.isPresent()) {
TsKvDictionary tsKvDictionary = new TsKvDictionary();
tsKvDictionary.setKey(strKey);
try {
TsKvDictionary saved = dictionaryRepository.save(tsKvDictionary);
tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId());
keyId = saved.getKeyId();
} catch (ConstraintViolationException e) {
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey));
TsKvDictionary dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!"));
tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId());
keyId = dictionary.getKeyId();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
}
} finally {
tsCreationLock.unlock();
}
} else {
keyId = tsKvDictionaryOptional.get().getKeyId();
tsKvDictionaryMap.put(strKey, keyId);
}
}
return keyId;
}
private ListenableFuture<List<Optional<TsKvEntry>>> findAndAggregateAsync(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());
SettableFuture<List<TimescaleTsKvEntity>> listenableFuture = SettableFuture.create();
listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> {
if (throwable != null) {
listenableFuture.setException(throwable);
} else {
listenableFuture.set(timescaleTsKvEntities);
}
});
return Futures.transform(listenableFuture, timescaleTsKvEntities -> {
if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) {
List<Optional<TsKvEntry>> result = new ArrayList<>();
timescaleTsKvEntities.forEach(entity -> {
if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId());
entity.setTenantId(tenantId.getId());
entity.setStrKey(key);
result.add(Optional.of(DaoUtil.getData(entity)));
} else {
result.add(Optional.empty());
}
});
return result;
} else {
return Collections.emptyList();
}
});
}
private CompletableFuture<List<TimescaleTsKvEntity>> switchAgregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId, UUID tenantId) {
switch (aggregation) {
case AVG:
@ -291,9 +230,9 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
}
}
private CompletableFuture<List<TimescaleTsKvEntity>> findAvg(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
private CompletableFuture<List<TimescaleTsKvEntity>> findCount(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findAvg(
return aggregationRepository.findCount(
tenantId,
entityId,
keyId,
@ -302,9 +241,9 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findMax(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
private CompletableFuture<List<TimescaleTsKvEntity>> findSum(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findMax(
return aggregationRepository.findSum(
tenantId,
entityId,
keyId,
@ -322,12 +261,11 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
timeBucket,
startTs,
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findSum(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
private CompletableFuture<List<TimescaleTsKvEntity>> findMax(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findSum(
return aggregationRepository.findMax(
tenantId,
entityId,
keyId,
@ -336,9 +274,9 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findCount(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
private CompletableFuture<List<TimescaleTsKvEntity>> findAvg(String key, long startTs, long endTs, long timeBucket, UUID entityId, UUID tenantId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findCount(
return aggregationRepository.findAvg(
tenantId,
entityId,
keyId,

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

@ -25,7 +25,7 @@ CREATE TABLE IF NOT EXISTS tenant_ts_kv (
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_pkey PRIMARY KEY (tenant_id, entity_id, key, ts)
CONSTRAINT tenant_ts_kv_pkey PRIMARY KEY (tenant_id, entity_id, key, ts)
);
CREATE TABLE IF NOT EXISTS ts_kv_dictionary (
@ -35,15 +35,14 @@ CREATE TABLE IF NOT EXISTS ts_kv_dictionary (
);
CREATE TABLE IF NOT EXISTS ts_kv_latest (
entity_type varchar(255) NOT NULL,
entity_id varchar(31) NOT NULL,
key varchar(255) NOT NULL,
entity_id uuid NOT NULL,
key int NOT NULL,
ts bigint NOT NULL,
bool_v boolean,
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_type, entity_id, key)
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);

22
dao/src/main/resources/sql/schema-ts-hsql.sql

@ -14,26 +14,32 @@
-- limitations under the License.
--
SET DATABASE SQL SYNTAX PGS TRUE;
CREATE TABLE IF NOT EXISTS ts_kv (
entity_type varchar(255) NOT NULL,
entity_id varchar(31) NOT NULL,
key varchar(255) NOT NULL,
entity_id uuid NOT NULL,
key int NOT NULL,
ts bigint NOT NULL,
bool_v boolean,
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_type, entity_id, key, ts)
CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_id, key, ts)
);
CREATE TABLE IF NOT EXISTS ts_kv_latest (
entity_type varchar(255) NOT NULL,
entity_id varchar(31) NOT NULL,
key varchar(255) NOT NULL,
entity_id uuid NOT NULL,
key int NOT NULL,
ts bigint NOT NULL,
bool_v boolean,
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_type, entity_id, key)
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key)
);
CREATE TABLE IF NOT EXISTS ts_kv_dictionary (
key varchar(255) NOT NULL,
key_id int GENERATED BY DEFAULT AS IDENTITY(start with 0 increment by 1) UNIQUE,
CONSTRAINT ts_key_id_pkey PRIMARY KEY (key)
);

19
dao/src/main/resources/sql/schema-ts-psql.sql

@ -24,20 +24,19 @@ CREATE TABLE IF NOT EXISTS ts_kv (
dbl_v double precision
) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS ts_kv_dictionary (
key varchar(255) NOT NULL,
key_id serial UNIQUE,
CONSTRAINT ts_key_id_pkey PRIMARY KEY (key)
);
CREATE TABLE IF NOT EXISTS ts_kv_latest (
entity_type varchar(255) NOT NULL,
entity_id varchar(31) NOT NULL,
key varchar(255) NOT NULL,
entity_id uuid NOT NULL,
key int NOT NULL,
ts bigint NOT NULL,
bool_v boolean,
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_type, entity_id, key)
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key)
);
CREATE TABLE IF NOT EXISTS ts_kv_dictionary (
key varchar(255) NOT NULL,
key_id serial UNIQUE,
CONSTRAINT ts_key_id_pkey PRIMARY KEY (key)
);
Loading…
Cancel
Save