From 9e408cd6f3a59ff03a3609aa7f04af0052e21ba7 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 27 Oct 2023 14:43:12 +0300 Subject: [PATCH 1/2] Distribute partitions between Rule Engines based only on tenant id --- .../discovery/HashPartitionServiceTest.java | 30 +++++++++++++++++++ .../queue/discovery/HashPartitionService.java | 1 - 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java index ef87794d96..64f902b8d3 100644 --- a/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java @@ -19,6 +19,7 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.ListUtils; +import org.apache.commons.lang3.RandomStringUtils; import org.junit.Assert; import org.junit.Before; import org.junit.Test; @@ -365,6 +366,35 @@ public class HashPartitionServiceTest { assertThat(clusterRoutingService.isManagedByCurrentService(regularTenantId)).isTrue(); } + @Test + public void testPartitionsDistribution_sameTenantDifferentQueues() { + List ruleEngines = new ArrayList<>(); + int serviceId = 0; + for (int i = 0; i < 5; i++) { + ServiceInfo commonServer = ServiceInfo.newBuilder() + .setServiceId("tb-rule-engine-" + serviceId) + .addAllServiceTypes(List.of(ServiceType.TB_RULE_ENGINE.name())) + .build(); + ruleEngines.add(commonServer); + serviceId++; + } + + Stream.concat(Stream.of(TenantId.SYS_TENANT_ID), Stream.generate(UUID::randomUUID).map(TenantId::new).limit(10)).forEach(tenantId -> { + List queues = Stream.generate(() -> RandomStringUtils.randomAlphabetic(10)) + .map(queueName -> new QueueKey(ServiceType.TB_RULE_ENGINE, queueName, tenantId)) + .limit(100).collect(Collectors.toList()); + + for (int partition = 0; partition < 10; partition++) { + ServiceInfo expectedAssignedRuleEngine = clusterRoutingService.resolveByPartitionIdx(ruleEngines, new QueueKey(ServiceType.TB_RULE_ENGINE, tenantId), partition); + for (QueueKey queueKey : queues) { + ServiceInfo assignedRuleEngine = clusterRoutingService.resolveByPartitionIdx(ruleEngines, queueKey, partition); + assertThat(assignedRuleEngine).as(queueKey + "[" + partition + "] should be assigned to " + expectedAssignedRuleEngine.getServiceId()) + .isEqualTo(expectedAssignedRuleEngine); + } + } + }); + } + private void verifyPartitionChangeEvent(Predicate predicate) { verify(applicationEventPublisher).publishEvent(argThat(event -> event instanceof PartitionChangeEvent && predicate.test((PartitionChangeEvent) event))); } 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 5be5caf6ff..3de2e9822a 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 @@ -529,7 +529,6 @@ public class HashPartitionService implements PartitionService { int hash = hashFunction.newHasher() .putLong(tenantId.getId().getMostSignificantBits()) .putLong(tenantId.getId().getLeastSignificantBits()) - .putString(queueKey.getQueueName(), StandardCharsets.UTF_8) .hash().asInt(); return servers.get(Math.abs((hash + partition) % servers.size())); } else { From 0812149524ffbf033153dfe0c7e2864ec8561506 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 27 Oct 2023 16:38:44 +0300 Subject: [PATCH 2/2] Minor refactoring for HashPartitionService --- .../server/queue/discovery/HashPartitionService.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) 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 3de2e9822a..e43713147f 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 @@ -526,10 +526,7 @@ public class HashPartitionService implements PartitionService { servers = responsible; } - int hash = hashFunction.newHasher() - .putLong(tenantId.getId().getMostSignificantBits()) - .putLong(tenantId.getId().getLeastSignificantBits()) - .hash().asInt(); + int hash = hash(tenantId.getId()); return servers.get(Math.abs((hash + partition) % servers.size())); } else { return servers.get(partition % servers.size());