Browse Source

Merge branch 'cacheCleanupFIx' into lts_4.3/cacheCleanupFix

# Conflicts:
#	application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
pull/15277/head
dashevchenko 5 months ago
parent
commit
7e9a970597
  1. 74
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  2. 31
      application/src/main/java/org/thingsboard/server/service/profile/DefaultTbAssetProfileCache.java
  3. 31
      application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java

74
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java

@ -20,6 +20,7 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Lazy;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
import org.springframework.util.ConcurrentReferenceHashMap;
import org.thingsboard.server.actors.ActorSystemContext;
@ -28,6 +29,8 @@ import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
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.AssetId;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
@ -45,8 +48,10 @@ import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CopyOnWriteArrayList;
@ -290,6 +295,75 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
});
}
@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<CalculatedFieldId>();
var removedCfEntityIds = new HashSet<EntityId>();
var removedLinkEntityIds = new HashSet<EntityId>();
for (Map.Entry<CalculatedFieldId, CalculatedField> entry : calculatedFields.entrySet()) {
CalculatedFieldId cfId = entry.getKey();
CalculatedField cf = entry.getValue();
if (cf.getTenantId().equals(tenantId)) {
calculatedFields.remove(cfId);
List<CalculatedFieldLink> links = calculatedFieldLinks.remove(cfId);
if (links != null) {
links.forEach(link -> removedLinkEntityIds.add(link.entityId()));
}
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<CalculatedField> cfs = entityIdCalculatedFields.get(entityId);
if (cfs != null) {
cfs.removeIf(cf -> removedCfIds.contains(cf.getId()));
if (cfs.isEmpty()) {
entityIdCalculatedFields.remove(entityId);
}
}
});
removedLinkEntityIds.forEach(entityId -> {
List<CalculatedFieldLink> entityLinks = entityIdCalculatedFieldLinks.get(entityId);
if (entityLinks != null) {
entityLinks.removeIf(link -> removedCfIds.contains(link.calculatedFieldId()));
if (entityLinks.isEmpty()) {
entityIdCalculatedFieldLinks.remove(entityId);
}
}
});
removedCfIds.forEach(calculatedFieldFetchLocks::remove);
break;
case DEVICE:
case ASSET:
case DEVICE_PROFILE:
case ASSET_PROFILE:
EntityId entityId = event.getEntityId();
List<CalculatedField> cfs = entityIdCalculatedFields.remove(entityId);
if (cfs != null) {
var cfIds = new HashSet<CalculatedFieldId>();
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.calculatedFieldId())));
cfIds.forEach(calculatedFieldFetchLocks::remove);
}
entityIdCalculatedFieldLinks.remove(entityId);
break;
}
}
private Lock getFetchLock(CalculatedFieldId id) {
return calculatedFieldFetchLocks.computeIfAbsent(id, __ -> new ReentrantLock());
}

31
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,13 @@ 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.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.Lock;
@ -143,6 +148,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<AssetProfileId>();
for (Map.Entry<AssetProfileId, AssetProfile> 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());
}
}
for (Map.Entry<AssetId, AssetProfileId> entry : assetsMap.entrySet()) {
if (removedProfileIds.contains(entry.getValue())) {
assetsMap.remove(entry.getKey());
}
}
profileListeners.remove(tenantId);
assetProfileListeners.remove(tenantId);
}
break;
}
}
private void notifyProfileListeners(AssetProfile profile) {
ConcurrentMap<EntityId, Consumer<AssetProfile>> tenantListeners = profileListeners.get(profile.getTenantId());
if (tenantListeners != null) {

31
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,13 @@ 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.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.Lock;
@ -143,6 +148,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<DeviceProfileId>();
for (Map.Entry<DeviceProfileId, DeviceProfile> 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());
}
}
for (Map.Entry<DeviceId, DeviceProfileId> entry : devicesMap.entrySet()) {
if (removedProfileIds.contains(entry.getValue())) {
devicesMap.remove(entry.getKey());
}
}
profileListeners.remove(tenantId);
deviceProfileListeners.remove(tenantId);
}
break;
}
}
private void notifyProfileListeners(DeviceProfile profile) {
ConcurrentMap<EntityId, Consumer<DeviceProfile>> tenantListeners = profileListeners.get(profile.getTenantId());
if (tenantListeners != null) {

Loading…
Cancel
Save