From 27e54c3ae07d355299ec082967e72e338f59b953 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 26 Mar 2025 16:36:03 +0200 Subject: [PATCH] fixed tenant cache computing when partition change event --- ...aultCalculatedFieldEntityProfileCache.java | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/cache/DefaultCalculatedFieldEntityProfileCache.java b/application/src/main/java/org/thingsboard/server/service/cf/cache/DefaultCalculatedFieldEntityProfileCache.java index 2f5772ae50..8a0c67ce29 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/cache/DefaultCalculatedFieldEntityProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/cache/DefaultCalculatedFieldEntityProfileCache.java @@ -24,14 +24,11 @@ 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.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.util.TbRuleEngineComponent; import java.util.Collection; -import java.util.Collections; -import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.stream.Collectors; @@ -46,15 +43,21 @@ public class DefaultCalculatedFieldEntityProfileCache extends TbApplicationEvent private static final Integer UNKNOWN = 0; private final ConcurrentMap tenantCache = new ConcurrentHashMap<>(); private final PartitionService partitionService; - private volatile List myPartitions = Collections.emptyList(); @Override protected void onTbApplicationEvent(PartitionChangeEvent event) { - myPartitions = event.getCfPartitions().stream() + var tenantPartitions = event.getCfPartitions().stream() .filter(TopicPartitionInfo::isMyPartition) - .map(tpi -> tpi.getPartition().orElse(UNKNOWN)).collect(Collectors.toList()); - //Naive approach that need to be improved. - tenantCache.values().forEach(cache -> cache.setMyPartitions(myPartitions)); + .filter(tpi -> tpi.getTenantId().isPresent()) + .collect(Collectors.groupingBy( + tpi -> tpi.getTenantId().get(), + Collectors.mapping(tpi -> tpi.getPartition().orElse(UNKNOWN), Collectors.toList()) + )); + + tenantPartitions.forEach((tenantId, partitions) -> { + var cache = tenantCache.computeIfAbsent(tenantId, id -> new TenantEntityProfileCache()); + cache.setMyPartitions(partitions); + }); } @Override