From f651f0a721c39fc6e1be59425aeff0fa206ccfa1 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Thu, 19 Mar 2026 17:27:32 +0200 Subject: [PATCH 1/2] fixed cache cleanup on tenant/entity deletion for DefaultCalculatedFieldCache, DefaultTbAssetProfileCache, DefaultTbDeviceProfileCache --- .../cf/DefaultCalculatedFieldCache.java | 47 +++++++++++++++++++ .../profile/DefaultTbAssetProfileCache.java | 30 ++++++++++++ .../profile/DefaultTbDeviceProfileCache.java | 30 ++++++++++++ 3 files changed, 107 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java index 1eba2c4549..57ce166014 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -19,11 +19,14 @@ import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.springframework.util.ConcurrentReferenceHashMap; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; @@ -35,6 +38,7 @@ import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -183,6 +187,49 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { log.debug("[{}] evict calculated field links from cached links by entity id: {}", calculatedFieldId, oldCalculatedField); } + @EventListener(ComponentLifecycleMsg.class) + public void onComponentLifecycleEvent(ComponentLifecycleMsg event) { + if (event.getEvent() != ComponentLifecycleEvent.DELETED) { + return; + } + switch (event.getEntityId().getEntityType()) { + case TENANT: + TenantId tenantId = event.getTenantId(); + var removedCfIds = new HashSet(); + calculatedFields.forEach((cfId, cf) -> { + if (cf.getTenantId().equals(tenantId)) { + calculatedFields.remove(cfId); + calculatedFieldLinks.remove(cfId); + calculatedFieldsCtx.remove(cfId); + removedCfIds.add(cfId); + log.debug("[{}] evict calculated field from cache on tenant deletion: {}", cfId, cf); + } + }); + entityIdCalculatedFields.values().forEach(list -> list.removeIf(cf -> removedCfIds.contains(cf.getId()))); + entityIdCalculatedFieldLinks.values().forEach(list -> list.removeIf(link -> removedCfIds.contains(link.getCalculatedFieldId()))); + break; + case DEVICE: + case ASSET: + case DEVICE_PROFILE: + case ASSET_PROFILE: + EntityId entityId = event.getEntityId(); + List cfs = entityIdCalculatedFields.remove(entityId); + if (cfs != null) { + var cfIds = new HashSet(); + cfs.forEach(cf -> { + calculatedFields.remove(cf.getId()); + calculatedFieldLinks.remove(cf.getId()); + calculatedFieldsCtx.remove(cf.getId()); + cfIds.add(cf.getId()); + log.debug("[{}] evict calculated field from cache on entity deletion: {}", cf.getId(), cf); + }); + entityIdCalculatedFieldLinks.values().forEach(list -> list.removeIf(link -> cfIds.contains(link.getCalculatedFieldId()))); + } + entityIdCalculatedFieldLinks.remove(entityId); + break; + } + } + private Lock getFetchLock(CalculatedFieldId id) { return calculatedFieldFetchLocks.computeIfAbsent(id, __ -> new ReentrantLock()); } diff --git a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java index 28fa68d803..a2795e6929 100644 --- a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.profile; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.AssetProfile; @@ -23,9 +24,12 @@ import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.dao.asset.AssetProfileService; import org.thingsboard.server.dao.asset.AssetService; +import java.util.HashSet; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.Lock; @@ -143,6 +147,32 @@ public class DefaultTbAssetProfileCache implements TbAssetProfileCache { } } + @EventListener(ComponentLifecycleMsg.class) + public void onComponentLifecycleEvent(ComponentLifecycleMsg event) { + switch (event.getEntityId().getEntityType()) { + case TENANT: + if (event.getEvent() == ComponentLifecycleEvent.DELETED) { + TenantId tenantId = event.getTenantId(); + var removedProfileIds = new HashSet(); + assetProfilesMap.forEach((assetProfileId, assetProfile) -> { + if (assetProfile.getTenantId().equals(tenantId)) { + assetProfilesMap.remove(assetProfileId); + removedProfileIds.add(assetProfileId); + log.debug("[{}] evict asset profile from cache: {}", assetProfileId, assetProfile); + } + }); + assetsMap.forEach((assetId, assetProfileId) -> { + if (removedProfileIds.contains(assetProfileId)) { + assetsMap.remove(assetId); + } + }); + profileListeners.remove(tenantId); + assetProfileListeners.remove(tenantId); + } + break; + } + } + private void notifyProfileListeners(AssetProfile profile) { ConcurrentMap> tenantListeners = profileListeners.get(profile.getTenantId()); if (tenantListeners != null) { diff --git a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java index 6b356adf94..93072a72c7 100644 --- a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.profile; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -23,9 +24,12 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; +import java.util.HashSet; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.Lock; @@ -143,6 +147,32 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache { } } + @EventListener(ComponentLifecycleMsg.class) + public void onComponentLifecycleEvent(ComponentLifecycleMsg event) { + switch (event.getEntityId().getEntityType()) { + case TENANT: + if (event.getEvent() == ComponentLifecycleEvent.DELETED) { + TenantId tenantId = event.getTenantId(); + var removedProfileIds = new HashSet(); + deviceProfilesMap.forEach((deviceProfileId, deviceProfile) -> { + if (deviceProfile.getTenantId().equals(tenantId)) { + deviceProfilesMap.remove(deviceProfileId); + removedProfileIds.add(deviceProfileId); + log.debug("[{}] evict device profile from cache: {}", deviceProfileId, deviceProfile); + } + }); + devicesMap.forEach((deviceId, deviceProfileId) -> { + if (removedProfileIds.contains(deviceProfileId)) { + devicesMap.remove(deviceId); + } + }); + profileListeners.remove(tenantId); + deviceProfileListeners.remove(tenantId); + } + break; + } + } + private void notifyProfileListeners(DeviceProfile profile) { ConcurrentMap> tenantListeners = profileListeners.get(profile.getTenantId()); if (tenantListeners != null) { From 262565a411dcfe305768798d00fafb17090b0acf Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Thu, 19 Mar 2026 19:56:07 +0200 Subject: [PATCH 2/2] refactoring --- .../cf/DefaultCalculatedFieldCache.java | 35 ++++++++++++++++--- .../profile/DefaultTbAssetProfileCache.java | 21 +++++------ .../profile/DefaultTbDeviceProfileCache.java | 21 +++++------ 3 files changed, 53 insertions(+), 24 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java index 57ce166014..05c60e8f9a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -40,6 +40,7 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import java.util.Collections; import java.util.HashSet; import java.util.List; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CopyOnWriteArrayList; @@ -196,17 +197,42 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { case TENANT: TenantId tenantId = event.getTenantId(); var removedCfIds = new HashSet(); - calculatedFields.forEach((cfId, cf) -> { + var removedCfEntityIds = new HashSet(); + var removedLinkEntityIds = new HashSet(); + for (Map.Entry entry : calculatedFields.entrySet()) { + CalculatedFieldId cfId = entry.getKey(); + CalculatedField cf = entry.getValue(); if (cf.getTenantId().equals(tenantId)) { calculatedFields.remove(cfId); - calculatedFieldLinks.remove(cfId); + List links = calculatedFieldLinks.remove(cfId); + if (links != null) { + links.forEach(link -> removedLinkEntityIds.add(link.getEntityId())); + } calculatedFieldsCtx.remove(cfId); removedCfIds.add(cfId); + removedCfEntityIds.add(cf.getEntityId()); log.debug("[{}] evict calculated field from cache on tenant deletion: {}", cfId, cf); } + } + removedCfEntityIds.forEach(entityId -> { + List cfs = entityIdCalculatedFields.get(entityId); + if (cfs != null) { + cfs.removeIf(cf -> removedCfIds.contains(cf.getId())); + if (cfs.isEmpty()) { + entityIdCalculatedFields.remove(entityId); + } + } + }); + removedLinkEntityIds.forEach(entityId -> { + List entityLinks = entityIdCalculatedFieldLinks.get(entityId); + if (entityLinks != null) { + entityLinks.removeIf(link -> removedCfIds.contains(link.getCalculatedFieldId())); + if (entityLinks.isEmpty()) { + entityIdCalculatedFieldLinks.remove(entityId); + } + } }); - entityIdCalculatedFields.values().forEach(list -> list.removeIf(cf -> removedCfIds.contains(cf.getId()))); - entityIdCalculatedFieldLinks.values().forEach(list -> list.removeIf(link -> removedCfIds.contains(link.getCalculatedFieldId()))); + removedCfIds.forEach(calculatedFieldFetchLocks::remove); break; case DEVICE: case ASSET: @@ -224,6 +250,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { log.debug("[{}] evict calculated field from cache on entity deletion: {}", cf.getId(), cf); }); entityIdCalculatedFieldLinks.values().forEach(list -> list.removeIf(link -> cfIds.contains(link.getCalculatedFieldId()))); + cfIds.forEach(calculatedFieldFetchLocks::remove); } entityIdCalculatedFieldLinks.remove(entityId); break; diff --git a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java index a2795e6929..a0bae27b68 100644 --- a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java @@ -30,6 +30,7 @@ import org.thingsboard.server.dao.asset.AssetProfileService; import org.thingsboard.server.dao.asset.AssetService; import java.util.HashSet; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.Lock; @@ -154,18 +155,18 @@ public class DefaultTbAssetProfileCache implements TbAssetProfileCache { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { TenantId tenantId = event.getTenantId(); var removedProfileIds = new HashSet(); - assetProfilesMap.forEach((assetProfileId, assetProfile) -> { - if (assetProfile.getTenantId().equals(tenantId)) { - assetProfilesMap.remove(assetProfileId); - removedProfileIds.add(assetProfileId); - log.debug("[{}] evict asset profile from cache: {}", assetProfileId, assetProfile); + for (Map.Entry entry : assetProfilesMap.entrySet()) { + if (entry.getValue().getTenantId().equals(tenantId)) { + assetProfilesMap.remove(entry.getKey()); + removedProfileIds.add(entry.getKey()); + log.debug("[{}] evict asset profile from cache: {}", entry.getKey(), entry.getValue()); } - }); - assetsMap.forEach((assetId, assetProfileId) -> { - if (removedProfileIds.contains(assetProfileId)) { - assetsMap.remove(assetId); + } + for (Map.Entry entry : assetsMap.entrySet()) { + if (removedProfileIds.contains(entry.getValue())) { + assetsMap.remove(entry.getKey()); } - }); + } profileListeners.remove(tenantId); assetProfileListeners.remove(tenantId); } diff --git a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java index 93072a72c7..34b5f365f5 100644 --- a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java @@ -30,6 +30,7 @@ import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; import java.util.HashSet; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.Lock; @@ -154,18 +155,18 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { TenantId tenantId = event.getTenantId(); var removedProfileIds = new HashSet(); - deviceProfilesMap.forEach((deviceProfileId, deviceProfile) -> { - if (deviceProfile.getTenantId().equals(tenantId)) { - deviceProfilesMap.remove(deviceProfileId); - removedProfileIds.add(deviceProfileId); - log.debug("[{}] evict device profile from cache: {}", deviceProfileId, deviceProfile); + for (Map.Entry entry : deviceProfilesMap.entrySet()) { + if (entry.getValue().getTenantId().equals(tenantId)) { + deviceProfilesMap.remove(entry.getKey()); + removedProfileIds.add(entry.getKey()); + log.debug("[{}] evict device profile from cache: {}", entry.getKey(), entry.getValue()); } - }); - devicesMap.forEach((deviceId, deviceProfileId) -> { - if (removedProfileIds.contains(deviceProfileId)) { - devicesMap.remove(deviceId); + } + for (Map.Entry entry : devicesMap.entrySet()) { + if (removedProfileIds.contains(entry.getValue())) { + devicesMap.remove(entry.getKey()); } - }); + } profileListeners.remove(tenantId); deviceProfileListeners.remove(tenantId); }