Browse Source

Merge pull request #9513 from thingsboard/fix/queue-name-hash-partition

Distribute partitions between Rule Engines based only on tenant id
pull/9558/head
Andrew Shvayka 3 years ago
committed by GitHub
parent
commit
b6581b0a61
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 30
      application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java
  2. 6
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

30
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.SneakyThrows;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections4.ListUtils; import org.apache.commons.collections4.ListUtils;
import org.apache.commons.lang3.RandomStringUtils;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
import org.junit.Test; import org.junit.Test;
@ -365,6 +366,35 @@ public class HashPartitionServiceTest {
assertThat(clusterRoutingService.isManagedByCurrentService(regularTenantId)).isTrue(); assertThat(clusterRoutingService.isManagedByCurrentService(regularTenantId)).isTrue();
} }
@Test
public void testPartitionsDistribution_sameTenantDifferentQueues() {
List<ServiceInfo> 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<QueueKey> 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<PartitionChangeEvent> predicate) { private void verifyPartitionChangeEvent(Predicate<PartitionChangeEvent> predicate) {
verify(applicationEventPublisher).publishEvent(argThat(event -> event instanceof PartitionChangeEvent && predicate.test((PartitionChangeEvent) event))); verify(applicationEventPublisher).publishEvent(argThat(event -> event instanceof PartitionChangeEvent && predicate.test((PartitionChangeEvent) event)));
} }

6
common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

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

Loading…
Cancel
Save