Browse Source

Rollback upgrade if schema update failed, refactoring

pull/11673/head
ViacheslavKlimov 2 years ago
parent
commit
6a9cbaf386
  1. 260
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java

260
application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java

@ -16,24 +16,22 @@
package org.thingsboard.server.service.install; package org.thingsboard.server.service.install;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.intellij.lang.annotations.Language;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Profile; import org.springframework.context.annotation.Profile;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.StatementCallback;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.service.install.update.DefaultDataUpdateService; import org.thingsboard.server.service.install.update.DefaultDataUpdateService;
import java.nio.charset.Charset; import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files; import java.nio.file.Files;
import java.nio.file.Path; import java.nio.file.Path;
import java.nio.file.Paths; import java.nio.file.Paths;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.SQLWarning; import java.sql.SQLWarning;
import java.sql.Statement;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
@Service @Service
@Profile("install") @Profile("install")
@ -42,140 +40,126 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
private static final String SCHEMA_UPDATE_SQL = "schema_update.sql"; private static final String SCHEMA_UPDATE_SQL = "schema_update.sql";
@Value("${spring.datasource.url}") private final InstallScripts installScripts;
private String dbUrl; private final JdbcTemplate jdbcTemplate;
private final TransactionTemplate transactionTemplate;
@Value("${spring.datasource.username}") public SqlDatabaseUpgradeService(InstallScripts installScripts, JdbcTemplate jdbcTemplate, PlatformTransactionManager transactionManager) {
private String dbUserName; this.installScripts = installScripts;
this.jdbcTemplate = jdbcTemplate;
@Value("${spring.datasource.password}") this.transactionTemplate = new TransactionTemplate(transactionManager);
private String dbPassword; this.transactionTemplate.setTimeout((int) TimeUnit.MINUTES.toSeconds(120));
}
@Autowired
private InstallScripts installScripts;
@Override @Override
public void upgradeDatabase(String fromVersion) throws Exception { public void upgradeDatabase(String fromVersion) {
switch (fromVersion) { switch (fromVersion) {
case "3.5.0": case "3.5.0" -> updateSchema("3.5.0", 3005000, "3.5.1", 3005001);
updateSchema("3.5.0", 3005000, "3.5.1", 3005001, null); case "3.5.1" -> {
break; updateSchema("3.5.1", 3005001, "3.6.0", 3006000);
case "3.5.1":
updateSchema("3.5.1", 3005001, "3.6.0", 3006000, conn -> { String[] tables = new String[]{"device", "component_descriptor", "customer", "dashboard", "rule_chain", "rule_node", "ota_package",
String[] entityNames = new String[]{"device", "component_descriptor", "customer", "dashboard", "rule_chain", "rule_node", "ota_package", "asset_profile", "asset", "device_profile", "tb_user", "tenant_profile", "tenant", "widgets_bundle", "entity_view", "edge"};
"asset_profile", "asset", "device_profile", "tb_user", "tenant_profile", "tenant", "widgets_bundle", "entity_view", "edge"}; for (String table : tables) {
for (String entityName : entityNames) { execute("ALTER TABLE " + table + " DROP COLUMN IF EXISTS search_text CASCADE");
try { }
conn.createStatement().execute("ALTER TABLE " + entityName + " DROP COLUMN search_text CASCADE"); execute(
} catch (Exception e) { "ALTER TABLE component_descriptor ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;",
} "ALTER TABLE rule_node ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;",
} "CREATE INDEX IF NOT EXISTS idx_rule_node_type_configuration_version ON rule_node(type, configuration_version);",
try { "UPDATE rule_node SET " +
conn.createStatement().execute("ALTER TABLE component_descriptor ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;");
} catch (Exception e) {
}
try {
conn.createStatement().execute("ALTER TABLE rule_node ADD COLUMN IF NOT EXISTS configuration_version int DEFAULT 0;");
} catch (Exception e) {
}
try {
conn.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_rule_node_type_configuration_version ON rule_node(type, configuration_version);");
} catch (Exception e) {
}
try {
conn.createStatement().execute("UPDATE rule_node SET " +
"configuration = (configuration::jsonb || '{\"updateAttributesOnlyOnValueChange\": \"false\"}'::jsonb)::varchar, " + "configuration = (configuration::jsonb || '{\"updateAttributesOnlyOnValueChange\": \"false\"}'::jsonb)::varchar, " +
"configuration_version = 1 " + "configuration_version = 1 " +
"WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode' AND configuration_version < 1;"); "WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode' AND configuration_version < 1;",
} catch (Exception e) { "CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';"
} );
try { }
conn.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';"); case "3.6.0" -> updateSchema("3.6.0", 3006000, "3.6.1", 3006001);
} catch (Exception e) { case "3.6.1" -> {
} updateSchema("3.6.1", 3006001, "3.6.2", 3006002);
});
break; try {
case "3.6.0": Path saveAttributesNodeUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.6.1", "save_attributes_node_update.sql");
updateSchema("3.6.0", 3006000, "3.6.1", 3006001, null); loadSql(saveAttributesNodeUpdateFile);
break; } catch (Exception e) {
case "3.6.1": log.warn("Failed to execute update script for save attributes rule nodes due to: ", e);
updateSchema("3.6.1", 3006001, "3.6.2", 3006002, connection -> { }
try { execute("CREATE INDEX IF NOT EXISTS idx_asset_profile_id ON asset(tenant_id, asset_profile_id);");
Path saveAttributesNodeUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.6.1", "save_attributes_node_update.sql"); }
loadSql(saveAttributesNodeUpdateFile, connection); case "3.6.2" -> updateSchema("3.6.2", 3006002, "3.6.3", 3006003);
} catch (Exception e) { case "3.6.3" -> updateSchema("3.6.3", 3006003, "3.6.4", 3006004);
log.warn("Failed to execute update script for save attributes rule nodes due to: ", e); case "3.6.4" -> updateSchema("3.6.4", 3006004, "3.7.0", 3007000);
} case "3.7.0" -> {
try { updateSchema("3.7.0", 3007000, "3.7.1", 3007001);
connection.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_asset_profile_id ON asset(tenant_id, asset_profile_id);");
} catch (Exception e) { try {
} execute("UPDATE rule_node SET " +
}); "configuration = CASE " +
break; " WHEN (configuration::jsonb ->> 'persistAlarmRulesState') = 'false'" +
case "3.6.2": " THEN (configuration::jsonb || '{\"fetchAlarmRulesStateOnStart\": \"false\"}'::jsonb)::varchar " +
updateSchema("3.6.2", 3006002, "3.6.3", 3006003, null); " ELSE configuration " +
break; "END, " +
case "3.6.3": "configuration_version = 1 " +
updateSchema("3.6.3", 3006003, "3.6.4", 3006004, null); "WHERE type = 'org.thingsboard.rule.engine.profile.TbDeviceProfileNode' " +
break; "AND configuration_version < 1;", false);
case "3.6.4": } catch (Exception e) {
updateSchema("3.6.4", 3006004, "3.7.0", 3007000, null); log.warn("Failed to execute update script for device profile rule nodes due to: ", e);
break; }
case "3.7.0": }
updateSchema("3.7.0", 3007000, "3.7.1", 3007001, connection -> { case "3.7.1" -> updateSchema("3.7.1", 3007001, "3.7.2", 3007002);
try { default -> throw new RuntimeException("Unsupported fromVersion '" + fromVersion + "'");
connection.createStatement().execute("UPDATE rule_node SET " +
"configuration = CASE " +
" WHEN (configuration::jsonb ->> 'persistAlarmRulesState') = 'false'" +
" THEN (configuration::jsonb || '{\"fetchAlarmRulesStateOnStart\": \"false\"}'::jsonb)::varchar " +
" ELSE configuration " +
"END, " +
"configuration_version = 1 " +
"WHERE type = 'org.thingsboard.rule.engine.profile.TbDeviceProfileNode' " +
"AND configuration_version < 1;");
} catch (Exception e) {
log.warn("Failed to execute update script for device profile rule nodes due to: ", e);
}
});
break;
case "3.7.1":
updateSchema("3.7.1", 3007001, "3.7.2", 3007002, null);
break;
default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
} }
} }
private void updateSchema(String oldVersionStr, int oldVersion, String newVersionStr, int newVersion, Consumer<Connection> additionalAction) { private void updateSchema(String oldVersionStr, int oldVersion, String newVersionStr, int newVersion) {
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { try {
log.info("Updating schema ..."); transactionTemplate.executeWithoutResult(ts -> {
if (isOldSchema(conn, oldVersion)) { log.info("Updating schema ...");
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", oldVersionStr, SCHEMA_UPDATE_SQL); if (isOldSchema(oldVersion)) {
loadSql(schemaUpdateFile, conn); Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", oldVersionStr, SCHEMA_UPDATE_SQL);
if (additionalAction != null) { loadSql(schemaUpdateFile);
additionalAction.accept(conn); jdbcTemplate.execute("UPDATE tb_schema_settings SET schema_version = " + newVersion);
log.info("Schema updated to version {}", newVersionStr);
} else {
log.info("Skip schema re-update to version {}. Use env flag 'SKIP_SCHEMA_VERSION_CHECK' to force the re-update.", newVersionStr);
} }
conn.createStatement().execute("UPDATE tb_schema_settings SET schema_version = " + newVersion + ";"); });
log.info("Schema updated to version {}", newVersionStr);
} else {
log.info("Skip schema re-update to version {}. Use env flag 'SKIP_SCHEMA_VERSION_CHECK' to force the re-update.", newVersionStr);
}
} catch (Exception e) { } catch (Exception e) {
log.error("Failed updating schema!!!", e); throw new RuntimeException("Failed to update schema", e);
}
}
private void loadSql(Path sqlFile) {
String sql;
try {
sql = Files.readString(sqlFile);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
jdbcTemplate.execute((StatementCallback<Object>) stmt -> {
stmt.execute(sql);
printWarnings(stmt.getWarnings());
return null;
});
}
private void execute(@Language("sql") String... statements) {
for (String statement : statements) {
execute(statement, true);
} }
} }
private void loadSql(Path sqlFile, Connection conn) throws Exception { private void execute(@Language("sql") String statement, boolean ignoreErrors) {
String sql = new String(Files.readAllBytes(sqlFile), Charset.forName("UTF-8")); try {
Statement st = conn.createStatement(); jdbcTemplate.execute(statement);
st.setQueryTimeout((int) TimeUnit.HOURS.toSeconds(3)); } catch (Exception e) {
st.execute(sql);//NOSONAR, ignoring because method used to execute thingsboard database upgrade script if (!ignoreErrors) {
printWarnings(st); throw e;
Thread.sleep(5000); }
}
} }
protected void printWarnings(Statement statement) throws SQLException { private void printWarnings(SQLWarning warnings) {
SQLWarning warnings = statement.getWarnings();
if (warnings != null) { if (warnings != null) {
log.info("{}", warnings.getMessage()); log.info("{}", warnings.getMessage());
SQLWarning nextWarning = warnings.getNextWarning(); SQLWarning nextWarning = warnings.getNextWarning();
@ -186,26 +170,18 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
} }
} }
protected boolean isOldSchema(Connection conn, long fromVersion) { private boolean isOldSchema(long fromVersion) {
if (DefaultDataUpdateService.getEnv("SKIP_SCHEMA_VERSION_CHECK", false)) { if (DefaultDataUpdateService.getEnv("SKIP_SCHEMA_VERSION_CHECK", false)) {
log.info("Skipped DB schema version check due to SKIP_SCHEMA_VERSION_CHECK set to true!"); log.info("Skipped DB schema version check due to SKIP_SCHEMA_VERSION_CHECK set to true!");
return true; return true;
} }
jdbcTemplate.execute("CREATE TABLE IF NOT EXISTS tb_schema_settings (schema_version bigint NOT NULL, CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version))");
Long schemaVersion = jdbcTemplate.queryForList("SELECT schema_version FROM tb_schema_settings", Long.class).stream().findFirst().orElse(null);
boolean isOldSchema = true; boolean isOldSchema = true;
try { if (schemaVersion != null) {
Statement statement = conn.createStatement(); isOldSchema = schemaVersion <= fromVersion;
statement.execute("CREATE TABLE IF NOT EXISTS tb_schema_settings ( schema_version bigint NOT NULL, CONSTRAINT tb_schema_settings_pkey PRIMARY KEY (schema_version));"); } else {
Thread.sleep(1000); jdbcTemplate.execute("INSERT INTO tb_schema_settings (schema_version) VALUES (" + fromVersion + ")");
ResultSet resultSet = statement.executeQuery("SELECT schema_version FROM tb_schema_settings;");
if (resultSet.next()) {
isOldSchema = resultSet.getLong(1) <= fromVersion;
} else {
resultSet.close();
statement.execute("INSERT INTO tb_schema_settings (schema_version) VALUES (" + fromVersion + ")");
}
statement.close();
} catch (InterruptedException | SQLException e) {
log.info("Failed to check current PostgreSQL schema due to: {}", e.getMessage());
} }
return isOldSchema; return isOldSchema;
} }

Loading…
Cancel
Save