diff --git a/application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java similarity index 66% rename from application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java rename to application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java index 9df6c16111..e20f40be46 100644 --- a/application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java @@ -13,9 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.cluster.routing; +package org.thingsboard.server.queue.discovery; +import com.datastax.driver.core.utils.UUIDs; import com.datastax.oss.driver.api.core.uuid.Uuids; +import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.junit.Assert; import org.junit.Before; @@ -29,17 +31,16 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.HashPartitionService; -import org.thingsboard.server.queue.discovery.QueueRoutingInfoService; -import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; -import org.thingsboard.server.queue.discovery.TenantRoutingInfoService; +import java.text.SimpleDateFormat; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Random; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static org.mockito.Mockito.mock; @@ -111,15 +112,56 @@ public class HashPartitionServiceTest { map.put(partition, map.getOrDefault(partition, 0) + 1); } - List> data = map.entrySet().stream().sorted(Comparator.comparingInt(Map.Entry::getValue)).collect(Collectors.toList()); + checkDispersion(start, map, ITERATIONS, 1.0); + } + + @SneakyThrows + @Test + public void testDispersionOnResolveByPartitionIdx() { + int serverCount = 5; + int tenantCount = 1000; + int queueCount = 3; + int partitionCount = 3; + + List services = new ArrayList<>(); + + for (int i = 0; i < serverCount; i++) { + services.add(TransportProtos.ServiceInfo.newBuilder().setServiceId("RE-" + i).build()); + } + + long start = System.currentTimeMillis(); + Map map = new HashMap<>(); + services.forEach(s -> map.put(s.getServiceId(), 0)); + + Random random = new Random(); + long ts = new SimpleDateFormat("dd-MM-yyyy").parse("06-12-2016").getTime() - TimeUnit.DAYS.toMillis(tenantCount); + for (int tenantIndex = 0; tenantIndex < tenantCount; tenantIndex++) { + TenantId tenantId = new TenantId(UUIDs.startOf(ts)); + ts += TimeUnit.DAYS.toMillis(1) + random.nextInt(1000); + for (int queueIndex = 0; queueIndex < queueCount; queueIndex++) { + QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, "queue" + queueIndex, tenantId); + for (int partition = 0; partition < partitionCount; partition++) { + TransportProtos.ServiceInfo serviceInfo = clusterRoutingService.resolveByPartitionIdx(services, queueKey, partition); + String serviceId = serviceInfo.getServiceId(); + map.put(serviceId, map.get(serviceId) + 1); + } + } + } + + checkDispersion(start, map, tenantCount * queueCount * partitionCount, 10.0); + } + + private void checkDispersion(long start, Map map, int iterations, double maxDiffPercent) { + List> data = map.entrySet().stream().sorted(Comparator.comparingInt(Map.Entry::getValue)).collect(Collectors.toList()); long end = System.currentTimeMillis(); - double diff = (data.get(data.size() - 1).getValue() - data.get(0).getValue()); - double diffPercent = (diff / ITERATIONS) * 100.0; + double ideal = ((double) iterations) / map.size(); + double diff = Math.max(data.get(data.size() - 1).getValue() - ideal, ideal - data.get(0).getValue()); + double diffPercent = (diff / ideal) * 100.0; System.out.println("Time: " + (end - start) + " Diff: " + diff + "(" + String.format("%f", diffPercent) + "%)"); - Assert.assertTrue(diffPercent < 0.5); - for (Map.Entry entry : data) { + for (Map.Entry entry : data) { System.out.println(entry.getKey() + ": " + entry.getValue()); } + Assert.assertTrue(diffPercent < maxDiffPercent); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index afbd96625d..809c08ac6f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -34,6 +34,7 @@ import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent; import org.thingsboard.server.queue.util.AfterStartUp; import javax.annotation.PostConstruct; +import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; @@ -254,7 +255,7 @@ public class HashPartitionService implements PartitionService { myPartitions = new ConcurrentHashMap<>(); partitionSizesMap.forEach((queueKey, size) -> { for (int i = 0; i < size; i++) { - ServiceInfo serviceInfo = resolveByPartitionIdx(queueServicesMap.get(queueKey), i); + ServiceInfo serviceInfo = resolveByPartitionIdx(queueServicesMap.get(queueKey), queueKey, i); if (currentService.equals(serviceInfo)) { myPartitions.computeIfAbsent(queueKey, key -> new ArrayList<>()).add(i); } @@ -434,11 +435,21 @@ public class HashPartitionService implements PartitionService { } } - private ServiceInfo resolveByPartitionIdx(List servers, Integer partitionIdx) { + protected ServiceInfo resolveByPartitionIdx(List servers, QueueKey queueKey, int partition) { if (servers == null || servers.isEmpty()) { return null; } - return servers.get(partitionIdx % servers.size()); + + if (!ServiceType.TB_RULE_ENGINE.equals(queueKey.getType()) || TenantId.SYS_TENANT_ID.equals(queueKey.getTenantId())) { + return servers.get(partition % servers.size()); + } else { + int hash = hashFunction.newHasher().putLong(queueKey.getTenantId().getId().getMostSignificantBits()) + .putLong(queueKey.getTenantId().getId().getLeastSignificantBits()) + .putString(queueKey.getQueueName(), StandardCharsets.UTF_8) + .hash().asInt(); + + return servers.get(Math.abs((hash + partition) % servers.size())); + } } public static HashFunction forName(String name) {