From 9a03fbadc75d6ff3f1a08467645a2314224caa65 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Thu, 10 Dec 2020 17:55:16 +0200 Subject: [PATCH 1/5] added ability to get attributes and timeseries keys by entity query --- .../controller/EntityQueryController.java | 51 +++++++++++++++++++ .../dao/attributes/AttributesService.java | 3 ++ .../dao/timeseries/TimeseriesService.java | 2 + .../server/dao/attributes/AttributesDao.java | 3 ++ .../dao/attributes/BaseAttributesService.java | 6 +++ .../sql/attributes/AttributeKvRepository.java | 3 ++ .../dao/sql/attributes/JpaAttributeDao.java | 7 +++ .../dao/sqlts/SqlTimeseriesLatestDao.java | 6 +++ .../sqlts/latest/TsKvLatestRepository.java | 5 ++ .../dao/timeseries/BaseTimeseriesService.java | 5 ++ .../CassandraBaseTimeseriesLatestDao.java | 5 ++ .../dao/timeseries/TimeseriesLatestDao.java | 2 + 12 files changed, 98 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 94417886d6..e58eee59a6 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -23,15 +23,25 @@ import org.springframework.web.bind.annotation.RequestMethod; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.RestController; import org.thingsboard.server.common.data.exception.ThingsboardException; +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.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.common.data.query.EntityDataPageLink; import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.query.EntityQueryService; +import java.util.Collections; +import java.util.List; +import java.util.function.Function; +import java.util.stream.Collectors; + @RestController @TbCoreComponent @RequestMapping("/api") @@ -40,6 +50,12 @@ public class EntityQueryController extends BaseController { @Autowired private EntityQueryService entityQueryService; + @Autowired + private AttributesService attributesService; + + @Autowired + private TimeseriesService timeseriesService; + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST) @@ -76,4 +92,39 @@ public class EntityQueryController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") + @RequestMapping(value = "/entitiesQuery/find/keys/timeseries", method = RequestMethod.POST) + @ResponseBody + public List findEntityTimeseriesKeysByQuery(@RequestBody EntityDataQuery query) throws ThingsboardException { + TenantId tenantId = getTenantId(); + return getKeys(query, entityIds -> timeseriesService.findAllKeysByEntityIds(tenantId, entityIds)); + } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") + @RequestMapping(value = "/entitiesQuery/find/keys/attributes", method = RequestMethod.POST) + @ResponseBody + public List findEntityAttributesKeysByQuery(@RequestBody EntityDataQuery query) throws ThingsboardException { + TenantId tenantId = getTenantId(); + return getKeys(query, entityIds -> attributesService.findAllKeysByEntityIds(tenantId, entityIds.get(0).getEntityType(), entityIds)); + } + + private List getKeys(EntityDataQuery query, Function, List> function) throws ThingsboardException { + checkNotNull(query); + try { + EntityDataPageLink pageLink = query.getPageLink(); + if (pageLink.getPageSize() > 100) { + pageLink.setPageSize(100); + } + List ids = this.entityQueryService.findEntityDataByQuery(getCurrentUser(), query).getData().stream() + .map(EntityData::getEntityId) + .collect(Collectors.toList()); + if (ids.isEmpty()) { + return Collections.emptyList(); + } + return function.apply(ids); + } catch (Exception e) { + throw handleException(e); + } + } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java index 2e5c895f84..47f1f91376 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.attributes; import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -42,4 +43,6 @@ public interface AttributesService { List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds); + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index dbfbe3d0f5..ccf933e03d 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -50,4 +50,6 @@ public interface TimeseriesService { ListenableFuture> removeAllLatest(TenantId tenantId, EntityId entityId); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, List entityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java index 458a251f8a..f1af5a1b1b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.attributes; import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -41,4 +42,6 @@ public interface AttributesDao { ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List keys); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java index 2e18006036..24f988c79d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java @@ -20,6 +20,7 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -65,6 +66,11 @@ public class BaseAttributesService implements AttributesService { return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds) { + return attributesDao.findAllKeysByEntityIds(tenantId, entityType, entityIds); + } + @Override public ListenableFuture> save(TenantId tenantId, EntityId entityId, String scope, List attributes) { validate(entityId, scope); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java index f6c3e195e7..3d14730b54 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java @@ -56,5 +56,8 @@ public interface AttributeKvRepository extends CrudRepository findAllKeysByTenantId(@Param("tenantId") UUID tenantId); + @Query(value = "SELECT DISTINCT attribute_key FROM attribute_kv WHERE entity_type = :entityType " + + "AND entity_id in :entityIds ORDER BY attribute_key", nativeQuery = true) + List findAllKeysByEntityIds(@Param("entityType") String entityType, @Param("entityIds") List entityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index a95bd4b612..a6a9348e8f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java @@ -22,6 +22,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -145,6 +146,12 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl } } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds) { + return attributeKvRepository + .findAllKeysByEntityIds(entityType.name(), entityIds.stream().map(EntityId::getId).collect(Collectors.toList())); + } + @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) { AttributeKvEntity entity = new AttributeKvEntity(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java index ff602e6eed..a97fa97a5c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java @@ -61,6 +61,7 @@ import java.util.Optional; import java.util.UUID; import java.util.concurrent.ExecutionException; import java.util.function.Function; +import java.util.stream.Collectors; @Slf4j @Component @@ -169,6 +170,11 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme } } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return tsKvLatestRepository.findAllKeysByEntityIds(entityIds.stream().map(EntityId::getId).collect(Collectors.toList())); + } + private ListenableFuture getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { ListenableFuture> future = findNewLatestEntryFuture(tenantId, entityId, query); return Futures.transformAsync(future, entryList -> { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java index cd3db69d70..da8c921487 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java @@ -36,4 +36,9 @@ public interface TsKvLatestRepository extends CrudRepository getKeysByTenantId(@Param("tenant_id") UUID tenantId); + @Query(value = "SELECT DISTINCT ts_kv_dictionary.key AS strKey FROM ts_kv_latest " + + "INNER JOIN ts_kv_dictionary ON ts_kv_latest.key = ts_kv_dictionary.key_id " + + "WHERE ts_kv_latest.entity_id IN :entityIds ORDER BY ts_kv_dictionary.key", nativeQuery = true) + List findAllKeysByEntityIds(@Param("entityIds") List entityIds); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 7160baf939..12e081d234 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -121,6 +121,11 @@ public class BaseTimeseriesService implements TimeseriesService { return timeseriesLatestDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return timeseriesLatestDao.findAllKeysByEntityIds(tenantId, entityIds); + } + @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { validate(entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java index f086f913c0..2cb67f1260 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java @@ -86,6 +86,11 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes return Collections.emptyList(); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return Collections.emptyList(); + } + @Override public ListenableFuture saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java index 0b7156e03e..b2cd346277 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java @@ -35,4 +35,6 @@ public interface TimeseriesLatestDao { ListenableFuture removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, List entityIds); } From f63b4b1f7c3b843b391cb244307538400439e9e4 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 15 Dec 2020 16:16:02 +0200 Subject: [PATCH 2/5] created findEntityTimeseriesAndAttributesKeysByQuery instead findEntityTimeseriesKeysByQuery and findEntityAttributesKeysByQuery --- .../controller/EntityQueryController.java | 45 ++---- .../query/DefaultEntityQueryService.java | 135 ++++++++++++++++++ .../service/query/EntityQueryService.java | 6 + 3 files changed, 151 insertions(+), 35 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index e58eee59a6..1fa802e18a 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -16,14 +16,16 @@ package org.thingsboard.server.controller; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.ResponseEntity; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; +import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.context.request.async.DeferredResult; import org.thingsboard.server.common.data.exception.ThingsboardException; -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.query.AlarmData; @@ -32,16 +34,9 @@ import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataPageLink; import org.thingsboard.server.common.data.query.EntityDataQuery; -import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.query.EntityQueryService; -import java.util.Collections; -import java.util.List; -import java.util.function.Function; -import java.util.stream.Collectors; - @RestController @TbCoreComponent @RequestMapping("/api") @@ -50,13 +45,6 @@ public class EntityQueryController extends BaseController { @Autowired private EntityQueryService entityQueryService; - @Autowired - private AttributesService attributesService; - - @Autowired - private TimeseriesService timeseriesService; - - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST) @ResponseBody @@ -96,35 +84,22 @@ public class EntityQueryController extends BaseController { @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/entitiesQuery/find/keys/timeseries", method = RequestMethod.POST) @ResponseBody - public List findEntityTimeseriesKeysByQuery(@RequestBody EntityDataQuery query) throws ThingsboardException { + public DeferredResult findEntityTimeseriesAndAttributesKeysByQuery(@RequestBody EntityDataQuery query, + @RequestParam("timeseries") boolean isTimeseries, + @RequestParam("attributes") boolean isAttributes) throws ThingsboardException { TenantId tenantId = getTenantId(); - return getKeys(query, entityIds -> timeseriesService.findAllKeysByEntityIds(tenantId, entityIds)); - } - - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") - @RequestMapping(value = "/entitiesQuery/find/keys/attributes", method = RequestMethod.POST) - @ResponseBody - public List findEntityAttributesKeysByQuery(@RequestBody EntityDataQuery query) throws ThingsboardException { - TenantId tenantId = getTenantId(); - return getKeys(query, entityIds -> attributesService.findAllKeysByEntityIds(tenantId, entityIds.get(0).getEntityType(), entityIds)); - } - - private List getKeys(EntityDataQuery query, Function, List> function) throws ThingsboardException { checkNotNull(query); try { EntityDataPageLink pageLink = query.getPageLink(); if (pageLink.getPageSize() > 100) { pageLink.setPageSize(100); } - List ids = this.entityQueryService.findEntityDataByQuery(getCurrentUser(), query).getData().stream() - .map(EntityData::getEntityId) - .collect(Collectors.toList()); - if (ids.isEmpty()) { - return Collections.emptyList(); - } - return function.apply(ids); + DeferredResult response = new DeferredResult<>(); + entityQueryService.getKeysByQueryCallback(getCurrentUser(), tenantId, query, isTimeseries, isAttributes, response); + return response; } catch (Exception e) { throw handleException(e); } } + } diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index b8ccad466d..337e2cdb07 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -15,11 +15,24 @@ */ package org.thingsboard.server.service.query; +import com.datastax.oss.driver.internal.core.util.CollectionsUtils; +import com.fasterxml.jackson.databind.node.ArrayNode; +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 lombok.extern.slf4j.Slf4j; +import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; import org.springframework.stereotype.Service; +import org.springframework.util.CollectionUtils; +import org.springframework.web.context.request.async.DeferredResult; +import org.thingsboard.server.common.data.EntityType; 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.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; @@ -31,12 +44,24 @@ import org.thingsboard.server.common.data.query.EntityDataSortOrder; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.dao.alarm.AlarmService; +import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.util.mapping.JacksonUtil; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.security.AccessValidator; import org.thingsboard.server.service.security.model.SecurityUser; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; @Service @Slf4j @@ -52,6 +77,15 @@ public class DefaultEntityQueryService implements EntityQueryService { @Value("${server.ws.max_entities_per_alarm_subscription:1000}") private int maxEntitiesPerAlarmSubscription; + @Autowired + private DbCallbackExecutorService dbCallbackExecutor; + + @Autowired + private TimeseriesService timeseriesService; + + @Autowired + private AttributesService attributesService; + @Override public long countEntitiesByQuery(SecurityUser securityUser, EntityCountQuery query) { return entityService.countEntitiesByQuery(securityUser.getTenantId(), securityUser.getCustomerId(), query); @@ -89,6 +123,107 @@ public class DefaultEntityQueryService implements EntityQueryService { } } + @Override + public void getKeysByQueryCallback(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes, DeferredResult response) { + if (!isAttributes && !isTimeseries) { + getEmptyResponseCallback(response); + return; + } + + List ids = this.findEntityDataByQuery(securityUser, query).getData().stream() + .map(EntityData::getEntityId) + .collect(Collectors.toList()); + if (ids.isEmpty()) { + getEmptyResponseCallback(response); + return; + } + + Set types = ids.stream().map(EntityId::getEntityType).collect(Collectors.toSet()); + ListenableFuture> timeseriesKeysFuture; + ListenableFuture> attributesKeysFuture; + + if (isTimeseries) { + timeseriesKeysFuture = dbCallbackExecutor.submit(() -> timeseriesService.findAllKeysByEntityIds(tenantId, ids)); + } else { + timeseriesKeysFuture = null; + } + + if (isAttributes) { + Map> typesMap = ids.stream().collect(Collectors.groupingBy(EntityId::getEntityType)); + List>> futures = new ArrayList<>(typesMap.size()); + typesMap.forEach((type, entityIds) -> futures.add(dbCallbackExecutor.submit(() -> attributesService.findAllKeysByEntityIds(tenantId, type, entityIds)))); + attributesKeysFuture = Futures.transform(Futures.allAsList(futures), lists -> { + if (CollectionUtils.isEmpty(lists)) { + return null; + } + + return lists.stream().flatMap(List::stream).distinct().sorted().collect(Collectors.toList()); + }, dbCallbackExecutor); + } else { + attributesKeysFuture = null; + } + + if (timeseriesKeysFuture != null && attributesKeysFuture != null) { + Futures.whenAllComplete(timeseriesKeysFuture, attributesKeysFuture).call(() -> { + try { + getResponseCallback(response, types, timeseriesKeysFuture.get(), attributesKeysFuture.get()); + } catch (Exception e) { + log.error("Failed to fetch timeseries and attributes keys!", e); + AccessValidator.handleError(e, response, HttpStatus.INTERNAL_SERVER_ERROR); + } + + return null; + }, dbCallbackExecutor); + } else if (timeseriesKeysFuture != null) { + Futures.addCallback(timeseriesKeysFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List keys) { + getResponseCallback(response, types, keys, null); + } + + @Override + public void onFailure(Throwable t) { + log.error("Failed to fetch timeseries keys!", t); + AccessValidator.handleError(t, response, HttpStatus.INTERNAL_SERVER_ERROR); + } + + }, dbCallbackExecutor); + } else { + Futures.addCallback(attributesKeysFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List keys) { + getResponseCallback(response, types, null, keys); + } + + @Override + public void onFailure(Throwable t) { + log.error("Failed to fetch attributes keys!", t); + AccessValidator.handleError(t, response, HttpStatus.INTERNAL_SERVER_ERROR); + } + }, dbCallbackExecutor); + } + } + + private void getResponseCallback(DeferredResult response, Set types, List timeseriesKeys, List attributesKeys) { + ObjectNode json = JacksonUtil.newObjectNode(); + addItemsToArrayNode(json.putArray("types"), types); + addItemsToArrayNode(json.putArray("timeseriesKeys"), timeseriesKeys); + addItemsToArrayNode(json.putArray("attributesKeys"), attributesKeys); + + response.setResult(new ResponseEntity(json, HttpStatus.OK)); + } + + private void getEmptyResponseCallback(DeferredResult response) { + getResponseCallback(response, null, null, null); + } + + private void addItemsToArrayNode(ArrayNode arrayNode, Collection collection) { + if (!CollectionUtils.isEmpty(collection)) { + collection.forEach(item -> arrayNode.add(item.toString())); + } + } + private EntityDataQuery buildEntityDataQuery(AlarmDataQuery query) { EntityDataSortOrder sortOrder = query.getPageLink().getSortOrder(); EntityDataSortOrder entitiesSortOrder; diff --git a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java index 15f7d86252..459fb144d1 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java @@ -15,6 +15,9 @@ */ package org.thingsboard.server.service.query; +import org.springframework.http.ResponseEntity; +import org.springframework.web.context.request.async.DeferredResult; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; @@ -31,4 +34,7 @@ public interface EntityQueryService { PageData findAlarmDataByQuery(SecurityUser securityUser, AlarmDataQuery query); + void getKeysByQueryCallback(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes, DeferredResult response); + } From dfb82bf28fb13ed7cf590eb37366487e7a6a80da Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 16 Dec 2020 11:05:35 +0200 Subject: [PATCH 3/5] findEntityTimeseriesAndAttributesKeysByQuery improvements --- .../controller/EntityQueryController.java | 10 +- .../query/DefaultEntityQueryService.java | 110 +++++++++--------- .../service/query/EntityQueryService.java | 4 +- 3 files changed, 61 insertions(+), 63 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 1fa802e18a..4325d6f2be 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -45,6 +45,8 @@ public class EntityQueryController extends BaseController { @Autowired private EntityQueryService entityQueryService; + private static final int MAX_PAGE_SIZE = 100; + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST) @ResponseBody @@ -91,12 +93,10 @@ public class EntityQueryController extends BaseController { checkNotNull(query); try { EntityDataPageLink pageLink = query.getPageLink(); - if (pageLink.getPageSize() > 100) { - pageLink.setPageSize(100); + if (pageLink.getPageSize() > MAX_PAGE_SIZE) { + pageLink.setPageSize(MAX_PAGE_SIZE); } - DeferredResult response = new DeferredResult<>(); - entityQueryService.getKeysByQueryCallback(getCurrentUser(), tenantId, query, isTimeseries, isAttributes, response); - return response; + return entityQueryService.getKeysByQuery(getCurrentUser(), tenantId, query, isTimeseries, isAttributes); } catch (Exception e) { throw handleException(e); } diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index 337e2cdb07..d8be5040c9 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.query; -import com.datastax.oss.driver.internal.core.util.CollectionsUtils; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; @@ -61,6 +60,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.function.Consumer; import java.util.stream.Collectors; @Service @@ -123,25 +123,38 @@ public class DefaultEntityQueryService implements EntityQueryService { } } + private EntityDataQuery buildEntityDataQuery(AlarmDataQuery query) { + EntityDataSortOrder sortOrder = query.getPageLink().getSortOrder(); + EntityDataSortOrder entitiesSortOrder; + if (sortOrder == null || sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) { + entitiesSortOrder = new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, ModelConstants.CREATED_TIME_PROPERTY)); + } else { + entitiesSortOrder = sortOrder; + } + EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder); + return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters()); + } + @Override - public void getKeysByQueryCallback(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, - boolean isTimeseries, boolean isAttributes, DeferredResult response) { + public DeferredResult getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes) { + final DeferredResult response = new DeferredResult<>(); if (!isAttributes && !isTimeseries) { - getEmptyResponseCallback(response); - return; + replyWithEmptyResponse(response); + return response; } List ids = this.findEntityDataByQuery(securityUser, query).getData().stream() .map(EntityData::getEntityId) .collect(Collectors.toList()); if (ids.isEmpty()) { - getEmptyResponseCallback(response); - return; + replyWithEmptyResponse(response); + return response; } Set types = ids.stream().map(EntityId::getEntityType).collect(Collectors.toSet()); - ListenableFuture> timeseriesKeysFuture; - ListenableFuture> attributesKeysFuture; + final ListenableFuture> timeseriesKeysFuture; + final ListenableFuture> attributesKeysFuture; if (isTimeseries) { timeseriesKeysFuture = dbCallbackExecutor.submit(() -> timeseriesService.findAllKeysByEntityIds(tenantId, ids)); @@ -155,67 +168,49 @@ public class DefaultEntityQueryService implements EntityQueryService { typesMap.forEach((type, entityIds) -> futures.add(dbCallbackExecutor.submit(() -> attributesService.findAllKeysByEntityIds(tenantId, type, entityIds)))); attributesKeysFuture = Futures.transform(Futures.allAsList(futures), lists -> { if (CollectionUtils.isEmpty(lists)) { - return null; + return Collections.emptyList(); } - return lists.stream().flatMap(List::stream).distinct().sorted().collect(Collectors.toList()); }, dbCallbackExecutor); } else { attributesKeysFuture = null; } - if (timeseriesKeysFuture != null && attributesKeysFuture != null) { - Futures.whenAllComplete(timeseriesKeysFuture, attributesKeysFuture).call(() -> { + if (isTimeseries && isAttributes) { + Futures.whenAllComplete(timeseriesKeysFuture, attributesKeysFuture).run(() -> { try { - getResponseCallback(response, types, timeseriesKeysFuture.get(), attributesKeysFuture.get()); + replyWithResponse(response, types, timeseriesKeysFuture.get(), attributesKeysFuture.get()); } catch (Exception e) { log.error("Failed to fetch timeseries and attributes keys!", e); AccessValidator.handleError(e, response, HttpStatus.INTERNAL_SERVER_ERROR); } - - return null; - }, dbCallbackExecutor); - } else if (timeseriesKeysFuture != null) { - Futures.addCallback(timeseriesKeysFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List keys) { - getResponseCallback(response, types, keys, null); - } - - @Override - public void onFailure(Throwable t) { - log.error("Failed to fetch timeseries keys!", t); - AccessValidator.handleError(t, response, HttpStatus.INTERNAL_SERVER_ERROR); - } - }, dbCallbackExecutor); + } else if (isTimeseries) { + addCallback(timeseriesKeysFuture, keys -> replyWithResponse(response, types, keys, null), + error -> { + log.error("Failed to fetch timeseries keys!", error); + AccessValidator.handleError(error, response, HttpStatus.INTERNAL_SERVER_ERROR); + }); } else { - Futures.addCallback(attributesKeysFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List keys) { - getResponseCallback(response, types, null, keys); - } - - @Override - public void onFailure(Throwable t) { - log.error("Failed to fetch attributes keys!", t); - AccessValidator.handleError(t, response, HttpStatus.INTERNAL_SERVER_ERROR); - } - }, dbCallbackExecutor); + addCallback(attributesKeysFuture, keys -> replyWithResponse(response, types, null, keys), + error -> { + log.error("Failed to fetch attributes keys!", error); + AccessValidator.handleError(error, response, HttpStatus.INTERNAL_SERVER_ERROR); + }); } + return response; } - private void getResponseCallback(DeferredResult response, Set types, List timeseriesKeys, List attributesKeys) { + private void replyWithResponse(DeferredResult response, Set types, List timeseriesKeys, List attributesKeys) { ObjectNode json = JacksonUtil.newObjectNode(); addItemsToArrayNode(json.putArray("types"), types); addItemsToArrayNode(json.putArray("timeseriesKeys"), timeseriesKeys); addItemsToArrayNode(json.putArray("attributesKeys"), attributesKeys); - response.setResult(new ResponseEntity(json, HttpStatus.OK)); } - private void getEmptyResponseCallback(DeferredResult response) { - getResponseCallback(response, null, null, null); + private void replyWithEmptyResponse(DeferredResult response) { + replyWithResponse(response, Collections.emptySet(), Collections.emptyList(), Collections.emptyList()); } private void addItemsToArrayNode(ArrayNode arrayNode, Collection collection) { @@ -224,15 +219,18 @@ public class DefaultEntityQueryService implements EntityQueryService { } } - private EntityDataQuery buildEntityDataQuery(AlarmDataQuery query) { - EntityDataSortOrder sortOrder = query.getPageLink().getSortOrder(); - EntityDataSortOrder entitiesSortOrder; - if (sortOrder == null || sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) { - entitiesSortOrder = new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, ModelConstants.CREATED_TIME_PROPERTY)); - } else { - entitiesSortOrder = sortOrder; - } - EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder); - return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters()); + private void addCallback(ListenableFuture> future, Consumer> success, Consumer error) { + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List keys) { + success.accept(keys); + } + + @Override + public void onFailure(Throwable t) { + error.accept(t); + } + }, dbCallbackExecutor); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java index 459fb144d1..763453c815 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java @@ -34,7 +34,7 @@ public interface EntityQueryService { PageData findAlarmDataByQuery(SecurityUser securityUser, AlarmDataQuery query); - void getKeysByQueryCallback(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, - boolean isTimeseries, boolean isAttributes, DeferredResult response); + DeferredResult getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes); } From 91bb1ed504af2c413863c9b95c4d5d1bf3c3563f Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 16 Dec 2020 11:21:28 +0200 Subject: [PATCH 4/5] refactored findEntityTimeseriesAndAttributesKeysByQuery --- .../thingsboard/server/controller/EntityQueryController.java | 2 +- .../server/service/query/DefaultEntityQueryService.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 4325d6f2be..2df8298d23 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -84,7 +84,7 @@ public class EntityQueryController extends BaseController { } @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") - @RequestMapping(value = "/entitiesQuery/find/keys/timeseries", method = RequestMethod.POST) + @RequestMapping(value = "/entitiesQuery/find/keys", method = RequestMethod.POST) @ResponseBody public DeferredResult findEntityTimeseriesAndAttributesKeysByQuery(@RequestBody EntityDataQuery query, @RequestParam("timeseries") boolean isTimeseries, diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index d8be5040c9..74c1b844ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -203,7 +203,7 @@ public class DefaultEntityQueryService implements EntityQueryService { private void replyWithResponse(DeferredResult response, Set types, List timeseriesKeys, List attributesKeys) { ObjectNode json = JacksonUtil.newObjectNode(); - addItemsToArrayNode(json.putArray("types"), types); + addItemsToArrayNode(json.putArray("entityTypes"), types); addItemsToArrayNode(json.putArray("timeseriesKeys"), timeseriesKeys); addItemsToArrayNode(json.putArray("attributesKeys"), attributesKeys); response.setResult(new ResponseEntity(json, HttpStatus.OK)); From 3aa97eec06b2fa94e56ce4b66dff542a73ee9804 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Thu, 17 Dec 2020 13:32:03 +0200 Subject: [PATCH 5/5] UI: Improvement autocomplete data keys in datasource widget --- .../query/DefaultEntityQueryService.java | 4 +- ui-ngx/src/app/core/http/entity.service.ts | 74 ++++++++++++++++++- .../widget/data-key-config.component.ts | 34 ++++++--- .../widget/data-keys.component.models.ts | 2 +- .../components/widget/data-keys.component.ts | 37 +++++++--- .../widget/widget-config.component.ts | 62 +++------------- ui-ngx/src/app/shared/models/entity.models.ts | 6 ++ 7 files changed, 139 insertions(+), 80 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index 74c1b844ef..c9726693bf 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -204,8 +204,8 @@ public class DefaultEntityQueryService implements EntityQueryService { private void replyWithResponse(DeferredResult response, Set types, List timeseriesKeys, List attributesKeys) { ObjectNode json = JacksonUtil.newObjectNode(); addItemsToArrayNode(json.putArray("entityTypes"), types); - addItemsToArrayNode(json.putArray("timeseriesKeys"), timeseriesKeys); - addItemsToArrayNode(json.putArray("attributesKeys"), attributesKeys); + addItemsToArrayNode(json.putArray("timeseries"), timeseriesKeys); + addItemsToArrayNode(json.putArray("attribute"), attributesKeys); response.setResult(new ResponseEntity(json, HttpStatus.OK)); } diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index 9263838ffa..fee7fb0094 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/ui-ngx/src/app/core/http/entity.service.ts @@ -41,10 +41,16 @@ import { AttributeScope, DataKeyType } from '@shared/models/telemetry/telemetry. import { defaultHttpOptionsFromConfig, RequestConfig } from '@core/http/http-utils'; import { RuleChainService } from '@core/http/rule-chain.service'; import { AliasInfo, StateParams, SubscriptionInfo } from '@core/api/widget-api.models'; -import { Datasource, DatasourceType, KeyInfo } from '@app/shared/models/widget.models'; +import { DataKey, Datasource, DatasourceType, KeyInfo } from '@app/shared/models/widget.models'; import { UtilsService } from '@core/services/utils.service'; import { AliasFilterType, EntityAlias, EntityAliasFilter, EntityAliasFilterResult } from '@shared/models/alias.models'; -import { entityFields, EntityInfo, ImportEntitiesResultInfo, ImportEntityData } from '@shared/models/entity.models'; +import { + EntitiesKeysByQuery, + entityFields, + EntityInfo, + ImportEntitiesResultInfo, + ImportEntityData +} from '@shared/models/entity.models'; import { EntityRelationService } from '@core/http/entity-relation.service'; import { deepClone, isDefined, isDefinedAndNotNull } from '@core/utils'; import { Asset } from '@shared/models/asset.models'; @@ -376,6 +382,13 @@ export class EntityService { return this.http.post>('/api/entitiesQuery/find', query, defaultHttpOptionsFromConfig(config)); } + public findEntityKeysByQuery(query: EntityDataQuery, attributes = true, timeseries = true, + config?: RequestConfig): Observable { + return this.http.post( + `/api/entitiesQuery/find/keys?attributes=${attributes}×eries=${timeseries}`, + query, defaultHttpOptionsFromConfig(config)); + } + public findAlarmDataByQuery(query: AlarmDataQuery, config?: RequestConfig): Observable> { return this.http.post>('/api/alarmsQuery/find', query, defaultHttpOptionsFromConfig(config)); } @@ -595,7 +608,7 @@ export class EntityService { return entityTypes; } - private getEntityFieldKeys(entityType: EntityType, searchText: string): Array { + private getEntityFieldKeys(entityType: EntityType, searchText: string = ''): Array { const entityFieldKeys: string[] = [entityFields.createdTime.keyName]; const query = searchText.toLowerCase(); switch (entityType) { @@ -637,7 +650,7 @@ export class EntityService { return query ? entityFieldKeys.filter((entityField) => entityField.toLowerCase().indexOf(query) === 0) : entityFieldKeys; } - private getAlarmKeys(searchText: string): Array { + private getAlarmKeys(searchText: string = ''): Array { const alarmKeys: string[] = Object.keys(alarmFields); const query = searchText.toLowerCase(); return query ? alarmKeys.filter((alarmField) => alarmField.toLowerCase().indexOf(query) === 0) : alarmKeys; @@ -672,6 +685,59 @@ export class EntityService { ); } + public getEntityKeysByEntityFilter(filter: EntityFilter, types: DataKeyType[], config?: RequestConfig): Observable> { + if (!types.length) { + return of([]); + } + let entitiesKeysByQuery$: Observable; + if (filter !== null && types.some(type => [DataKeyType.timeseries, DataKeyType.attribute].includes(type))) { + const dataQuery = { + entityFilter: filter, + pageLink: createDefaultEntityDataPageLink(100), + }; + entitiesKeysByQuery$ = this.findEntityKeysByQuery(dataQuery, types.includes(DataKeyType.attribute), + types.includes(DataKeyType.timeseries), config); + } else { + entitiesKeysByQuery$ = of({ + attribute: [], + timeseries: [], + entityTypes: [], + }); + } + return entitiesKeysByQuery$.pipe( + map((entitiesKeys) => { + const dataKeys: Array = []; + types.forEach(type => { + let keys: Array; + switch (type) { + case DataKeyType.entityField: + if (entitiesKeys.entityTypes.length) { + const entitiesFields = []; + entitiesKeys.entityTypes.forEach(entityType => entitiesFields.push(...this.getEntityFieldKeys(entityType))); + keys = Array.from(new Set(entitiesFields)); + } + break; + case DataKeyType.alarm: + keys = this.getAlarmKeys(); + break; + case DataKeyType.attribute: + case DataKeyType.timeseries: + if (entitiesKeys[type].length) { + keys = entitiesKeys[type]; + } + break; + } + if (keys) { + dataKeys.push(...keys.map(key => { + return {name: key, type}; + })); + } + }); + return dataKeys; + }) + ); + } + public createDatasourcesFromSubscriptionsInfo(subscriptionsInfo: Array): Array { const datasources = subscriptionsInfo.map(subscriptionInfo => this.createDatasourceFromSubscriptionInfo(subscriptionInfo)); this.utils.generateColors(datasources); diff --git a/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts b/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts index fde4b1a32f..a4187d9b6c 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts @@ -36,7 +36,7 @@ import { EntityService } from '@core/http/entity.service'; import { DataKeysCallbacks } from '@home/components/widget/data-keys.component.models'; import { DataKeyType } from '@shared/models/telemetry/telemetry.models'; import { Observable, of } from 'rxjs'; -import { map, mergeMap, tap } from 'rxjs/operators'; +import { map, mergeMap, publishReplay, refCount, tap } from 'rxjs/operators'; import { alarmFields } from '@shared/models/alarm.models'; import { JsFuncComponent } from '@shared/components/js-func.component'; import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models'; @@ -95,6 +95,7 @@ export class DataKeyConfigComponent extends PageComponent implements OnInit, Con filteredKeys: Observable>; private latestKeySearchResult: Array = null; + private fetchObservable$: Observable> = null; keySearchText = ''; @@ -205,31 +206,42 @@ export class DataKeyConfigComponent extends PageComponent implements OnInit, Con } private fetchKeys(searchText?: string): Observable> { - if (this.latestKeySearchResult === null || this.keySearchText !== searchText) { + if (this.keySearchText !== searchText || this.latestKeySearchResult === null) { this.keySearchText = searchText; - let fetchObservable: Observable> = null; + const dataKeyFilter = this.createKeyFilter(this.keySearchText); + return this.getKeys().pipe( + map(name => name.filter(dataKeyFilter)), + tap(res => this.latestKeySearchResult = res) + ); + } + return of(this.latestKeySearchResult); + } + + private getKeys() { + if (this.fetchObservable$ === null) { + let fetchObservable: Observable>; if (this.modelValue.type === DataKeyType.alarm) { - const dataKeyFilter = this.createDataKeyFilter(this.keySearchText); - fetchObservable = of(this.alarmKeys.filter(dataKeyFilter)); + fetchObservable = of(this.alarmKeys); } else { if (this.entityAliasId) { const dataKeyTypes = [this.modelValue.type]; - fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, this.keySearchText, dataKeyTypes); + fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, dataKeyTypes); } else { fetchObservable = of([]); } } - return fetchObservable.pipe( + this.fetchObservable$ = fetchObservable.pipe( map((dataKeys) => dataKeys.map((dataKey) => dataKey.name)), - tap(res => this.latestKeySearchResult = res) + publishReplay(1), + refCount() ); } - return of(this.latestKeySearchResult); + return this.fetchObservable$; } - private createDataKeyFilter(query: string): (key: DataKey) => boolean { + private createKeyFilter(query: string): (key: string) => boolean { const lowercaseQuery = query.toLowerCase(); - return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; + return key => key.toLowerCase().startsWith(lowercaseQuery); } public validateOnSubmit() { diff --git a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts index 452e390faa..f5a937a4be 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts +++ b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts @@ -20,5 +20,5 @@ import { Observable } from 'rxjs'; export interface DataKeysCallbacks { generateDataKey: (chip: any, type: DataKeyType) => DataKey; - fetchEntityKeys: (entityAliasId: string, query: string, types: Array) => Observable>; + fetchEntityKeys: (entityAliasId: string, types: Array) => Observable>; } diff --git a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts index 22c9663be8..e31ad036c9 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts @@ -38,7 +38,7 @@ import { Validators } from '@angular/forms'; import { Observable, of } from 'rxjs'; -import { filter, map, mergeMap, share, tap } from 'rxjs/operators'; +import { filter, map, mergeMap, publishReplay, refCount, share, tap } from 'rxjs/operators'; import { Store } from '@ngrx/store'; import { AppState } from '@app/core/core.state'; import { TranslateService } from '@ngx-translate/core'; @@ -142,6 +142,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie searchText = ''; private latestSearchTextResult: Array = null; + private fetchObservable$: Observable> = null; private dirty = false; @@ -260,6 +261,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie if (!change.firstChange && change.currentValue !== change.previousValue) { if (propName === 'entityAliasId') { this.searchText = ''; + this.fetchObservable$ = null; this.latestSearchTextResult = null; this.dirty = true; } else if (['widgetType', 'datasourceType'].includes(propName)) { @@ -405,14 +407,24 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie return key ? key.name : undefined; } - fetchKeys(searchText?: string): Observable> { - if (this.latestSearchTextResult === null || this.searchText !== searchText) { + private fetchKeys(searchText?: string): Observable> { + if (this.searchText !== searchText || this.latestSearchTextResult === null) { this.searchText = searchText; - let fetchObservable: Observable> = null; + const dataKeyFilter = this.createDataKeyFilter(this.searchText); + return this.getKeys().pipe( + map(name => name.filter(dataKeyFilter)), + tap(res => this.latestSearchTextResult = res) + ); + } + return of(this.latestSearchTextResult); + } + + private getKeys(): Observable> { + if (this.fetchObservable$ === null) { + let fetchObservable: Observable>; if (this.datasourceType === DatasourceType.function) { - const dataKeyFilter = this.createDataKeyFilter(this.searchText); const targetKeysList = this.widgetType === widgetType.alarm ? this.alarmKeys : this.functionTypeKeys; - fetchObservable = of(targetKeysList.filter(dataKeyFilter)); + fetchObservable = of(targetKeysList); } else { if (this.entityAliasId) { const dataKeyTypes = [DataKeyType.timeseries]; @@ -420,24 +432,25 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie dataKeyTypes.push(DataKeyType.attribute); dataKeyTypes.push(DataKeyType.entityField); if (this.widgetType === widgetType.alarm) { - dataKeyTypes.push(DataKeyType.alarm); + dataKeyTypes.push(DataKeyType.alarm); } } - fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, this.searchText, dataKeyTypes); + fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, dataKeyTypes); } else { fetchObservable = of([]); } } - return fetchObservable.pipe( - tap(res => this.latestSearchTextResult = res) + this.fetchObservable$ = fetchObservable.pipe( + publishReplay(1), + refCount() ); } - return of(this.latestSearchTextResult); + return this.fetchObservable$; } private createDataKeyFilter(query: string): (key: DataKey) => boolean { const lowercaseQuery = query.toLowerCase(); - return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; + return key => key.name.toLowerCase().startsWith(lowercaseQuery); } textIsNotEmpty(text: string): boolean { diff --git a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts index ba8cdaf5d2..b1984f997a 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts @@ -54,13 +54,13 @@ import { UtilsService } from '@core/services/utils.service'; import { DataKeyType } from '@shared/models/telemetry/telemetry.models'; import { TranslateService } from '@ngx-translate/core'; import { EntityType } from '@shared/models/entity-type.models'; -import { forkJoin, Observable, of, Subscription } from 'rxjs'; +import { Observable, of, Subscription } from 'rxjs'; import { WidgetConfigCallbacks } from '@home/components/widget/widget-config.component.models'; import { EntityAliasDialogComponent, EntityAliasDialogData } from '@home/components/alias/entity-alias-dialog.component'; -import { catchError, map, mergeMap, tap } from 'rxjs/operators'; +import { catchError, mergeMap, tap } from 'rxjs/operators'; import { MatDialog } from '@angular/material/dialog'; import { EntityService } from '@core/http/entity.service'; import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models'; @@ -792,54 +792,16 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont ); } - private fetchEntityKeys(entityAliasId: string, query: string, dataKeyTypes: Array): Observable> { - return this.aliasController.resolveSingleEntityInfo(entityAliasId).pipe( - mergeMap((entity) => { - if (entity) { - const fetchEntityTasks: Array>> = []; - for (const dataKeyType of dataKeyTypes) { - fetchEntityTasks.push( - this.entityService.getEntityKeys( - {entityType: entity.entityType, id: entity.id}, - query, - dataKeyType, - {ignoreLoading: true, ignoreErrors: true} - ).pipe( - map((keys) => { - const dataKeys: Array = []; - for (const key of keys) { - dataKeys.push({name: key, type: dataKeyType}); - } - return dataKeys; - } - ), - catchError(() => of([])) - )); - } - return forkJoin(fetchEntityTasks).pipe( - map(arrayOfDataKeys => { - const result = new Array(); - arrayOfDataKeys.forEach((dataKeyArray) => { - result.push(...dataKeyArray); - }); - return result; - } - )); - } else if (dataKeyTypes.includes(DataKeyType.alarm)) { - return this.entityService.getEntityKeys(null, query, DataKeyType.alarm).pipe( - map((keys) => { - const dataKeys: Array = []; - for (const key of keys) { - dataKeys.push({name: key, type: DataKeyType.alarm}); - } - return dataKeys; - } - ), - catchError(() => of([])) - ); - } else { - return of([]); - } + private fetchEntityKeys(entityAliasId: string, dataKeyTypes: Array): Observable> { + return this.aliasController.getAliasInfo(entityAliasId).pipe( + mergeMap((aliasInfo) => { + return this.entityService.getEntityKeysByEntityFilter( + aliasInfo.entityFilter, + dataKeyTypes, + {ignoreLoading: true, ignoreErrors: true} + ).pipe( + catchError(() => of([])) + ); }), catchError(() => of([] as Array)) ); diff --git a/ui-ngx/src/app/shared/models/entity.models.ts b/ui-ngx/src/app/shared/models/entity.models.ts index 8f1d4f417e..300de484c1 100644 --- a/ui-ngx/src/app/shared/models/entity.models.ts +++ b/ui-ngx/src/app/shared/models/entity.models.ts @@ -64,6 +64,12 @@ export interface EntityField { time?: boolean; } +export interface EntitiesKeysByQuery { + attribute: Array; + timeseries: Array; + entityTypes: EntityType[]; +} + export const entityFields: {[fieldName: string]: EntityField} = { createdTime: { keyName: 'createdTime',