diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index 4bb29426c9..9e942b481e 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -756,6 +756,13 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService 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_version = 1 " + + "WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode' AND configuration_version < 1;"); + } catch (Exception e) { + } try { conn.createStatement().execute("CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';"); } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index a1048d6a37..8c3b78040b 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -78,6 +78,7 @@ import org.thingsboard.server.dao.model.sql.DeviceProfileEntity; import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; +import org.thingsboard.server.dao.sql.JpaExecutorService; import org.thingsboard.server.dao.sql.device.DeviceProfileRepository; import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.dao.tenant.TenantService; @@ -101,6 +102,8 @@ import static org.thingsboard.server.common.data.StringUtils.isBlank; @Slf4j public class DefaultDataUpdateService implements DataUpdateService { + private static final int MAX_PENDING_SAVE_RULE_NODE_FUTURES = 100; + @Autowired private TenantService tenantService; @@ -153,6 +156,9 @@ public class DefaultDataUpdateService implements DataUpdateService { @Autowired private EdgeEventDao edgeEventDao; + @Autowired + JpaExecutorService jpaExecutorService; + @Override public void updateData(String fromVersion) throws Exception { switch (fromVersion) { @@ -225,13 +231,15 @@ public class DefaultDataUpdateService implements DataUpdateService { @Override public void upgradeRuleNodes() { try { + var futures = new ArrayList>(100); + int totalRuleNodesUpgraded = 0; log.info("Starting rule nodes upgrade ..."); var nodeClassToVersionMap = componentDiscoveryService.getVersionedNodes(); log.debug("Found {} versioned nodes to check for upgrade!", nodeClassToVersionMap.size()); - nodeClassToVersionMap.forEach(clazz -> { - var ruleNodeType = clazz.getClassName(); - var ruleNodeTypeForLogs = clazz.getSimpleName(); - var toVersion = clazz.getCurrentVersion(); + for (var ruleNodeClassInfo : nodeClassToVersionMap) { + var ruleNodeType = ruleNodeClassInfo.getClassName(); + var ruleNodeTypeForLogs = ruleNodeClassInfo.getSimpleName(); + var toVersion = ruleNodeClassInfo.getCurrentVersion(); log.debug("Going to check for nodes with type: {} to upgrade to version: {}.", ruleNodeTypeForLogs, toVersion); var ruleNodesToUpdate = new PageDataIterable<>( pageLink -> ruleChainService.findAllRuleNodesByTypeAndVersionLessThan(ruleNodeType, toVersion, pageLink), 1024 @@ -246,28 +254,44 @@ public class DefaultDataUpdateService implements DataUpdateService { log.debug("Going to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion); try { - var tbVersionedNode = (TbNode) clazz.getClazz().getDeclaredConstructor().newInstance(); + var tbVersionedNode = (TbNode) ruleNodeClassInfo.getClazz().getDeclaredConstructor().newInstance(); TbPair upgradeRuleNodeConfigurationResult = tbVersionedNode.upgrade(fromVersion, oldConfiguration); if (upgradeRuleNodeConfigurationResult.getFirst()) { ruleNode.setConfiguration(upgradeRuleNodeConfigurationResult.getSecond()); } ruleNode.setConfigurationVersion(toVersion); - ruleChainService.saveRuleNode(TenantId.SYS_TENANT_ID, ruleNode); - log.debug("Successfully upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", - ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion); + futures.add(jpaExecutorService.submit(() -> { + ruleChainService.saveRuleNode(TenantId.SYS_TENANT_ID, ruleNode); + log.debug("Successfully upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", + ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion); + })); + if (futures.size() >= MAX_PENDING_SAVE_RULE_NODE_FUTURES) { + log.info("{} upgraded rule nodes so far ...", + totalRuleNodesUpgraded += awaitFuturesToCompleteAndGetCount(futures)); + futures.clear(); + } } catch (Exception e) { log.warn("Failed to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {} due to: ", ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion, e); } } } - }); - log.info("Finished rule nodes upgrade!"); + } + log.info("Finished rule nodes upgrade. Upgraded rule nodes count: {}", + totalRuleNodesUpgraded + awaitFuturesToCompleteAndGetCount(futures)); } catch (Exception e) { log.error("Unexpected error during rule nodes upgrade: ", e); } } + private int awaitFuturesToCompleteAndGetCount(List> futures) { + try { + return Futures.allAsList(futures).get().size(); + } catch (ExecutionException | InterruptedException e) { + throw new RuntimeException("Failed to process save rule nodes requests due to: ", e); + } + } + private final PaginatedUpdater deviceProfileEntityDynamicConditionsUpdater = new PaginatedUpdater<>() { diff --git a/application/src/main/resources/logback.xml b/application/src/main/resources/logback.xml index eaf6ace29e..b0a4174012 100644 --- a/application/src/main/resources/logback.xml +++ b/application/src/main/resources/logback.xml @@ -30,6 +30,8 @@ + +