From 3f10ba464e96106f77688c641829e7874c146a29 Mon Sep 17 00:00:00 2001 From: Shvaika Dmytro Date: Tue, 12 Sep 2023 12:16:31 +0300 Subject: [PATCH] added executor to speed up rule nodes upgrade & added SQL upgrade for TbMsgAttributesNode & updated logback.xml with new debug entry to view the rule node upgrade logs (#9224) * added executor to speed up rule nodes upgrade & added SQL upgrade for TbMsgAttributesNode & updated logback.xml with new debug entry to view the rule node upgrade logs * refactoring after review * create index before upsert * added missing logic to await futures to be completed with max pending futures size set to 100 * added better logging * refacotring after review --- .../install/SqlDatabaseUpgradeService.java | 7 +++ .../update/DefaultDataUpdateService.java | 44 ++++++++++++++----- application/src/main/resources/logback.xml | 2 + 3 files changed, 43 insertions(+), 10 deletions(-) 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 @@ + +