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!");