From b0bfdfff8e8e3826fefa145856e8211892d5a35b Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 25 May 2021 10:38:12 +0300 Subject: [PATCH] Pagedata usage for get edges methods --- .../server/controller/BaseController.java | 22 +++- .../controller/RuleChainController.java | 4 +- .../rpc/processor/AlarmEdgeProcessor.java | 40 +++--- .../edge/rpc/processor/BaseEdgeProcessor.java | 4 +- .../rpc/processor/CustomerEdgeProcessor.java | 2 +- .../rpc/processor/DeviceEdgeProcessor.java | 60 ++++----- .../edge/rpc/processor/EdgeProcessor.java | 2 +- .../rpc/processor/EntityEdgeProcessor.java | 83 ++++++------- .../rpc/processor/RelationEdgeProcessor.java | 62 +++++----- .../rpc/sync/DefaultEdgeRequestsService.java | 4 +- .../server/dao/edge/EdgeService.java | 6 +- .../thingsboard/server/dao/edge/EdgeDao.java | 18 +-- .../server/dao/edge/EdgeServiceImpl.java | 117 ++++++------------ .../dao/entity/AbstractEntityService.java | 4 +- .../server/dao/rule/BaseRuleChainService.java | 22 ++-- .../server/dao/sql/edge/EdgeRepository.java | 11 ++ .../server/dao/sql/edge/JpaEdgeDao.java | 23 ++-- .../rule/engine/edge/TbMsgPushToEdgeNode.java | 105 +++++++--------- 18 files changed, 257 insertions(+), 332 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 9f6379cf6c..0543fdd47f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -83,6 +83,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; @@ -145,6 +146,7 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.mail.MessagingException; import javax.servlet.http.HttpServletResponse; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Optional; @@ -164,6 +166,8 @@ public abstract class BaseController { protected static final String DEFAULT_DASHBOARD = "defaultDashboardId"; protected static final String HOME_DASHBOARD = "homeDashboardId"; + private static final int DEFAULT_PAGE_SIZE = 1000; + private static final ObjectMapper json = new ObjectMapper(); @Autowired @@ -1082,12 +1086,18 @@ public abstract class BaseController { if (!edgesEnabled) { return null; } - List result = null; - try { - result = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId).get(); - } catch (Exception e) { - log.error("[{}] can't find related edge ids for entity [{}]", tenantId, entityId, e); - } + List result = new ArrayList<>(); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + result.addAll(pageData.getData()); + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } + } + } while (pageData != null && pageData.hasNext()); return result; } diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index d3035b4f3b..1f09e29e6e 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -86,7 +86,7 @@ public class RuleChainController extends BaseController { public static final String RULE_CHAIN_ID = "ruleChainId"; public static final String RULE_NODE_ID = "ruleNodeId"; - private static final int DEFAULT_LIMIT = 100; + private static final int DEFAULT_PAGE_SIZE = 1000; private static final ObjectMapper objectMapper = new ObjectMapper(); @@ -643,7 +643,7 @@ public class RuleChainController extends BaseController { try { TenantId tenantId = getCurrentUser().getTenantId(); List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = ruleChainService.findAutoAssignToEdgeRuleChainsByTenantId(tenantId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java index 3ff245b0d9..5893840728 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java @@ -34,6 +34,8 @@ import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.DownlinkMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; @@ -41,7 +43,6 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; import java.util.Collections; -import java.util.List; import java.util.UUID; @Component @@ -138,34 +139,31 @@ public class AlarmEdgeProcessor extends BaseEdgeProcessor { if (alarm != null) { EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType()); if (type != null) { - ListenableFuture> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); - Futures.addCallback(relatedEdgeIdsByEntityIdFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List relatedEdgeIdsByEntityId) { - if (relatedEdgeIdsByEntityId != null) { - for (EdgeId edgeId : relatedEdgeIdsByEntityId) { - saveEdgeEvent(tenantId, - edgeId, - EdgeEventType.ALARM, - EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), - alarmId, - null); - } + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (EdgeId edgeId : pageData.getData()) { + saveEdgeEvent(tenantId, + edgeId, + EdgeEventType.ALARM, + EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), + alarmId, + null); + } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); } } - - @Override - public void onFailure(Throwable t) { - log.warn("[{}] can't find related edge ids by entity id [{}]", tenantId.getId(), alarm.getOriginator(), t); - } - }, dbCallbackExecutorService); + } while (pageData != null && pageData.hasNext()); } } } @Override public void onFailure(Throwable t) { - log.warn("[{}] can't find alarm by id [{}]", tenantId.getId(), alarmId.getId(), t); + log.warn("[{}] can't find alarm by id [{}] {}", tenantId.getId(), alarmId.getId(), t); } }, dbCallbackExecutorService); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index e982cca088..f32576499e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -74,7 +74,7 @@ public abstract class BaseEdgeProcessor { protected static final ObjectMapper mapper = new ObjectMapper(); - protected static final int DEFAULT_LIMIT = 100; + protected static final int DEFAULT_PAGE_SIZE = 1000; @Autowired protected RuleChainService ruleChainService; @@ -221,7 +221,7 @@ public abstract class BaseEdgeProcessor { } protected void processActionForAllEdges(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) { - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = edgeService.findEdgesByTenantId(tenantId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java index 63b2c5d577..00de4eb576 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java @@ -75,7 +75,7 @@ public class CustomerEdgeProcessor extends BaseEdgeProcessor { CustomerId customerId = new CustomerId(EntityIdFactory.getByEdgeEventTypeAndUuid(type, uuid).getId()); switch (actionType) { case UPDATED: - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = edgeService.findEdgesByTenantIdAndCustomerId(tenantId, customerId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index c889b3a49c..f540186976 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -17,14 +17,12 @@ package org.thingsboard.server.service.edge.rpc.processor; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.node.ObjectNode; -import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; import org.apache.commons.lang3.StringUtils; -import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.stereotype.Component; import org.thingsboard.rule.engine.api.RpcError; import org.thingsboard.server.common.data.Customer; @@ -40,6 +38,8 @@ import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.security.DeviceCredentials; @@ -61,7 +61,6 @@ import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; import java.util.Collections; -import java.util.List; import java.util.UUID; import java.util.concurrent.locks.ReentrantLock; @@ -79,41 +78,34 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { String deviceName = deviceUpdateMsg.getName(); Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); if (device != null) { - ListenableFuture> future = edgeService.findRelatedEdgeIdsByEntityId(tenantId, device.getId()); - SettableFuture futureToSet = SettableFuture.create(); - Futures.addCallback(future, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List edgeIds) { - boolean update = false; - if (edgeIds != null && !edgeIds.isEmpty()) { - if (edgeIds.contains(edge.getId())) { - update = true; - } + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, device.getId(), pageLink); + boolean update = false; + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + if (pageData.getData().contains(edge.getId())) { + update = true; } - Device device; - if (update) { - log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " + - "deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg); - updateDevice(tenantId, edge, deviceUpdateMsg); - } else { - log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." + - "Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg); - String newDeviceName = deviceUpdateMsg.getName() + "_" + RandomStringUtils.randomAlphabetic(15); - device = createDevice(tenantId, edge, deviceUpdateMsg, newDeviceName); - ObjectNode body = mapper.createObjectNode(); - body.put("conflictName", deviceName); - saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, device.getId(), body); + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); } - futureToSet.set(null); } - - @Override - public void onFailure(Throwable t) { - log.error("[{}] Failed to get related edge ids by device id [{}], edge [{}]", tenantId, deviceUpdateMsg, edge.getId(), t); - futureToSet.setException(t); + Device newDevice; + if (update) { + log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " + + "deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg); + updateDevice(tenantId, edge, deviceUpdateMsg); + } else { + log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." + + "Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg); + String newDeviceName = deviceUpdateMsg.getName() + "_" + RandomStringUtils.randomAlphabetic(15); + newDevice = createDevice(tenantId, edge, deviceUpdateMsg, newDeviceName); + ObjectNode body = mapper.createObjectNode(); + body.put("conflictName", deviceName); + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, newDevice.getId(), body); } - }, dbCallbackExecutorService); - return futureToSet; + } while (pageData != null && pageData.hasNext()); } else { log.info("[{}] Creating new device and replacing device entity on the edge [{}]", tenantId, deviceUpdateMsg); device = createDevice(tenantId, edge, deviceUpdateMsg, deviceUpdateMsg.getName()); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java index c9cd0f4019..c225fb9c60 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java @@ -54,7 +54,7 @@ public class EdgeProcessor extends BaseEdgeProcessor { public void onSuccess(@Nullable Edge edge) { if (edge != null && !customerId.isNullUid()) { saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, customerId, null); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = userService.findCustomerUsers(tenantId, customerId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java index 5bf2e77044..feb50ce750 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java @@ -93,63 +93,56 @@ public class EntityEdgeProcessor extends BaseEdgeProcessor { EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); - ListenableFuture> edgeIdsFuture; + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; switch (actionType) { case ADDED: // used only for USER entity case UPDATED: case CREDENTIALS_UPDATED: - edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId); - Futures.addCallback(edgeIdsFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List edgeIds) { - if (edgeIds != null && !edgeIds.isEmpty()) { - for (EdgeId edgeId : edgeIds) { - saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); - } + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (EdgeId relatedEdgeId : pageData.getData()) { + saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null); + } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); } } - @Override - public void onFailure(Throwable throwable) { - log.error("Failed to find related edge ids [{}]", edgeNotificationMsg, throwable); - } - }, dbCallbackExecutorService); + } while (pageData != null && pageData.hasNext()); break; case ASSIGNED_TO_CUSTOMER: case UNASSIGNED_FROM_CUSTOMER: - edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId); - Futures.addCallback(edgeIdsFuture, new FutureCallback<>() { - @Override - public void onSuccess(@Nullable List edgeIds) { - if (edgeIds != null && !edgeIds.isEmpty()) { - for (EdgeId edgeId : edgeIds) { - try { - CustomerId customerId = mapper.readValue(edgeNotificationMsg.getBody(), CustomerId.class); - ListenableFuture future = edgeService.findEdgeByIdAsync(tenantId, edgeId); - Futures.addCallback(future, new FutureCallback() { - @Override - public void onSuccess(@Nullable Edge edge) { - if (edge != null && edge.getCustomerId() != null && - !edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(customerId)) { - saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); - } + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (EdgeId relatedEdgeId : pageData.getData()) { + try { + CustomerId customerId = mapper.readValue(edgeNotificationMsg.getBody(), CustomerId.class); + ListenableFuture future = edgeService.findEdgeByIdAsync(tenantId, relatedEdgeId); + Futures.addCallback(future, new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Edge edge) { + if (edge != null && edge.getCustomerId() != null && + !edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(customerId)) { + saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null); } - @Override - public void onFailure(Throwable throwable) { - log.error("Failed to find edge by id [{}]", edgeNotificationMsg, throwable); - } - }, dbCallbackExecutorService); - } catch (Exception e) { - log.error("Can't parse customer id from entity body [{}]", edgeNotificationMsg, e); - } + } + + @Override + public void onFailure(Throwable t) { + log.error("Failed to find edge by id [{}] {}", edgeNotificationMsg, t); + } + }, dbCallbackExecutorService); + } catch (Exception e) { + log.error("Can't parse customer id from entity body [{}]", edgeNotificationMsg, e); } } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } } - - @Override - public void onFailure(Throwable throwable) { - log.error("Failed to find related edge ids [{}]", edgeNotificationMsg, throwable); - } - }, dbCallbackExecutorService); + } while (pageData != null && pageData.hasNext()); break; case DELETED: saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); @@ -165,7 +158,7 @@ public class EntityEdgeProcessor extends BaseEdgeProcessor { } private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java index 190d58fcef..ddaf55a73b 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java @@ -16,11 +16,9 @@ package org.thingsboard.server.service.edge.rpc.processor; import com.fasterxml.jackson.core.JsonProcessingException; -import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.EdgeEvent; @@ -38,6 +36,8 @@ import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.gen.edge.DownlinkMsg; @@ -127,39 +127,35 @@ public class RelationEdgeProcessor extends BaseEdgeProcessor { EntityRelation relation = mapper.readValue(edgeNotificationMsg.getBody(), EntityRelation.class); if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && !relation.getTo().getEntityType().equals(EntityType.EDGE)) { - List>> futures = new ArrayList<>(); - futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo())); - futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom())); - ListenableFuture>> combinedFuture = Futures.allAsList(futures); - Futures.addCallback(combinedFuture, new FutureCallback>>() { - @Override - public void onSuccess(@Nullable List> listOfListsEdgeIds) { - Set uniqueEdgeIds = new HashSet<>(); - if (listOfListsEdgeIds != null && !listOfListsEdgeIds.isEmpty()) { - for (List listOfListsEdgeId : listOfListsEdgeIds) { - if (listOfListsEdgeId != null) { - uniqueEdgeIds.addAll(listOfListsEdgeId); - } - } - } - if (!uniqueEdgeIds.isEmpty()) { - for (EdgeId edgeId : uniqueEdgeIds) { - saveEdgeEvent(tenantId, - edgeId, - EdgeEventType.RELATION, - EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), - null, - mapper.valueToTree(relation)); - } - } + Set uniqueEdgeIds = new HashSet<>(); + uniqueEdgeIds.addAll(findRelatedEdgeIds(tenantId, relation.getTo())); + uniqueEdgeIds.addAll(findRelatedEdgeIds(tenantId, relation.getFrom())); + if (!uniqueEdgeIds.isEmpty()) { + for (EdgeId edgeId : uniqueEdgeIds) { + saveEdgeEvent(tenantId, + edgeId, + EdgeEventType.RELATION, + EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), + null, + mapper.valueToTree(relation)); } + } + } + } - @Override - public void onFailure(Throwable t) { - log.warn("[{}] can't find related edge ids by relation to id [{}] and relation from id [{}]" , - tenantId.getId(), relation.getTo().getId(), relation.getFrom().getId(), t); + private List findRelatedEdgeIds(TenantId tenantId, EntityId entityId) { + List result = new ArrayList<>(); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + result.addAll(pageData.getData()); + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); } - }, dbCallbackExecutorService); - } + } + } while (pageData != null && pageData.hasNext()); + return result; } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java index bccc6cedd0..bad62b0945 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java @@ -87,7 +87,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { private static final ObjectMapper mapper = new ObjectMapper(); - private static final int DEFAULT_LIMIT = 100; + private static final int DEFAULT_PAGE_SIZE = 1000; @Autowired private EdgeEventService edgeEventService; @@ -351,7 +351,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { List> futures = new ArrayList<>(); log.trace("[{}] syncDevices [{}][{}]", tenantId, edge.getName(), deviceType); try { - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = deviceService.findDevicesByTenantIdAndEdgeIdAndType(tenantId, edge.getId(), deviceType, pageLink); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index a7b9145c01..93e1a0c51f 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -82,11 +82,9 @@ public interface EdgeService { void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId); - ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId); + PageData findEdgesByTenantIdAndEntityId(TenantId tenantId, EntityId ruleChainId, PageLink pageLink); - ListenableFuture> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId); - - ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId); + PageData findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink); Object checkInstance(Object request); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java index d254c5ee9a..49cbfce81c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.edge; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.EntitySubtype; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeInfo; import org.thingsboard.server.common.data.id.TenantId; @@ -160,20 +161,13 @@ public interface EdgeDao extends Dao { PageData findEdgeInfosByTenantId(UUID tenantId, PageLink pageLink); /** - * Find edges by tenantId and ruleChainId. + * Find edges by tenantId and entityId. * * @param tenantId the tenantId - * @param ruleChainId the ruleChainId - * @return the list of rule chain objects + * @param entityId the entityId + * @param entityType the entityType + * @return the list of edge objects */ - ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId); + PageData findEdgesByTenantIdAndEntityId(UUID tenantId, UUID entityId, EntityType entityType, PageLink pageLink); - /** - * Find edges by tenantId and dashboardId. - * - * @param tenantId the tenantId - * @param dashboardId the dashboardId - * @return the list of rule chain objects - */ - ListenableFuture> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId); } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index aee1fe7d34..b19733ff17 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -48,7 +48,6 @@ import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeInfo; import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; -import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.IdBased; @@ -59,7 +58,6 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; -import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; import org.thingsboard.server.dao.customer.CustomerDao; @@ -106,7 +104,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic private static final ObjectMapper mapper = new ObjectMapper(); - private static final int DEFAULT_LIMIT = 100; + private static final int DEFAULT_PAGE_SIZE = 1000; private RestTemplate restTemplate; @@ -380,7 +378,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic @Override public void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId) { log.trace("Executing assignDefaultRuleChainsToEdge, tenantId [{}], edgeId [{}]", tenantId, edgeId); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = ruleChainService.findAutoAssignToEdgeRuleChainsByTenantId(tenantId, pageLink); @@ -396,19 +394,11 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } @Override - public ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId) { - log.trace("Executing findEdgesByTenantIdAndRuleChainId, tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId); + public PageData findEdgesByTenantIdAndEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) { + log.trace("Executing findEdgesByTenantIdAndEntityId, tenantId [{}], entityId [{}], pageLink [{}]", tenantId, entityId, pageLink); Validator.validateId(tenantId, "Incorrect tenantId " + tenantId); - Validator.validateId(ruleChainId, "Incorrect ruleChainId " + ruleChainId); - return edgeDao.findEdgesByTenantIdAndRuleChainId(tenantId.getId(), ruleChainId.getId()); - } - - @Override - public ListenableFuture> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId) { - log.trace("Executing findEdgesByTenantIdAndDashboardId, tenantId [{}], dashboardId [{}]", tenantId, dashboardId); - Validator.validateId(tenantId, "Incorrect tenantId " + tenantId); - Validator.validateId(dashboardId, "Incorrect dashboardId " + dashboardId); - return edgeDao.findEdgesByTenantIdAndDashboardId(tenantId.getId(), dashboardId.getId()); + validatePageLink(pageLink); + return edgeDao.findEdgesByTenantIdAndEntityId(tenantId.getId(), entityId.getId(), entityId.getEntityType(), pageLink); } private DataValidator edgeValidator = @@ -496,88 +486,55 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic }; @Override - public ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) { - // TODO: @voba - rewrite 'find' to use native SQL queries instead of fetching relations - - log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}]", tenantId, entityId); + public PageData findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) { + log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}] [{}]", tenantId, entityId, pageLink); if (EntityType.TENANT.equals(entityId.getEntityType()) || EntityType.CUSTOMER.equals(entityId.getEntityType()) || EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) { - List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); - PageData pageData; - do { - if (EntityType.TENANT.equals(entityId.getEntityType()) || - EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) { - pageData = findEdgesByTenantId(tenantId, pageLink); - } else { - pageData = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), pageLink); - } - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - for (Edge edge : pageData.getData()) { - result.add(edge.getId()); - } - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData != null && pageData.hasNext()); - return Futures.immediateFuture(result); + if (EntityType.TENANT.equals(entityId.getEntityType()) || + EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) { + return convertToEdgeIds(findEdgesByTenantId(tenantId, pageLink)); + } else { + return convertToEdgeIds(findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), pageLink)); + } } else { switch (entityId.getEntityType()) { case DEVICE: case ASSET: case ENTITY_VIEW: - ListenableFuture> originatorEdgeRelationsFuture = - relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> { - if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0 && - originatorEdgeRelations.get(0).getFrom() != null) { - return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId())); - } else { - return Collections.emptyList(); - } - }, MoreExecutors.directExecutor()); case DASHBOARD: - return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId()))); case RULE_CHAIN: - return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId()))); + return convertToEdgeIds(findEdgesByTenantIdAndEntityId(tenantId, entityId, pageLink)); case USER: User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); if (userById == null) { - return Futures.immediateFuture(Collections.emptyList()); + return createEmptyEdgeIdPageData(); + } + if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) { + return convertToEdgeIds(findEdgesByTenantId(tenantId, pageLink)); + } else { + return convertToEdgeIds(findEdgesByTenantIdAndCustomerId(tenantId, userById.getCustomerId(), pageLink)); } - List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); - PageData pageData; - do { - if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) { - pageData = findEdgesByTenantId(tenantId, pageLink); - } else { - pageData = findEdgesByTenantIdAndCustomerId(tenantId, userById.getCustomerId(), pageLink); - } - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - result.addAll(pageData.getData()); - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData != null && pageData.hasNext()); - return convertToEdgeIds(Futures.immediateFuture(result)); default: - return Futures.immediateFuture(Collections.emptyList()); + log.warn("[{}] Unsupported entity type {}", tenantId, entityId.getEntityType()); + return createEmptyEdgeIdPageData(); } } } - private ListenableFuture> convertToEdgeIds(ListenableFuture> future) { - return Futures.transform(future, edges -> { - if (edges != null && !edges.isEmpty()) { - return edges.stream().map(IdBased::getId).collect(Collectors.toList()); - } else { - return Collections.emptyList(); - } - }, MoreExecutors.directExecutor()); + private PageData createEmptyEdgeIdPageData() { + return new PageData<>(new ArrayList<>(), 0, 0, false); + } + + private PageData convertToEdgeIds(PageData pageData) { + if (pageData == null) { + return createEmptyEdgeIdPageData(); + } + List edgeIds = new ArrayList<>(); + if (pageData.getData() != null && !pageData.getData().isEmpty()) { + edgeIds = pageData.getData().stream().map(IdBased::getId).collect(Collectors.toList()); + } + return new PageData<>(edgeIds, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); } @Override @@ -625,7 +582,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic private List findEdgeRuleChains(TenantId tenantId, EdgeId edgeId) { List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_LIMIT); + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; do { pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java index 5da1d56181..4d2a299576 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java @@ -77,8 +77,8 @@ public abstract class AbstractEntityService { List entityViews = entityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId).get(); if (entityViews != null && !entityViews.isEmpty()) { EntityView entityView = entityViews.get(0); - // TODO: @voba - refactor this blocking operation in 3.3+ - Boolean relationExists = relationService.checkRelation(tenantId,edgeId, entityView.getId(), + // TODO: @voba - refactor this blocking operation + Boolean relationExists = relationService.checkRelation(tenantId, edgeId, entityView.getId(), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE).get(); if (relationExists) { throw new DataValidationException("Can't unassign device/asset from edge that is related to entity view and entity view is assigned to edge!"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 38eeab9546..2ead3f4334 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -66,7 +66,6 @@ import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.TENANT; @@ -375,18 +374,21 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC throw new DataValidationException("Deletion of Root Tenant Rule Chain is prohibited!"); } if (RuleChainType.EDGE.equals(ruleChain.getType())) { - try { - List edges = edgeService.findEdgesByTenantIdAndRuleChainId(tenantId, ruleChainId).get(); - if (edges != null && !edges.isEmpty()) { - for (Edge edge : edges) { + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = edgeService.findEdgesByTenantIdAndEntityId(tenantId, ruleChainId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (Edge edge : pageData.getData()) { if (edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) { throw new DataValidationException("Can't delete rule chain that is root for edge [" + edge.getName() + "]. Please assign another root rule chain first to the edge!"); } } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } } - } catch (InterruptedException | ExecutionException e) { - log.error("Can't get edges by tenant id [{}] and rule chain id [{}]", tenantId.getId(), ruleChainId.getId(), e); - } + } while (pageData != null && pageData.hasNext()); } } checkRuleNodesAndDelete(tenantId, ruleChainId); @@ -638,7 +640,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC } private void checkRuleNodesAndDelete(TenantId tenantId, RuleChainId ruleChainId) { - try{ + try { ruleChainDao.removeById(tenantId, ruleChainId.getId()); } catch (Exception t) { ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); @@ -673,7 +675,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Override protected void validateCreate(TenantId tenantId, RuleChain data) { DefaultTenantProfileConfiguration profileConfiguration = - (DefaultTenantProfileConfiguration)tenantProfileCache.get(tenantId).getProfileData().getConfiguration(); + (DefaultTenantProfileConfiguration) tenantProfileCache.get(tenantId).getProfileData().getConfiguration(); long maxRuleChains = profileConfiguration.getMaxRuleChains(); validateNumberOfEntitiesPerTenant(tenantId, ruleChainDao, maxRuleChains, EntityType.RULE_CHAIN); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java index a5216033b3..f12dda790d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java @@ -22,6 +22,7 @@ import org.springframework.data.repository.PagingAndSortingRepository; import org.springframework.data.repository.query.Param; import org.thingsboard.server.dao.model.sql.EdgeEntity; import org.thingsboard.server.dao.model.sql.EdgeInfoEntity; +import org.thingsboard.server.dao.model.sql.RuleChainEntity; import java.util.List; import java.util.UUID; @@ -110,6 +111,16 @@ public interface EdgeRepository extends PagingAndSortingRepository findByTenantIdAndEntityId(@Param("tenantId") UUID tenantId, + @Param("entityId") UUID entityId, + @Param("entityType") String entityType, + @Param("searchText") String searchText, + Pageable pageable); + @Query("SELECT DISTINCT d.type FROM EdgeEntity d WHERE d.tenantId = :tenantId") List findTenantEdgeTypes(@Param("tenantId") UUID tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java index f17196fe93..eec6ee7f29 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java @@ -26,13 +26,10 @@ import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeInfo; -import org.thingsboard.server.common.data.id.DashboardId; -import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.relation.EntityRelation; -import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.edge.EdgeDao; import org.thingsboard.server.dao.model.sql.EdgeEntity; @@ -181,17 +178,15 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao imple } @Override - public ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId) { - log.debug("Try to find edges by tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId); - ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return transformFromRelationToEdge(tenantId, relations); - } - - @Override - public ListenableFuture> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId) { - log.debug("Try to find edges by tenantId [{}], dashboardId [{}]", tenantId, dashboardId); - ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new DashboardId(dashboardId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return transformFromRelationToEdge(tenantId, relations); + public PageData findEdgesByTenantIdAndEntityId(UUID tenantId, UUID entityId, EntityType entityType, PageLink pageLink) { + log.debug("Try to find edges by tenantId [{}], entityId [{}], entityType [{}], pageLink [{}]", tenantId, entityId, entityType, pageLink); + return DaoUtil.toPageData( + edgeRepository.findByTenantIdAndEntityId( + tenantId, + entityId, + entityType.name(), + Objects.toString(pageLink.getTextSearch(), ""), + DaoUtil.toPageable(pageLink))); } private ListenableFuture> transformFromRelationToEdge(UUID tenantId, ListenableFuture> relations) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index f208213d63..ce46ccffbf 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -15,7 +15,6 @@ */ package org.thingsboard.rule.engine.edge; -import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.FutureCallback; @@ -39,6 +38,8 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.msg.TbMsg; @@ -90,6 +91,8 @@ public class TbMsgPushToEdgeNode implements TbNode { private static final String SCOPE = "scope"; + private static final int DEFAULT_PAGE_SIZE = 1000; + @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = TbNodeUtils.convert(configuration, EmptyNodeConfiguration.class); @@ -117,77 +120,53 @@ public class TbMsgPushToEdgeNode implements TbNode { private void processMsg(TbContext ctx, TbMsg msg) { if (EntityType.EDGE.equals(msg.getOriginator().getEntityType())) { - try { - EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx); - if (edgeEvent != null) { - EdgeId edgeId = new EdgeId(msg.getOriginator().getId()); - edgeEvent.setEdgeId(edgeId); - ListenableFuture saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent); - Futures.addCallback(saveFuture, new FutureCallback() { - @Override - public void onSuccess(@Nullable EdgeEvent event) { - ctx.tellNext(msg, SUCCESS); - ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId); - } - - @Override - public void onFailure(Throwable th) { - log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th); - ctx.tellFailure(msg, th); - } - }, ctx.getDbCallbackExecutor()); - } - } catch (JsonProcessingException e) { - log.error("Failed to build edge event", e); - ctx.tellFailure(msg, e); + EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx); + if (edgeEvent != null) { + EdgeId edgeId = new EdgeId(msg.getOriginator().getId()); + notifyEdge(ctx, msg, edgeEvent, edgeId); } } else { - ListenableFuture> getEdgeIdsFuture = ctx.getEdgeService().findRelatedEdgeIdsByEntityId(ctx.getTenantId(), msg.getOriginator()); - Futures.addCallback(getEdgeIdsFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List edgeIds) { - if (edgeIds != null && !edgeIds.isEmpty()) { - for (EdgeId edgeId : edgeIds) { - try { - EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx); - if (edgeEvent == null) { - log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType()); - ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'")); - } else { - edgeEvent.setEdgeId(edgeId); - ListenableFuture saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent); - Futures.addCallback(saveFuture, new FutureCallback() { - @Override - public void onSuccess(@Nullable EdgeEvent event) { - ctx.tellNext(msg, SUCCESS); - ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId); - } - - @Override - public void onFailure(Throwable th) { - log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th); - ctx.tellFailure(msg, th); - } - }, ctx.getDbCallbackExecutor()); - } - } catch (JsonProcessingException e) { - log.error("Failed to build edge event", e); - ctx.tellFailure(msg, e); - } + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + do { + pageData = ctx.getEdgeService().findRelatedEdgeIdsByEntityId(ctx.getTenantId(), msg.getOriginator(), pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (EdgeId edgeId : pageData.getData()) { + EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx); + if (edgeEvent == null) { + log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType()); + ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'")); + } else { + notifyEdge(ctx, msg, edgeEvent, edgeId); } } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } } + } while (pageData != null && pageData.hasNext()); + } + } - @Override - public void onFailure(Throwable t) { - ctx.tellFailure(msg, t); - } + private void notifyEdge(TbContext ctx, TbMsg msg, EdgeEvent edgeEvent, EdgeId edgeId) { + edgeEvent.setEdgeId(edgeId); + ListenableFuture saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent); + Futures.addCallback(saveFuture, new FutureCallback() { + @Override + public void onSuccess(@Nullable EdgeEvent event) { + ctx.tellNext(msg, SUCCESS); + ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId); + } - }, ctx.getDbCallbackExecutor()); - } + @Override + public void onFailure(Throwable th) { + log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th); + ctx.tellFailure(msg, th); + } + }, ctx.getDbCallbackExecutor()); } - private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) throws JsonProcessingException { + private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) { String msgType = msg.getType(); if (DataConstants.ALARM.equals(msgType)) { return buildEdgeEvent(ctx.getTenantId(), EdgeEventActionType.ADDED, getUUIDFromMsgData(msg), EdgeEventType.ALARM, null);