From f63b4b1f7c3b843b391cb244307538400439e9e4 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 15 Dec 2020 16:16:02 +0200 Subject: [PATCH] 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); + }