From ee7c4f6e7fc322728f9e9bf72f3f65eca404a159 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Mon, 25 May 2020 12:51:00 +0300 Subject: [PATCH 1/3] added inserts strategy for upgrade to 2.5 --- .../upgrade/2.4.3/schema_update_psql_ts.sql | 136 +++++++++++++++--- .../2.4.3/schema_update_timescale_ts.sql | 44 ++++++ .../install/PsqlTsDatabaseUpgradeService.java | 86 +++++++---- .../TimescaleTsDatabaseUpgradeService.java | 59 +++++--- 4 files changed, 267 insertions(+), 58 deletions(-) diff --git a/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql b/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql index 670900ea81..e5cdafa03f 100644 --- a/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql +++ b/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql @@ -36,6 +36,7 @@ BEGIN ALTER COLUMN key TYPE integer USING key::integer; ALTER TABLE ts_kv ADD CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_id, key, ts); + CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT; END; $$; @@ -44,22 +45,40 @@ $$; CREATE OR REPLACE PROCEDURE create_new_ts_kv_latest_table() LANGUAGE plpgsql 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); + IF NOT EXISTS(SELECT FROM pg_tables WHERE schemaname = 'public' AND tablename = 'ts_kv_latest_old') THEN + 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); + ELSE + RAISE NOTICE 'ts_kv_latest_old table already exists!'; + IF NOT EXISTS(SELECT FROM pg_tables WHERE schemaname = 'public' AND tablename = 'ts_kv_latest') THEN + 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, + json_v json, + CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key) + ); + END IF; + END IF; END; $$; @@ -218,4 +237,89 @@ BEGIN END; $$; +-- call insert_into_ts_kv_cursor(); + +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_cursor() LANGUAGE plpgsql AS $$ +DECLARE + insert_size CONSTANT integer := 10000; + insert_counter integer DEFAULT 0; + insert_record RECORD; + insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, + ts_kv_records.key AS key, + ts_kv_records.ts AS ts, + ts_kv_records.bool_v AS bool_v, + ts_kv_records.str_v AS str_v, + ts_kv_records.long_v AS long_v, + ts_kv_records.dbl_v AS dbl_v + FROM (SELECT entity_id AS entity_id, + key_id AS key, + ts, + bool_v, + str_v, + long_v, + dbl_v + FROM ts_kv_old + INNER JOIN ts_kv_dictionary ON (ts_kv_old.key = ts_kv_dictionary.key)) AS ts_kv_records; +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 partitioned ts_kv!',insert_counter - 1; + EXIT; + END IF; + INSERT INTO ts_kv(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 partitioned ts_kv!',insert_counter; + END IF; + END LOOP; + CLOSE insert_cursor; +END; +$$; + +-- call insert_into_ts_kv_latest_cursor(); + +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_latest_cursor() LANGUAGE plpgsql AS $$ +DECLARE + insert_size CONSTANT integer := 10000; + insert_counter integer DEFAULT 0; + insert_record RECORD; + insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, + ts_kv_latest_records.key AS key, + ts_kv_latest_records.ts AS ts, + ts_kv_latest_records.bool_v AS bool_v, + ts_kv_latest_records.str_v AS str_v, + ts_kv_latest_records.long_v AS long_v, + ts_kv_latest_records.dbl_v AS dbl_v + FROM (SELECT entity_id AS entity_id, + 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 ts_kv_latest_records; +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; +$$; diff --git a/application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql b/application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql index 30a76aeb4c..6124910efd 100644 --- a/application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql +++ b/application/src/main/data/upgrade/2.4.3/schema_update_timescale_ts.sql @@ -162,3 +162,47 @@ BEGIN CLOSE insert_cursor; END; $$; + +-- call insert_into_ts_kv_cursor(); + +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_cursor() LANGUAGE plpgsql AS $$ + +DECLARE + insert_size CONSTANT integer := 10000; + insert_counter integer DEFAULT 0; + insert_record RECORD; + insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, + new_ts_kv_records.key AS key, + new_ts_kv_records.ts AS ts, + new_ts_kv_records.bool_v AS bool_v, + new_ts_kv_records.str_v AS str_v, + new_ts_kv_records.long_v AS long_v, + new_ts_kv_records.dbl_v AS dbl_v + FROM (SELECT entity_id AS entity_id, + 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 new_ts_kv_records; +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 ts_kv table!',insert_counter - 1; + EXIT; + END IF; + INSERT INTO ts_kv(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 new ts_kv table!',insert_counter; + END IF; + END LOOP; + CLOSE insert_cursor; +END; +$$; \ No newline at end of file diff --git a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java index bca71fc539..3cca294f20 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java @@ -57,11 +57,15 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe private static final String INSERT_INTO_DICTIONARY = "insert_into_dictionary()"; private static final String INSERT_INTO_TS_KV = "insert_into_ts_kv(IN path_to_file varchar)"; private static final String INSERT_INTO_TS_KV_LATEST = "insert_into_ts_kv_latest(IN path_to_file varchar)"; + private static final String INSERT_INTO_TS_KV_CURSOR = "insert_into_ts_kv_cursor()"; + private static final String INSERT_INTO_TS_KV_LATEST_CURSOR = "insert_into_ts_kv_latest_cursor()"; private static final String CALL_CREATE_PARTITION_TS_KV_TABLE = CALL_REGEX + CREATE_PARTITION_TS_KV_TABLE; private static final String CALL_CREATE_NEW_TS_KV_LATEST_TABLE = CALL_REGEX + CREATE_NEW_TS_KV_LATEST_TABLE; private static final String CALL_CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + CREATE_TS_KV_DICTIONARY_TABLE; private static final String CALL_INSERT_INTO_DICTIONARY = CALL_REGEX + INSERT_INTO_DICTIONARY; + private static final String CALL_INSERT_INTO_TS_KV_CURSOR = CALL_REGEX + INSERT_INTO_TS_KV_CURSOR; + private static final String CALL_INSERT_INTO_TS_KV_LATEST_CURSOR = CALL_REGEX + INSERT_INTO_TS_KV_LATEST_CURSOR; 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; @@ -73,6 +77,8 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe private static final String DROP_PROCEDURE_INSERT_INTO_DICTIONARY = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_DICTIONARY; private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV; private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV_LATEST; + private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV_CURSOR = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV_CURSOR; + private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST_CURSOR = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV_LATEST_CURSOR; private static final String DROP_FUNCTION_GET_PARTITION_DATA = "DROP FUNCTION IF EXISTS get_partitions_data;"; @Override @@ -118,39 +124,44 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe Path tsKvLatestFile = Files.createTempFile(pathToDir, "ts_kv_latest", ".sql"); pathToTempTsKvFile = tsKvFile.toAbsolutePath(); pathToTempTsKvLatestFile = tsKvLatestFile.toAbsolutePath(); - executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); - executeQuery(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE); - executeQuery(conn, "call insert_into_ts_kv_latest('" + pathToTempTsKvLatestFile + "');"); + try { + copyTimeseries(conn, pathToTempTsKvFile, pathToTempTsKvLatestFile); + } catch (Exception e) { + log.info("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); + insertTimeseries(conn); + } } catch (IOException | SecurityException e) { throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); } } else { - Path tempDirPath = Files.createTempDirectory("ts_kv"); - File tempDirAsFile = tempDirPath.toFile(); - boolean writable = tempDirAsFile.setWritable(true, false); - boolean readable = tempDirAsFile.setReadable(true, false); - boolean executable = tempDirAsFile.setExecutable(true, false); - if (writable && readable && executable) { + try { + Path tempDirPath = Files.createTempDirectory("ts_kv"); + File tempDirAsFile = tempDirPath.toFile(); + boolean writable = tempDirAsFile.setWritable(true, false); + boolean readable = tempDirAsFile.setReadable(true, false); + boolean executable = tempDirAsFile.setExecutable(true, false); pathToTempTsKvFile = tempDirPath.resolve(TS_KV_SQL).toAbsolutePath(); pathToTempTsKvLatestFile = tempDirPath.resolve(TS_KV_LATEST_SQL).toAbsolutePath(); - executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); - executeQuery(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE); - executeQuery(conn, "call insert_into_ts_kv_latest('" + pathToTempTsKvLatestFile + "');"); - } else { - throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); - } - } - if (pathToTempTsKvFile.toFile().exists() && pathToTempTsKvLatestFile.toFile().exists()) { - boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); - if (deleteTsKvFile) { - log.info("Successfully deleted the temp file for ts_kv table upgrade!"); - } - boolean deleteTsKvLatestFile = pathToTempTsKvLatestFile.toFile().delete(); - if (deleteTsKvLatestFile) { - log.info("Successfully deleted the temp file for ts_kv_latest table upgrade!"); + try { + if (writable && readable && executable) { + copyTimeseries(conn, pathToTempTsKvFile, pathToTempTsKvLatestFile); + } else { + throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); + } + } catch (Exception e) { + log.info(e.getMessage()); + log.info("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); + insertTimeseries(conn); + } + } catch (IOException | SecurityException e) { + throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); } } + removeUpgradeFiles(pathToTempTsKvFile, pathToTempTsKvLatestFile); + executeQuery(conn, DROP_TABLE_TS_KV_OLD); executeQuery(conn, DROP_TABLE_TS_KV_LATEST_OLD); @@ -161,6 +172,8 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV); executeQuery(conn, DROP_PROCEDURE_CREATE_NEW_TS_KV_LATEST_TABLE); executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST); + executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV_CURSOR); + executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST_CURSOR); executeQuery(conn, DROP_FUNCTION_GET_PARTITION_DATA); executeQuery(conn, "ALTER TABLE ts_kv ADD COLUMN IF NOT EXISTS json_v json;"); @@ -186,6 +199,31 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe } } + private void removeUpgradeFiles(Path pathToTempTsKvFile, Path pathToTempTsKvLatestFile) { + if (pathToTempTsKvFile.toFile().exists() && pathToTempTsKvLatestFile.toFile().exists()) { + boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); + if (deleteTsKvFile) { + log.info("Successfully deleted the temp file for ts_kv table upgrade!"); + } + boolean deleteTsKvLatestFile = pathToTempTsKvLatestFile.toFile().delete(); + if (deleteTsKvLatestFile) { + log.info("Successfully deleted the temp file for ts_kv_latest table upgrade!"); + } + } + } + + private void copyTimeseries(Connection conn, Path pathToTempTsKvFile, Path pathToTempTsKvLatestFile) { + executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); + executeQuery(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE); + executeQuery(conn, "call insert_into_ts_kv_latest('" + pathToTempTsKvLatestFile + "')"); + } + + private void insertTimeseries(Connection conn) { + executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + executeQuery(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE); + executeQuery(conn, CALL_INSERT_INTO_TS_KV_LATEST_CURSOR); + } + @Override protected void loadSql(Connection conn, String fileName) { Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.4.3", fileName); diff --git a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java index a7f243d7b5..625140e47c 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java @@ -53,6 +53,7 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr private static final String CREATE_TS_KV_DICTIONARY_TABLE = "create_ts_kv_dictionary_table()"; private static final String INSERT_INTO_DICTIONARY = "insert_into_dictionary()"; private static final String INSERT_INTO_TS_KV = "insert_into_ts_kv(IN path_to_file varchar)"; + private static final String INSERT_INTO_TS_KV_CURSOR = "insert_into_ts_kv_cursor()"; private static final String INSERT_INTO_TS_KV_LATEST = "insert_into_ts_kv_latest()"; private static final String CALL_CREATE_TS_KV_LATEST_TABLE = CALL_REGEX + CREATE_TS_KV_LATEST_TABLE; @@ -60,6 +61,7 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr private static final String CALL_CREATE_TS_KV_DICTIONARY_TABLE = CALL_REGEX + CREATE_TS_KV_DICTIONARY_TABLE; private static final String CALL_INSERT_INTO_DICTIONARY = CALL_REGEX + INSERT_INTO_DICTIONARY; private static final String CALL_INSERT_INTO_TS_KV_LATEST = CALL_REGEX + INSERT_INTO_TS_KV_LATEST; + private static final String CALL_INSERT_INTO_TS_KV_CURSOR = CALL_REGEX + INSERT_INTO_TS_KV_CURSOR; private static final String DROP_OLD_TENANT_TS_KV_TABLE = DROP_TABLE + TENANT_TS_KV_OLD_TABLE; @@ -68,6 +70,7 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr private static final String DROP_PROCEDURE_CREATE_TS_KV_DICTIONARY_TABLE = DROP_PROCEDURE_IF_EXISTS + CREATE_TS_KV_DICTIONARY_TABLE; private static final String DROP_PROCEDURE_INSERT_INTO_DICTIONARY = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_DICTIONARY; private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV; + private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV_CURSOR = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV_CURSOR; private static final String DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST = DROP_PROCEDURE_IF_EXISTS + INSERT_INTO_TS_KV_LATEST; @Autowired @@ -112,31 +115,41 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr try { Path tsKvFile = Files.createTempFile(pathToDir, "ts_kv", ".sql"); pathToTempTsKvFile = tsKvFile.toAbsolutePath(); - executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); - pathToTempTsKvFile.toFile().deleteOnExit(); + try { + executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); + } catch (Exception e) { + log.info("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); + executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + } } catch (IOException | SecurityException e) { throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); } } else { - Path tempDirPath = Files.createTempDirectory("ts_kv"); - File tempDirAsFile = tempDirPath.toFile(); - boolean writable = tempDirAsFile.setWritable(true, false); - boolean readable = tempDirAsFile.setReadable(true, false); - boolean executable = tempDirAsFile.setExecutable(true, false); - if (writable && readable && executable) { + try { + Path tempDirPath = Files.createTempDirectory("ts_kv"); + File tempDirAsFile = tempDirPath.toFile(); + boolean writable = tempDirAsFile.setWritable(true, false); + boolean readable = tempDirAsFile.setReadable(true, false); + boolean executable = tempDirAsFile.setExecutable(true, false); pathToTempTsKvFile = tempDirPath.resolve(TS_KV_SQL).toAbsolutePath(); - executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); - } else { - throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); - } - } - - if (pathToTempTsKvFile.toFile().exists()) { - boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); - if (deleteTsKvFile) { - log.info("Successfully deleted the temp file for ts_kv table upgrade!"); + try { + if (writable && readable && executable) { + executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); + } else { + throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); + } + } catch (Exception e) { + log.info(e.getMessage()); + log.info("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); + executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + } + } catch (IOException | SecurityException e) { + throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); } } + removeUpgradeFile(pathToTempTsKvFile); executeQuery(conn, CALL_INSERT_INTO_TS_KV_LATEST); @@ -147,6 +160,7 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr executeQuery(conn, DROP_PROCEDURE_CREATE_TS_KV_DICTIONARY_TABLE); executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_DICTIONARY); executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV); + executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV_CURSOR); executeQuery(conn, DROP_PROCEDURE_INSERT_INTO_TS_KV_LATEST); executeQuery(conn, "ALTER TABLE ts_kv ADD COLUMN IF NOT EXISTS json_v json;"); @@ -166,6 +180,15 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr } } + private void removeUpgradeFile(Path pathToTempTsKvFile) { + if (pathToTempTsKvFile.toFile().exists()) { + boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); + if (deleteTsKvFile) { + log.info("Successfully deleted the temp file for ts_kv table upgrade!"); + } + } + } + @Override protected void loadSql(Connection conn, String fileName) { Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.4.3", fileName); From f8a355fec884dc615548bc9e919fc48994885d72 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Mon, 25 May 2020 15:28:55 +0300 Subject: [PATCH 2/3] added upgrade from version 2.5.0 --- .../server/install/ThingsboardInstallService.java | 9 ++++++++- .../service/install/PsqlTsDatabaseSchemaService.java | 4 +--- .../service/install/PsqlTsDatabaseUpgradeService.java | 8 ++++++-- .../install/TimescaleTsDatabaseUpgradeService.java | 5 +++++ .../server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java | 2 +- dao/src/main/resources/sql/schema-timescale.sql | 2 +- dao/src/main/resources/sql/schema-ts-psql.sql | 2 +- 7 files changed, 23 insertions(+), 9 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 8dd45bc41b..e281c0958e 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -133,13 +133,20 @@ public class ThingsboardInstallService { databaseEntitiesUpgradeService.upgradeDatabase("2.4.2"); case "2.4.3": - log.info("Upgrading ThingsBoard from version 2.4.3 to 2.5 ..."); + log.info("Upgrading ThingsBoard from version 2.4.3 to 2.5.0 ..."); if (databaseTsUpgradeService != null) { databaseTsUpgradeService.upgradeDatabase("2.4.3"); } databaseEntitiesUpgradeService.upgradeDatabase("2.4.3"); + case "2.5.0": + log.info("Upgrading ThingsBoard from version 2.5.0 to 2.5.1 ..."); + if (databaseTsUpgradeService != null) { + databaseTsUpgradeService.upgradeDatabase("2.5.0"); + } + + log.info("Updating system data..."); systemDataLoaderService.deleteSystemWidgetBundle("charts"); diff --git a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseSchemaService.java b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseSchemaService.java index bcc3d9bb81..1f2d2f55a8 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseSchemaService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseSchemaService.java @@ -37,8 +37,6 @@ public class PsqlTsDatabaseSchemaService extends SqlAbstractDatabaseSchemaServic @Override public void createDatabaseSchema() throws Exception { super.createDatabaseSchema(); - if (partitionType.equals("INDEFINITE")) { - executeQuery("CREATE TABLE ts_kv_indefinite PARTITION OF ts_kv DEFAULT;"); - } + executeQuery("CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT;"); } } \ No newline at end of file diff --git a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java index 3cca294f20..205973e6eb 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java @@ -99,8 +99,6 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe executeQuery(conn, CALL_CREATE_PARTITION_TS_KV_TABLE); if (!partitionType.equals("INDEFINITE")) { executeQuery(conn, "call create_partitions('" + partitionType + "')"); - } else { - executeQuery(conn, "CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT;"); } executeQuery(conn, CALL_CREATE_TS_KV_DICTIONARY_TABLE); executeQuery(conn, CALL_INSERT_INTO_DICTIONARY); @@ -194,6 +192,12 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe } } break; + case "2.5.0": + try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + executeQuery(conn, "CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT;"); + executeQuery(conn, "UPDATE tb_schema_settings SET schema_version = 2005001"); + } + break; default: throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); } diff --git a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java index 625140e47c..ada1f7ad42 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java @@ -175,6 +175,11 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr } } break; + case "2.5.0": + try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + executeQuery(conn, "UPDATE tb_schema_settings SET schema_version = 2005001"); + } + break; default: throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java index e0854a0a19..7d0aefb62c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java @@ -90,7 +90,7 @@ public class JpaPsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDa } private void savePartitionIfNotExist(long ts) { - if (!tsFormat.equals(SqlTsPartitionDate.INDEFINITE)) { + if (!tsFormat.equals(SqlTsPartitionDate.INDEFINITE) && ts >= 0) { LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); LocalDateTime localDateTimeStart = tsFormat.trancateTo(time); long partitionStartTs = toMills(localDateTimeStart); diff --git a/dao/src/main/resources/sql/schema-timescale.sql b/dao/src/main/resources/sql/schema-timescale.sql index bb0a964c13..926f9ee1a6 100644 --- a/dao/src/main/resources/sql/schema-timescale.sql +++ b/dao/src/main/resources/sql/schema-timescale.sql @@ -52,7 +52,7 @@ CREATE TABLE IF NOT EXISTS tb_schema_settings CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version) ); -INSERT INTO tb_schema_settings (schema_version) VALUES (2005000) ON CONFLICT (schema_version) DO UPDATE SET schema_version = 2005000; +INSERT INTO tb_schema_settings (schema_version) VALUES (2005001) ON CONFLICT (schema_version) DO UPDATE SET schema_version = 2005001; CREATE OR REPLACE FUNCTION to_uuid(IN entity_id varchar, OUT uuid_id uuid) AS $$ diff --git a/dao/src/main/resources/sql/schema-ts-psql.sql b/dao/src/main/resources/sql/schema-ts-psql.sql index 6f6177de01..28420a8957 100644 --- a/dao/src/main/resources/sql/schema-ts-psql.sql +++ b/dao/src/main/resources/sql/schema-ts-psql.sql @@ -53,7 +53,7 @@ CREATE TABLE IF NOT EXISTS tb_schema_settings CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version) ); -INSERT INTO tb_schema_settings (schema_version) VALUES (2005000) ON CONFLICT (schema_version) DO UPDATE SET schema_version = 2005000; +INSERT INTO tb_schema_settings (schema_version) VALUES (2005001) ON CONFLICT (schema_version) DO UPDATE SET schema_version = 2005001; CREATE OR REPLACE PROCEDURE drop_partitions_by_max_ttl(IN partition_type varchar, IN system_ttl bigint, INOUT deleted bigint) LANGUAGE plpgsql AS From 2e08805e07f1558b5779a30308736307c32ca076 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Tue, 26 May 2020 16:14:37 +0300 Subject: [PATCH 3/3] fix upgrade on RDS using insert strategy --- .../upgrade/2.4.3/schema_update_psql_ts.sql | 174 +++++++++++------- .../install/PsqlTsDatabaseUpgradeService.java | 21 ++- .../TimescaleTsDatabaseUpgradeService.java | 25 +-- 3 files changed, 129 insertions(+), 91 deletions(-) diff --git a/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql b/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql index e5cdafa03f..49823fb5af 100644 --- a/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql +++ b/application/src/main/data/upgrade/2.4.3/schema_update_psql_ts.sql @@ -16,51 +16,67 @@ -- call create_partition_ts_kv_table(); -CREATE OR REPLACE PROCEDURE create_partition_ts_kv_table() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE create_partition_ts_kv_table() + LANGUAGE plpgsql AS +$$ BEGIN - ALTER TABLE ts_kv - RENAME TO ts_kv_old; - ALTER TABLE ts_kv_old - RENAME CONSTRAINT ts_kv_pkey TO ts_kv_pkey_old; - CREATE TABLE IF NOT EXISTS ts_kv - ( - LIKE ts_kv_old - ) - PARTITION BY RANGE (ts); - ALTER TABLE ts_kv - DROP COLUMN entity_type; - ALTER TABLE ts_kv - ALTER COLUMN entity_id TYPE uuid USING entity_id::uuid; - ALTER TABLE ts_kv - ALTER COLUMN key TYPE integer USING key::integer; - ALTER TABLE ts_kv - ADD CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_id, key, ts); - CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT; + ALTER TABLE ts_kv + DROP CONSTRAINT IF EXISTS ts_kv_unq_key; + ALTER TABLE ts_kv + DROP CONSTRAINT IF EXISTS ts_kv_pkey; + ALTER TABLE ts_kv + ADD CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_type, entity_id, key, ts); + ALTER TABLE ts_kv + RENAME TO ts_kv_old; + ALTER TABLE ts_kv_old + RENAME CONSTRAINT ts_kv_pkey TO ts_kv_pkey_old; + CREATE TABLE IF NOT EXISTS ts_kv + ( + LIKE ts_kv_old + ) + PARTITION BY RANGE (ts); + ALTER TABLE ts_kv + DROP COLUMN entity_type; + ALTER TABLE ts_kv + ALTER COLUMN entity_id TYPE uuid USING entity_id::uuid; + ALTER TABLE ts_kv + ALTER COLUMN key TYPE integer USING key::integer; + ALTER TABLE ts_kv + ADD CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_id, key, ts); + CREATE TABLE IF NOT EXISTS ts_kv_indefinite PARTITION OF ts_kv DEFAULT; END; $$; -- call create_new_ts_kv_latest_table(); -CREATE OR REPLACE PROCEDURE create_new_ts_kv_latest_table() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE create_new_ts_kv_latest_table() + LANGUAGE plpgsql AS +$$ BEGIN IF NOT EXISTS(SELECT FROM pg_tables WHERE schemaname = 'public' AND tablename = 'ts_kv_latest_old') THEN - ALTER TABLE ts_kv_latest + ALTER TABLE ts_kv_latest + DROP CONSTRAINT IF EXISTS ts_kv_latest_unq_key; + ALTER TABLE ts_kv_latest + DROP CONSTRAINT IF EXISTS ts_kv_latest_pkey; + ALTER TABLE ts_kv_latest + ADD CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_type, entity_id, key); + 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 - ( + 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 + ); + ALTER TABLE ts_kv_latest DROP COLUMN entity_type; - ALTER TABLE ts_kv_latest + ALTER TABLE ts_kv_latest ALTER COLUMN entity_id TYPE uuid USING entity_id::uuid; - ALTER TABLE ts_kv_latest + ALTER TABLE ts_kv_latest ALTER COLUMN key TYPE integer USING key::integer; - ALTER TABLE ts_kv_latest + ALTER TABLE ts_kv_latest ADD CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key); ELSE RAISE NOTICE 'ts_kv_latest_old table already exists!'; @@ -112,8 +128,9 @@ BEGIN RETURN QUERY SELECT SUBSTRING(year_date.year, 1, 4) AS partition_date, (extract(epoch from (year_date.year)::timestamp) * 1000)::bigint AS from_ts, (extract(epoch from (year_date.year::date + INTERVAL '1 YEAR')::timestamp) * - 1000)::bigint AS to_ts - FROM (SELECT DISTINCT TO_CHAR(TO_TIMESTAMP(ts / 1000), 'YYYY_01_01') AS year FROM ts_kv_old) AS year_date; + 1000)::bigint AS to_ts + FROM (SELECT DISTINCT TO_CHAR(TO_TIMESTAMP(ts / 1000), 'YYYY_01_01') AS year + FROM ts_kv_old) AS year_date; ELSE RAISE EXCEPTION 'Failed to parse partitioning property: % !', partition_type; END CASE; @@ -122,13 +139,16 @@ $$ LANGUAGE plpgsql; -- call create_partitions(); -CREATE OR REPLACE PROCEDURE create_partitions(IN partition_type varchar) LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE create_partitions(IN partition_type varchar) + LANGUAGE plpgsql AS +$$ DECLARE partition_date varchar; from_ts bigint; to_ts bigint; - partitions_cursor CURSOR FOR SELECT * FROM get_partitions_data(partition_type); + partitions_cursor CURSOR FOR SELECT * + FROM get_partitions_data(partition_type); BEGIN OPEN partitions_cursor; LOOP @@ -146,21 +166,25 @@ $$; -- call create_ts_kv_dictionary_table(); -CREATE OR REPLACE PROCEDURE create_ts_kv_dictionary_table() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE create_ts_kv_dictionary_table() + LANGUAGE plpgsql 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) - ); + 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; $$; -- call insert_into_dictionary(); -CREATE OR REPLACE PROCEDURE insert_into_dictionary() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE insert_into_dictionary() + LANGUAGE plpgsql AS +$$ DECLARE insert_record RECORD; @@ -191,9 +215,11 @@ BEGIN END; $$ LANGUAGE plpgsql; -CREATE OR REPLACE PROCEDURE insert_into_ts_kv(IN path_to_file varchar) LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE insert_into_ts_kv(IN path_to_file varchar) + LANGUAGE plpgsql AS +$$ BEGIN - EXECUTE format ('COPY (SELECT to_uuid(entity_id) AS entity_id, + EXECUTE format('COPY (SELECT to_uuid(entity_id) AS entity_id, ts_kv_records.key AS key, ts_kv_records.ts AS ts, ts_kv_records.bool_v AS bool_v, @@ -208,16 +234,19 @@ BEGIN long_v, dbl_v FROM ts_kv_old - INNER JOIN ts_kv_dictionary ON (ts_kv_old.key = ts_kv_dictionary.key)) AS ts_kv_records) TO %L;', path_to_file); - EXECUTE format ('COPY ts_kv FROM %L', path_to_file); + INNER JOIN ts_kv_dictionary ON (ts_kv_old.key = ts_kv_dictionary.key)) AS ts_kv_records) TO %L;', + path_to_file); + EXECUTE format('COPY ts_kv FROM %L', path_to_file); END $$; -- call insert_into_ts_kv_latest(); -CREATE OR REPLACE PROCEDURE insert_into_ts_kv_latest(IN path_to_file varchar) LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_latest(IN path_to_file varchar) + LANGUAGE plpgsql AS +$$ BEGIN - EXECUTE format ('COPY (SELECT to_uuid(entity_id) AS entity_id, + EXECUTE format('COPY (SELECT to_uuid(entity_id) AS entity_id, ts_kv_latest_records.key AS key, ts_kv_latest_records.ts AS ts, ts_kv_latest_records.bool_v AS bool_v, @@ -232,27 +261,30 @@ BEGIN 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 ts_kv_latest_records) TO %L;', path_to_file); - EXECUTE format ('COPY ts_kv_latest FROM %L', path_to_file); + INNER JOIN ts_kv_dictionary ON (ts_kv_latest_old.key = ts_kv_dictionary.key)) AS ts_kv_latest_records) TO %L;', + path_to_file); + EXECUTE format('COPY ts_kv_latest FROM %L', path_to_file); END; $$; -- call insert_into_ts_kv_cursor(); -CREATE OR REPLACE PROCEDURE insert_into_ts_kv_cursor() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_cursor() + LANGUAGE plpgsql AS +$$ DECLARE insert_size CONSTANT integer := 10000; insert_counter integer DEFAULT 0; insert_record RECORD; - insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, - ts_kv_records.key AS key, - ts_kv_records.ts AS ts, - ts_kv_records.bool_v AS bool_v, - ts_kv_records.str_v AS str_v, - ts_kv_records.long_v AS long_v, - ts_kv_records.dbl_v AS dbl_v - FROM (SELECT entity_id AS entity_id, - key_id AS key, + insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, + ts_kv_records.key AS key, + ts_kv_records.ts AS ts, + ts_kv_records.bool_v AS bool_v, + ts_kv_records.str_v AS str_v, + ts_kv_records.long_v AS long_v, + ts_kv_records.dbl_v AS dbl_v + FROM (SELECT entity_id AS entity_id, + key_id AS key, ts, bool_v, str_v, @@ -282,20 +314,22 @@ $$; -- call insert_into_ts_kv_latest_cursor(); -CREATE OR REPLACE PROCEDURE insert_into_ts_kv_latest_cursor() LANGUAGE plpgsql AS $$ +CREATE OR REPLACE PROCEDURE insert_into_ts_kv_latest_cursor() + LANGUAGE plpgsql AS +$$ DECLARE insert_size CONSTANT integer := 10000; insert_counter integer DEFAULT 0; insert_record RECORD; - insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, - ts_kv_latest_records.key AS key, - ts_kv_latest_records.ts AS ts, - ts_kv_latest_records.bool_v AS bool_v, - ts_kv_latest_records.str_v AS str_v, - ts_kv_latest_records.long_v AS long_v, - ts_kv_latest_records.dbl_v AS dbl_v - FROM (SELECT entity_id AS entity_id, - key_id AS key, + insert_cursor CURSOR FOR SELECT to_uuid(entity_id) AS entity_id, + ts_kv_latest_records.key AS key, + ts_kv_latest_records.ts AS ts, + ts_kv_latest_records.bool_v AS bool_v, + ts_kv_latest_records.str_v AS str_v, + ts_kv_latest_records.long_v AS long_v, + ts_kv_latest_records.dbl_v AS dbl_v + FROM (SELECT entity_id AS entity_id, + key_id AS key, ts, bool_v, str_v, diff --git a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java index 205973e6eb..7a8174af16 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/PsqlTsDatabaseUpgradeService.java @@ -103,8 +103,8 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe executeQuery(conn, CALL_CREATE_TS_KV_DICTIONARY_TABLE); executeQuery(conn, CALL_INSERT_INTO_DICTIONARY); - Path pathToTempTsKvFile; - Path pathToTempTsKvLatestFile; + Path pathToTempTsKvFile = null; + Path pathToTempTsKvLatestFile = null; if (SystemUtils.IS_OS_WINDOWS) { log.info("Lookup for environment variable: {} ...", THINGSBOARD_WINDOWS_UPGRADE_DIR); Path pathToDir; @@ -125,12 +125,11 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe try { copyTimeseries(conn, pathToTempTsKvFile, pathToTempTsKvLatestFile); } catch (Exception e) { - log.info("Upgrade script failed using the copy to/from files strategy!" + - " Trying to perfrom the upgrade using Inserts strategy ..."); insertTimeseries(conn); } } catch (IOException | SecurityException e) { - throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); + log.warn("Failed to create time-series upgrade files due to: {}", e.getMessage()); + insertTimeseries(conn); } } else { try { @@ -148,13 +147,11 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); } } catch (Exception e) { - log.info(e.getMessage()); - log.info("Upgrade script failed using the copy to/from files strategy!" + - " Trying to perfrom the upgrade using Inserts strategy ..."); insertTimeseries(conn); } } catch (IOException | SecurityException e) { - throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); + log.warn("Failed to create time-series upgrade files due to: {}", e.getMessage()); + insertTimeseries(conn); } } @@ -204,11 +201,13 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe } private void removeUpgradeFiles(Path pathToTempTsKvFile, Path pathToTempTsKvLatestFile) { - if (pathToTempTsKvFile.toFile().exists() && pathToTempTsKvLatestFile.toFile().exists()) { + if (pathToTempTsKvFile != null && pathToTempTsKvFile.toFile().exists()) { boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); if (deleteTsKvFile) { log.info("Successfully deleted the temp file for ts_kv table upgrade!"); } + } + if (pathToTempTsKvLatestFile != null && pathToTempTsKvLatestFile.toFile().exists()) { boolean deleteTsKvLatestFile = pathToTempTsKvLatestFile.toFile().delete(); if (deleteTsKvLatestFile) { log.info("Successfully deleted the temp file for ts_kv_latest table upgrade!"); @@ -223,6 +222,8 @@ public class PsqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSe } private void insertTimeseries(Connection conn) { + log.warn("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); executeQuery(conn, CALL_CREATE_NEW_TS_KV_LATEST_TABLE); executeQuery(conn, CALL_INSERT_INTO_TS_KV_LATEST_CURSOR); diff --git a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java index ada1f7ad42..d8f7ea61f9 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/TimescaleTsDatabaseUpgradeService.java @@ -99,7 +99,7 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr executeQuery(conn, CALL_CREATE_TS_KV_DICTIONARY_TABLE); executeQuery(conn, CALL_INSERT_INTO_DICTIONARY); - Path pathToTempTsKvFile; + Path pathToTempTsKvFile = null; if (SystemUtils.IS_OS_WINDOWS) { Path pathToDir; log.info("Lookup for environment variable: {} ...", THINGSBOARD_WINDOWS_UPGRADE_DIR); @@ -118,12 +118,11 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr try { executeQuery(conn, "call insert_into_ts_kv('" + pathToTempTsKvFile + "')"); } catch (Exception e) { - log.info("Upgrade script failed using the copy to/from files strategy!" + - " Trying to perfrom the upgrade using Inserts strategy ..."); - executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + insertTimeseries(conn); } } catch (IOException | SecurityException e) { - throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); + log.warn("Failed to create time-series upgrade files due to: {}", e.getMessage()); + insertTimeseries(conn); } } else { try { @@ -140,13 +139,11 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!"); } } catch (Exception e) { - log.info(e.getMessage()); - log.info("Upgrade script failed using the copy to/from files strategy!" + - " Trying to perfrom the upgrade using Inserts strategy ..."); - executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + insertTimeseries(conn); } } catch (IOException | SecurityException e) { - throw new RuntimeException("Failed to create time-series upgrade files due to: " + e); + log.warn("Failed to create time-series upgrade files due to: {}", e.getMessage()); + insertTimeseries(conn); } } removeUpgradeFile(pathToTempTsKvFile); @@ -185,8 +182,14 @@ public class TimescaleTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgr } } + private void insertTimeseries(Connection conn) { + log.warn("Upgrade script failed using the copy to/from files strategy!" + + " Trying to perfrom the upgrade using Inserts strategy ..."); + executeQuery(conn, CALL_INSERT_INTO_TS_KV_CURSOR); + } + private void removeUpgradeFile(Path pathToTempTsKvFile) { - if (pathToTempTsKvFile.toFile().exists()) { + if (pathToTempTsKvFile != null && pathToTempTsKvFile.toFile().exists()) { boolean deleteTsKvFile = pathToTempTsKvFile.toFile().delete(); if (deleteTsKvFile) { log.info("Successfully deleted the temp file for ts_kv table upgrade!");