From 2fced834fdde91773dee94366b9fdb7ad6a789ae Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 4 Apr 2024 14:25:08 +0300 Subject: [PATCH] Improvement to entity data query: add optimization, add fetching by 1024 max --- .../server/dao/entity/BaseEntityService.java | 121 ++++++++++++------ .../server/dao/service/EntityServiceTest.java | 116 ++++++++++++++++- 2 files changed, 199 insertions(+), 38 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java index 5edb59f5c8..59475213b3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java @@ -40,11 +40,16 @@ import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityFilterType; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityListFilter; +import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.RelationsQueryFilter; import org.thingsboard.server.dao.exception.IncorrectParameterException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Optional; +import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -63,6 +68,10 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId "; public static final CustomerId NULL_CUSTOMER_ID = new CustomerId(NULL_UUID); + private static final int MAX_ENTITY_IDS_SIZE = 1024; + private static final Set EXCLUDED_TYPES_FROM_OPTIMIZATION = Set.of( + EntityFilterType.ENTITY_LIST, EntityFilterType.SINGLE_ENTITY, EntityFilterType.RELATIONS_QUERY); + @Autowired private EntityQueryDao entityQueryDao; @@ -86,9 +95,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe validateId(customerId, id -> INCORRECT_CUSTOMER_ID + id); validateEntityDataQuery(query); - if (EntityFilterType.RELATIONS_QUERY.equals(query.getEntityFilter().getType()) - || EntityFilterType.SINGLE_ENTITY.equals(query.getEntityFilter().getType()) - || StringUtils.isNotEmpty(query.getPageLink().getTextSearch())) { + if (isOptimizationExcluded(query)) { return this.entityQueryDao.findEntityDataByQuery(tenantId, customerId, query); } @@ -97,41 +104,10 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe if (entityDataByQuery == null || entityDataByQuery.getData().isEmpty()) { return entityDataByQuery; } - // 2 step - find entity data by entity ids from the 1st step - PageData result = findEntityDataByEntityIds(tenantId, customerId, query, entityDataByQuery.getData()); - return new PageData<>(result.getData(), entityDataByQuery.getTotalPages(), entityDataByQuery.getTotalElements(), entityDataByQuery.hasNext()); - } - - private PageData findEntityIdsByFilterAndSorterColumns(TenantId tenantId, CustomerId customerId, EntityDataQuery query) { - List entityFields = null; - List latestValues = null; - if (query.getPageLink().getSortOrder() != null) { - if (query.getEntityFields() != null) { - entityFields = query.getEntityFields().stream() - .filter(entityKey -> entityKey.getKey().equals(query.getPageLink().getSortOrder().getKey().getKey())) - .collect(Collectors.toList()); - } - if (query.getLatestValues() != null) { - latestValues = query.getLatestValues().stream() - .filter(entityKey -> entityKey.getKey().equals(query.getPageLink().getSortOrder().getKey().getKey())) - .collect(Collectors.toList()); - } - } - EntityDataQuery entityQuery = new EntityDataQuery(query.getEntityFilter(), query.getPageLink(), entityFields, latestValues, query.getKeyFilters()); - return this.entityQueryDao.findEntityDataByQuery(tenantId, customerId, entityQuery); - } - - private PageData findEntityDataByEntityIds(TenantId tenantId, CustomerId customerId, EntityDataQuery query, List data) { - List entityIds = data.stream().map(d -> d.getEntityId().getId().toString()).toList(); - EntityType entityType = data.isEmpty() ? null : data.get(0).getEntityId().getEntityType(); - - EntityListFilter filter = new EntityListFilter(); - filter.setEntityType(entityType); - filter.setEntityList(entityIds); - EntityDataPageLink pageLink = new EntityDataPageLink(query.getPageLink().getPageSize(), 0, null, query.getPageLink().getSortOrder()); - EntityDataQuery entityQuery = new EntityDataQuery(filter, pageLink, query.getEntityFields(), query.getLatestValues(), null); - return this.entityQueryDao.findEntityDataByQuery(tenantId, customerId, entityQuery); + // 2 step - find entity data by entity ids from the 1st step + List result = fetchEntityDataByIdsFromInitialQuery(tenantId, customerId, query, entityDataByQuery.getData()); + return new PageData<>(result, entityDataByQuery.getTotalPages(), entityDataByQuery.getTotalElements(), entityDataByQuery.hasNext()); } @Override @@ -228,4 +204,75 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } } + private boolean isOptimizationExcluded(EntityDataQuery query) { + if (StringUtils.isNotEmpty(query.getPageLink().getTextSearch())) { + return true; + } + + if (EXCLUDED_TYPES_FROM_OPTIMIZATION.contains(query.getEntityFilter().getType())) { + return true; + } + + if ((query.getEntityFields() == null || query.getEntityFields().isEmpty()) && + (query.getLatestValues() == null || query.getLatestValues().isEmpty())) { + return true; + } + + Set entityKeys = new HashSet<>(Optional.ofNullable(query.getKeyFilters()).orElse(Collections.emptyList()).stream().map(KeyFilter::getKey).toList()); + Set entityFields = new HashSet<>(Optional.ofNullable(query.getEntityFields()).orElse(Collections.emptyList())); + Set latestValues = new HashSet<>(Optional.ofNullable(query.getLatestValues()).orElse(Collections.emptyList())); + + return entityKeys.equals(entityFields) && entityKeys.equals(latestValues); + } + + private PageData findEntityIdsByFilterAndSorterColumns(TenantId tenantId, CustomerId customerId, EntityDataQuery query) { + List entityFields = null; + List latestValues = null; + if (query.getPageLink().getSortOrder() != null) { + if (query.getEntityFields() != null) { + entityFields = query.getEntityFields().stream() + .filter(entityKey -> entityKey.getKey().equals(query.getPageLink().getSortOrder().getKey().getKey())) + .collect(Collectors.toList()); + } + if (query.getLatestValues() != null) { + latestValues = query.getLatestValues().stream() + .filter(entityKey -> entityKey.getKey().equals(query.getPageLink().getSortOrder().getKey().getKey())) + .collect(Collectors.toList()); + } + } + EntityDataQuery entityQuery = new EntityDataQuery(query.getEntityFilter(), query.getPageLink(), entityFields, latestValues, query.getKeyFilters()); + return this.entityQueryDao.findEntityDataByQuery(tenantId, customerId, entityQuery); + } + + private List fetchEntityDataByIdsFromInitialQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query, List initialQueryResult) { + List result = new ArrayList<>(); + + List entityIds = initialQueryResult.stream().map(d -> d.getEntityId().getId().toString()).collect(Collectors.toList()); + EntityType entityType = initialQueryResult.get(0).getEntityId().getEntityType(); + + if (entityIds.size() > MAX_ENTITY_IDS_SIZE) { + List> chunks = new ArrayList<>(); + for (int i = 0; i < entityIds.size(); i += MAX_ENTITY_IDS_SIZE) { + chunks.add(entityIds.subList(i, Math.min(entityIds.size(), i + MAX_ENTITY_IDS_SIZE))); + } + for (List chunk : chunks) { + result.addAll(findEntityDataByEntityIds(tenantId, customerId, query, chunk, entityType, chunk.size())); + } + } else { + result.addAll(findEntityDataByEntityIds(tenantId, customerId, query, entityIds, entityType, query.getPageLink().getPageSize())); + } + return result; + } + + private List findEntityDataByEntityIds(TenantId tenantId, CustomerId customerId, EntityDataQuery query, + List entityIds, EntityType entityType, int pageSize) { + EntityListFilter filter = new EntityListFilter(); + filter.setEntityType(entityType); + filter.setEntityList(entityIds); + + EntityDataPageLink pageLink = new EntityDataPageLink(pageSize, 0, null, query.getPageLink().getSortOrder()); + EntityDataQuery entityQuery = new EntityDataQuery(filter, pageLink, query.getEntityFields(), query.getLatestValues(), null); + return this.entityQueryDao.findEntityDataByQuery(tenantId, customerId, entityQuery).getData(); + } + } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/EntityServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/EntityServiceTest.java index 216d2abb07..58ee94f8b7 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/EntityServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/EntityServiceTest.java @@ -26,7 +26,6 @@ import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.jdbc.core.ResultSetExtractor; import org.thingsboard.server.common.data.AttributeScope; -import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; @@ -99,6 +98,7 @@ import java.util.stream.Stream; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertEquals; +import static org.thingsboard.server.common.data.query.EntityKeyType.ENTITY_FIELD; @Slf4j @DaoSqlTest @@ -2113,6 +2113,120 @@ public class EntityServiceTest extends AbstractServiceTest { deviceService.deleteDevicesByTenantId(tenantId); } + @Test + public void testFindEntityQueryWith_3000_pageSize() throws InterruptedException { + int pageSize = 3000; + + List devices = new ArrayList<>(); + + for (int i = 0; i < pageSize; i++) { + Device device = new Device(); + device.setTenantId(tenantId); + device.setName("Device_" + i); + device.setType("default"); + device.setLabel("testLabel" + (int) (Math.random() * 1000)); + devices.add(deviceService.saveDevice(device)); + //TO make sure devices have different created time + Thread.sleep(1); + } + + DeviceTypeFilter filter = new DeviceTypeFilter(); + filter.setDeviceTypes(List.of("default")); + filter.setDeviceNameFilter("D%"); + + EntityDataSortOrder sortOrder = new EntityDataSortOrder(new EntityKey(ENTITY_FIELD, "name"), EntityDataSortOrder.Direction.DESC); + + List deviceTypeFilters = createStringKeyFilters("type", ENTITY_FIELD, StringFilterPredicate.StringOperation.EQUAL, "default"); + + KeyFilter createdTimeFilter = createNumericKeyFilter("createdTime", ENTITY_FIELD, NumericFilterPredicate.NumericOperation.GREATER, 1L); + List createdTimeFilters = Collections.singletonList(createdTimeFilter); + + List nameFilters = createStringKeyFilters("name", ENTITY_FIELD, StringFilterPredicate.StringOperation.CONTAINS, "Device"); + + List entityFields = Arrays.asList(new EntityKey(ENTITY_FIELD, "name"), + new EntityKey(ENTITY_FIELD, "type")); + + // 1. Device type filters: + + // query with textSearch - optimization is not performing + EntityDataPageLink originalPageLink = new EntityDataPageLink(pageSize, 0, "Device", sortOrder); + EntityDataQuery originalQuery = new EntityDataQuery(filter, originalPageLink, entityFields, null, deviceTypeFilters); + PageData originalData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), originalQuery); + + // query without textSearch - optimization is performing + EntityDataPageLink optimizedPageLink = new EntityDataPageLink(pageSize, 0, null, sortOrder); + EntityDataQuery optimizedQuery = new EntityDataQuery(filter, optimizedPageLink, entityFields, null, deviceTypeFilters); + PageData optimizedData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), optimizedQuery); + List loadedEntities = getLoadedEntities(optimizedData, optimizedQuery); + Assert.assertEquals(devices.size(), loadedEntities.size()); + + for (int i = 0; i < devices.size(); i++) { + var originalElement = originalData.getData().get(i); + var optimizedElement = optimizedData.getData().get(i); + Assert.assertEquals(originalElement.getEntityId(), optimizedElement.getEntityId()); + originalElement.getLatest().get(ENTITY_FIELD).forEach((key, value) -> { + Assert.assertEquals(value.getValue(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getValue()); + Assert.assertEquals(value.getCount(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getCount()); + }); + } + Assert.assertEquals(originalData.getTotalPages(), optimizedData.getTotalPages()); + Assert.assertEquals(originalData.getTotalElements(), optimizedData.getTotalElements()); + + // 2. Device create time filters + + // query with textSearch - optimization is not performing + originalPageLink = new EntityDataPageLink(pageSize, 0, "Device", sortOrder); + originalQuery = new EntityDataQuery(filter, originalPageLink, entityFields, null, createdTimeFilters); + originalData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), originalQuery); + + // query without textSearch - optimization is performing + optimizedPageLink = new EntityDataPageLink(pageSize, 0, null, sortOrder); + optimizedQuery = new EntityDataQuery(filter, optimizedPageLink, entityFields, null, createdTimeFilters); + optimizedData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), optimizedQuery); + loadedEntities = getLoadedEntities(optimizedData, optimizedQuery); + Assert.assertEquals(devices.size(), loadedEntities.size()); + + for (int i = 0; i < devices.size(); i++) { + var originalElement = originalData.getData().get(i); + var optimizedElement = optimizedData.getData().get(i); + Assert.assertEquals(originalElement.getEntityId(), optimizedElement.getEntityId()); + originalElement.getLatest().get(ENTITY_FIELD).forEach((key, value) -> { + Assert.assertEquals(value.getValue(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getValue()); + Assert.assertEquals(value.getCount(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getCount()); + }); + } + Assert.assertEquals(originalData.getTotalPages(), optimizedData.getTotalPages()); + Assert.assertEquals(originalData.getTotalElements(), optimizedData.getTotalElements()); + + // 3. Device name filters + + // query with textSearch - optimization is not performing + originalPageLink = new EntityDataPageLink(pageSize, 0, "Device", sortOrder); + originalQuery = new EntityDataQuery(filter, originalPageLink, entityFields, null, nameFilters); + originalData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), originalQuery); + + // query without textSearch - optimization is performing + optimizedPageLink = new EntityDataPageLink(pageSize, 0, null, sortOrder); + optimizedQuery = new EntityDataQuery(filter, optimizedPageLink, entityFields, null, nameFilters); + optimizedData = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), optimizedQuery); + loadedEntities = getLoadedEntities(optimizedData, optimizedQuery); + Assert.assertEquals(devices.size(), loadedEntities.size()); + + for (int i = 0; i < devices.size(); i++) { + var originalElement = originalData.getData().get(i); + var optimizedElement = optimizedData.getData().get(i); + Assert.assertEquals(originalElement.getEntityId(), optimizedElement.getEntityId()); + originalElement.getLatest().get(ENTITY_FIELD).forEach((key, value) -> { + Assert.assertEquals(value.getValue(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getValue()); + Assert.assertEquals(value.getCount(), optimizedElement.getLatest().get(EntityKeyType.ENTITY_FIELD).get(key).getCount()); + }); + } + Assert.assertEquals(originalData.getTotalPages(), optimizedData.getTotalPages()); + Assert.assertEquals(originalData.getTotalElements(), optimizedData.getTotalElements()); + + deviceService.deleteDevicesByTenantId(tenantId); + } + private Boolean listEqualWithoutOrder(List A, List B) { return A.containsAll(B) && B.containsAll(A); }