Browse Source

Merge with develop/3.4

pull/6654/head
Igor Kulikov 4 years ago
parent
commit
d9410e5330
  1. 41
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

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

@ -94,10 +94,6 @@ public class HashPartitionService implements PartitionService {
@PostConstruct @PostConstruct
public void init() { public void init() {
this.hashFunction = forName(hashFunctionName); this.hashFunction = forName(hashFunctionName);
}
@AfterStartUp(order = AfterStartUp.QUEUE_INFO_INITIALIZATION)
public void partitionsInit() {
QueueKey coreKey = new QueueKey(ServiceType.TB_CORE); QueueKey coreKey = new QueueKey(ServiceType.TB_CORE);
partitionSizesMap.put(coreKey, corePartitions); partitionSizesMap.put(coreKey, corePartitions);
partitionTopicsMap.put(coreKey, coreTopic); partitionTopicsMap.put(coreKey, coreTopic);
@ -106,12 +102,33 @@ public class HashPartitionService implements PartitionService {
partitionSizesMap.put(vcKey, vcPartitions); partitionSizesMap.put(vcKey, vcPartitions);
partitionTopicsMap.put(vcKey, vcTopic); partitionTopicsMap.put(vcKey, vcTopic);
List<QueueRoutingInfo> queueRoutingInfoList; if (!isTransport(serviceInfoProvider.getServiceType())) {
doInitRuleEnginePartitions();
}
}
String serviceType = serviceInfoProvider.getServiceType(); @AfterStartUp(order = AfterStartUp.QUEUE_INFO_INITIALIZATION)
public void partitionsInit() {
if (isTransport(serviceInfoProvider.getServiceType())) {
doInitRuleEnginePartitions();
}
}
private void doInitRuleEnginePartitions() {
List<QueueRoutingInfo> queueRoutingInfoList = getQueueRoutingInfos();
queueRoutingInfoList.forEach(queue -> {
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queue);
partitionTopicsMap.put(queueKey, queue.getQueueTopic());
partitionSizesMap.put(queueKey, queue.getPartitions());
queuesById.put(queue.getQueueId(), queue);
});
}
private List<QueueRoutingInfo> getQueueRoutingInfos() {
List<QueueRoutingInfo> queueRoutingInfoList;
String serviceType = serviceInfoProvider.getServiceType();
if ("tb-transport".equals(serviceType)) { if (isTransport(serviceType)) {
//If transport started earlier than tb-core //If transport started earlier than tb-core
int getQueuesRetries = 10; int getQueuesRetries = 10;
while (true) { while (true) {
@ -136,13 +153,11 @@ public class HashPartitionService implements PartitionService {
} else { } else {
queueRoutingInfoList = queueRoutingInfoService.getAllQueuesRoutingInfo(); queueRoutingInfoList = queueRoutingInfoService.getAllQueuesRoutingInfo();
} }
return queueRoutingInfoList;
}
queueRoutingInfoList.forEach(queue -> { private boolean isTransport(String serviceType) {
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queue); return "tb-transport".equals(serviceType);
partitionTopicsMap.put(queueKey, queue.getQueueTopic());
partitionSizesMap.put(queueKey, queue.getPartitions());
queuesById.put(queue.getQueueId(), queue);
});
} }
@Override @Override

Loading…
Cancel
Save