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..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,11 +526,7 @@ public class HashPartitionService implements PartitionService { servers = responsible; } - int hash = hashFunction.newHasher() - .putLong(tenantId.getId().getMostSignificantBits()) - .putLong(tenantId.getId().getLeastSignificantBits()) - .putString(queueKey.getQueueName(), StandardCharsets.UTF_8) - .hash().asInt(); + int hash = hash(tenantId.getId()); return servers.get(Math.abs((hash + partition) % servers.size())); } else { return servers.get(partition % servers.size());