From 44775f6f1c063462854d3054f5b65987a6fc19f6 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 12 Mar 2024 12:41:41 +0200 Subject: [PATCH] DataQuery tests --- .../server/dao/edq/DeviceRepoData.java | 22 ++++ .../server/dao/edq/EntityDataQueryTest.java | 103 ++++++++++++++++++ .../server/dao/edq/InMemoryRepository.java | 56 ++++++++++ .../thingsboard/server/dao/edq/SortPair.java | 14 +++ 4 files changed, 195 insertions(+) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/edq/DeviceRepoData.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/edq/EntityDataQueryTest.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/edq/InMemoryRepository.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/edq/SortPair.java diff --git a/dao/src/main/java/org/thingsboard/server/dao/edq/DeviceRepoData.java b/dao/src/main/java/org/thingsboard/server/dao/edq/DeviceRepoData.java new file mode 100644 index 0000000000..0fdd8bfe5f --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edq/DeviceRepoData.java @@ -0,0 +1,22 @@ +package org.thingsboard.server.dao.edq; + +import lombok.Getter; +import org.thingsboard.server.common.data.Device; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +public class DeviceRepoData { + @Getter + private final Device device; + @Getter + private final Map attrs = new ConcurrentHashMap<>(); + + public DeviceRepoData(Device device) { + this.device = device; + } + + public void putAttr(Integer keyId, String value){ + attrs.put(keyId, value); + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edq/EntityDataQueryTest.java b/dao/src/main/java/org/thingsboard/server/dao/edq/EntityDataQueryTest.java new file mode 100644 index 0000000000..4f286749f6 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edq/EntityDataQueryTest.java @@ -0,0 +1,103 @@ +package org.thingsboard.server.dao.edq; + +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.id.DeviceId; + +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.UUID; +import java.util.function.Function; +import java.util.function.Predicate; + +@Slf4j +public class EntityDataQueryTest { + + private static final int DEVICE_COUNT = 500000; + private static final int ATTR_COUNT = 100; + private static final Comparator CREATED_TIME_COMPARATOR = Comparator.comparingLong(o -> o.getDevice().getCreatedTime()); + + private static final Comparator ATTR_COMPARATOR = Comparator.comparing(o -> o.getAttrs().get(10)); + + private static final Comparator SORT_ASC = Comparator.comparing(SortPair::getKey); + private static final Comparator SORT_DESC = Comparator.comparing(SortPair::getKey).reversed(); + + public static void main(String[] args) { + InMemoryRepository repository = new InMemoryRepository(); + long startTs = System.currentTimeMillis(); + long ts = System.currentTimeMillis() - DEVICE_COUNT; + for (int i = 0; i < DEVICE_COUNT; i++) { + DeviceId deviceId = new DeviceId(UUID.randomUUID()); + Device device = new Device(); + device.setId(deviceId); + device.setCreatedTime(ts + i); + device.setName("Device " + i); + device.setLabel("Device Label" + i); + device.setType("Device Type " + (i % 100)); + repository.add(device); + String random = StringUtils.randomAlphanumeric(5); + for (int j = 0; j < ATTR_COUNT; j++) { + repository.add(deviceId, j, random); + } + } + log.error("Repository created in {}", System.currentTimeMillis() - startTs); + + + for (int i = 0; i < 1; i++) { + test("DeviceName filter, sort by attribute desc", repository, 10, d -> d.getDevice().getName().startsWith("Device 9"), ATTR_COMPARATOR.reversed()); + test2("DeviceName filter, sort by attribute desc", repository, 10, d -> d.getDevice().getName().startsWith("Device 8"), d -> d.getAttrs().get(10), SORT_DESC); + + test("DeviceType filter, no sort", repository, 3, d -> d.getDevice().getType().equals("Device Type 10"), null); + test2("DeviceType filter, no sort", repository, 3, d -> d.getDevice().getType().equals("Device Type 10"), null, null); + + test("Attribute filter, no sort", repository, 3, d -> d.getAttrs().get(10).contains("1"), null); + test2("Attribute filter, no sort", repository, 3, d -> d.getAttrs().get(9).contains("2"), null, null); + + test("No filter, sort by attribute", repository, 3, null, ATTR_COMPARATOR); + test2("No filter, sort by attribute", repository, 3, null, d -> d.getAttrs().get(10), SORT_ASC); + + test("Attribute filter, createdTime", repository, 3, d -> d.getAttrs().get(10).contains("1"), CREATED_TIME_COMPARATOR); + test("No filter, sort by createdTime", repository, 3, null, CREATED_TIME_COMPARATOR); + } + } + + private static void test(String testName, InMemoryRepository repository, int pages, Predicate predicate, Comparator comparator) { + log.error("==================================================================================================="); + log.error("============================ 1 = {} =======================================", testName); + List execTimes = new ArrayList<>(); + for (int i = 0; i < pages; i++) { + long startTs = System.nanoTime(); + var results = repository.findDevices(i, 10, predicate, comparator); + if (results.isEmpty()) { + log.error("Empty results"); + } + long endTs = System.nanoTime(); + execTimes.add(endTs - startTs); + } + log.error("Run queries in {}", execTimes); + log.error("Average {}", execTimes.stream().mapToDouble(v -> (double) v).average().orElse(0.0) / 1000000.0); + log.error("==================================================================================================="); + } + + private static void test2(String testName, InMemoryRepository repository, int pages, Predicate predicate, Function toSortKey, Comparator comparator) { + log.error("==================================================================================================="); + log.error("=========================== 2 = {} =======================================", testName); + List execTimes = new ArrayList<>(); + for (int i = 0; i < pages; i++) { + long startTs = System.nanoTime(); + var results = repository.findDevices(i, 10, predicate, toSortKey, comparator); + if (results.isEmpty()) { + log.error("Empty results"); + } + long endTs = System.nanoTime(); + execTimes.add(endTs - startTs); + } + log.error("Run queries in {}", execTimes); + log.error("Average {}", execTimes.stream().mapToDouble(v -> (double) v).average().orElse(0.0) / 1000000.0); + log.error("==================================================================================================="); + } + + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edq/InMemoryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/edq/InMemoryRepository.java new file mode 100644 index 0000000000..f32b5b86ae --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edq/InMemoryRepository.java @@ -0,0 +1,56 @@ +package org.thingsboard.server.dao.edq; + +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.id.DeviceId; + +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ConcurrentSkipListSet; +import java.util.function.Function; +import java.util.function.Predicate; +import java.util.stream.Collectors; + +public class InMemoryRepository { + + private final static ConcurrentSkipListSet devices = new ConcurrentSkipListSet<>(Comparator.comparingLong(o -> o.getDevice().getCreatedTime())); + private final static ConcurrentMap devicesById = new ConcurrentHashMap<>(); + + public void add(Device device) { + var dd = new DeviceRepoData(device); + devicesById.put(dd.getDevice().getId(), dd); + devices.add(dd); + } + + public void add(DeviceId deviceId, Integer attrKey, String attrValue) { + devicesById.get(deviceId).putAttr(attrKey, attrValue); + } + + public List findDevices(int page, int pageSize, Predicate predicate, Comparator comparator) { + var stream = devices.stream(); + if (predicate != null) { + stream = stream.filter(predicate); + } + if (comparator != null) { + stream = stream.sorted(comparator); + } + return stream.skip(page * pageSize).limit(pageSize).collect(Collectors.toList()); + } + + public List findDevices(int page, int pageSize, Predicate predicate, Function sortFieldFunction, Comparator comparator) { + var stream = devices.stream(); + if (predicate != null) { + stream = stream.filter(predicate); + } + + if (comparator != null) { + var sortedStream = stream.map(d -> new SortPair(sortFieldFunction.apply(d), d)); + sortedStream = sortedStream.sorted(comparator); + return sortedStream.skip(page * pageSize).limit(pageSize).map(v -> v.data).collect(Collectors.toList()); + } else { + return stream.skip(page * pageSize).limit(pageSize).collect(Collectors.toList()); + } + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edq/SortPair.java b/dao/src/main/java/org/thingsboard/server/dao/edq/SortPair.java new file mode 100644 index 0000000000..3323f4c829 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edq/SortPair.java @@ -0,0 +1,14 @@ +package org.thingsboard.server.dao.edq; + +import lombok.Getter; +import lombok.RequiredArgsConstructor; + +@RequiredArgsConstructor +public class SortPair { + + @Getter + final String key; + @Getter + final DeviceRepoData data; + +}