Browse Source

Improvements to queries and cache cleanup

pull/9542/head
Andrii Shvaika 3 years ago
parent
commit
c09caa22dc
  1. 6
      application/src/main/data/upgrade/3.6.1/schema_update.sql
  2. 16
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  3. 46
      application/src/main/java/org/thingsboard/server/service/resource/DefaultTbImageService.java
  4. 3
      application/src/main/java/org/thingsboard/server/service/resource/TbImageService.java
  5. 2
      application/src/main/resources/thingsboard.yml
  6. 7
      common/cluster-api/src/main/proto/queue.proto
  7. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/dashboard/DashboardInfoRepository.java
  8. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/widget/WidgetTypeInfoRepository.java

6
application/src/main/data/upgrade/3.6.1/schema_update.sql

@ -27,9 +27,9 @@ $$
END IF;
END;
$$;
ALTER TABLE resource
ADD COLUMN IF NOT EXISTS descriptor varchar,
ADD COLUMN IF NOT EXISTS preview bytea;
ALTER TABLE resource ADD COLUMN IF NOT EXISTS descriptor varchar;
ALTER TABLE resource ADD COLUMN IF NOT EXISTS preview bytea;
ALTER TABLE resource ADD COLUMN IF NOT EXISTS external_id uuid
-- RESOURCES UPDATE END

16
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -82,6 +82,8 @@ import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.queue.processing.IdMsgPair;
import org.thingsboard.server.service.resource.ImageCacheKey;
import org.thingsboard.server.service.resource.TbImageService;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg;
import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
@ -140,6 +142,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private final TbCoreConsumerStats stats;
protected final TbQueueConsumer<TbProtoQueueMsg<ToUsageStatsServiceMsg>> usageStatsConsumer;
private final TbQueueConsumer<TbProtoQueueMsg<ToOtaPackageStateServiceMsg>> firmwareStatesConsumer;
private final TbImageService imageService;
protected volatile ExecutorService consumersExecutor;
protected volatile ExecutorService usageStatsExecutor;
@ -165,7 +168,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
ApplicationEventPublisher eventPublisher,
Optional<JwtSettingsService> jwtSettingsService,
NotificationSchedulerService notificationSchedulerService,
NotificationRuleProcessor notificationRuleProcessor) {
NotificationRuleProcessor notificationRuleProcessor,
TbImageService imageService) {
super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, eventPublisher, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer(), jwtSettingsService);
this.mainConsumer = tbCoreQueueFactory.createToCoreMsgConsumer();
this.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer();
@ -181,6 +185,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
this.vcQueueService = vcQueueService;
this.notificationSchedulerService = notificationSchedulerService;
this.notificationRuleProcessor = notificationRuleProcessor;
this.imageService = imageService;
}
@PostConstruct
@ -407,6 +412,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
.getNotificationRuleProcessorMsg().getTrigger().toByteArray());
notificationRuleTrigger.ifPresent(notificationRuleProcessor::process);
callback.onSuccess();
} else if (toCoreNotification.hasResourceCacheInvalidateMsg()) {
forwardToResourceService(toCoreNotification.getResourceCacheInvalidateMsg(), callback);
}
if (statsEnabled) {
stats.log(toCoreNotification);
@ -550,6 +557,13 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
callback.onSuccess();
}
private void forwardToResourceService(TransportProtos.ResourceCacheInvalidateMsg msg, TbCallback callback) {
var tenantId = new TenantId(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB()));
imageService.evictETag(new ImageCacheKey(tenantId, msg.getResourceKey(), false));
imageService.evictETag(new ImageCacheKey(tenantId, msg.getResourceKey(), true));
callback.onSuccess();
}
private void forwardToSubMgrService(SubscriptionMgrMsgProto msg, TbCallback callback) {
if (msg.hasSubEvent()) {
TbEntitySubEventProto subEvent = msg.getSubEvent();

46
application/src/main/java/org/thingsboard/server/service/resource/DefaultTbImageService.java

@ -15,11 +15,16 @@
*/
package org.thingsboard.server.service.resource;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.util.Pair;
import org.springframework.stereotype.Service;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.ImageDescriptor;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TbImageDeleteResult;
import org.thingsboard.server.common.data.TbResource;
import org.thingsboard.server.common.data.TbResourceInfo;
@ -28,21 +33,25 @@ import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.id.TbResourceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.resource.ImageService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.entitiy.AbstractTbEntityService;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
@Service
@TbCoreComponent
public class DefaultTbImageService extends AbstractTbEntityService implements TbImageService {
private final TbClusterService clusterService;
private final ImageService imageService;
private final Cache<ImageCacheKey, String> cache;
public DefaultTbImageService(ImageService imageService,
@Value("${cache.image.etag.timeToLiveInMinutes:120}") int cacheTtl,
@Value("${cache.image.etag.maxSize:200000}") int cacheMaxSize) {
public DefaultTbImageService(TbClusterService clusterService, ImageService imageService,
@Value("${cache.image.etag.timeToLiveInMinutes:44640}") int cacheTtl,
@Value("${cache.image.etag.maxSize:10000}") int cacheMaxSize) {
this.clusterService = clusterService;
this.imageService = imageService;
this.cache = Caffeine.newBuilder()
.expireAfterAccess(cacheTtl, TimeUnit.MINUTES)
@ -60,13 +69,33 @@ public class DefaultTbImageService extends AbstractTbEntityService implements Tb
cache.put(imageCacheKey, etag);
}
@Override
public void evictETag(ImageCacheKey imageCacheKey) {
cache.invalidate(imageCacheKey);
}
@Override
public TbResourceInfo save(TbResource image, User user) throws Exception {
ActionType actionType = image.getId() == null ? ActionType.ADDED : ActionType.UPDATED;
TenantId tenantId = image.getTenantId();
try {
var oldEtag = getEtag(image);
TbResourceInfo savedImage = imageService.saveImage(image);
notificationEntityService.logEntityAction(tenantId, savedImage.getId(), savedImage, actionType, user);
if (oldEtag.isPresent()) {
var newEtag = getEtag(savedImage);
if (newEtag.isPresent() && !oldEtag.get().equals(newEtag.get())) {
evictETag(new ImageCacheKey(image.getTenantId(), image.getResourceKey(), false));
evictETag(new ImageCacheKey(image.getTenantId(), image.getResourceKey(), true));
clusterService.broadcastToCore(TransportProtos.ToCoreNotificationMsg.newBuilder()
.setResourceCacheInvalidateMsg(TransportProtos.ResourceCacheInvalidateMsg.newBuilder()
.setTenantIdMSB(tenantId.getId().getMostSignificantBits())
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.setResourceKey(image.getResourceKey())
.build())
.build());
}
}
return savedImage;
} catch (Exception e) {
image.setData(null);
@ -75,6 +104,17 @@ public class DefaultTbImageService extends AbstractTbEntityService implements Tb
}
}
private Optional<String> getEtag(TbResourceInfo image) throws JsonProcessingException {
var descriptor = image.getDescriptor(ImageDescriptor.class);
return Optional.ofNullable(descriptor != null ? descriptor.getEtag() : null);
}
private Optional<String> getPreviewEtag(TbResourceInfo image) throws JsonProcessingException {
var descriptor = image.getDescriptor(ImageDescriptor.class);
descriptor = descriptor != null ? descriptor.getPreviewDescriptor() : null;
return Optional.ofNullable(descriptor != null ? descriptor.getEtag() : null);
}
@Override
public TbResourceInfo save(TbResourceInfo imageInfo, User user) {
TenantId tenantId = imageInfo.getTenantId();

3
application/src/main/java/org/thingsboard/server/service/resource/TbImageService.java

@ -19,6 +19,7 @@ import org.thingsboard.server.common.data.TbImageDeleteResult;
import org.thingsboard.server.common.data.TbResource;
import org.thingsboard.server.common.data.TbResourceInfo;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.TenantId;
public interface TbImageService {
@ -31,4 +32,6 @@ public interface TbImageService {
String getETag(ImageCacheKey imageCacheKey);
void putETag(ImageCacheKey imageCacheKey, String etag);
void evictETag(ImageCacheKey imageCacheKey);
}

2
application/src/main/resources/thingsboard.yml

@ -589,7 +589,7 @@ cache:
image:
etag:
timeToLiveInMinutes: "${CACHE_SPECS_IMAGE_ETAGS_TTL:44640}" # Image ETags cache TTL
maxSize: "${CACHE_SPECS_IMAGE_ETAGS_MAX_SIZE:1000000}" # 0 means the cache is disabled
maxSize: "${CACHE_SPECS_IMAGE_ETAGS_MAX_SIZE:10000}" # 0 means the cache is disabled
systemImagesBrowserTtlInMinutes: "${CACHE_SPECS_IMAGE_SYSTEM_BROWSER_TTL:0}" # Browser cache TTL for system images in minutes. 0 means the cache is disabled
tenantImagesBrowserTtlInMinutes: "${CACHE_SPECS_IMAGE_TENANT_BROWSER_TTL:0}" # Browser cache TTL for tenant images in minutes. 0 means the cache is disabled

7
common/cluster-api/src/main/proto/queue.proto

@ -306,6 +306,12 @@ message CoreStartupMsg {
int64 ts = 3;
}
message ResourceCacheInvalidateMsg {
int64 tenantIdMSB = 1;
int64 tenantIdLSB = 2;
string resourceKey = 3;
}
message LwM2MRegistrationRequestMsg {
string tenantId = 1;
string endpoint = 2;
@ -1273,6 +1279,7 @@ message ToCoreNotificationMsg {
EdgeEventUpdateMsgProto edgeEventUpdate = 14;
ToEdgeSyncRequestMsgProto toEdgeSyncRequest = 15;
FromEdgeSyncResponseMsgProto fromEdgeSyncResponse = 16;
ResourceCacheInvalidateMsg resourceCacheInvalidateMsg = 17;
}
/* Messages that are handled by ThingsBoard RuleEngine Service */

4
dao/src/main/java/org/thingsboard/server/dao/sql/dashboard/DashboardInfoRepository.java

@ -78,12 +78,12 @@ public interface DashboardInfoRepository extends JpaRepository<DashboardInfoEnti
@Query(nativeQuery = true,
value = "SELECT * FROM dashboard d WHERE d.tenant_id = :tenantId " +
"and (d.image = :imageLink or d.configuration ILIKE CONCAT('%', :imageLink, '%')) limit :lmt"
"and (d.image = :imageLink or d.configuration ILIKE CONCAT('%\"', :imageLink, '\"%')) limit :lmt"
)
List<DashboardInfoEntity> findByTenantAndImageLink(@Param("tenantId") UUID tenantId, @Param("imageLink") String imageLink, @Param("lmt") int lmt);
@Query(nativeQuery = true,
value = "SELECT * FROM dashboard d WHERE d.image = :imageLink or d.configuration ILIKE CONCAT('%', :imageLink, '%') limit :lmt"
value = "SELECT * FROM dashboard d WHERE d.image = :imageLink or d.configuration ILIKE CONCAT('%\"', :imageLink, '\"%') limit :lmt"
)
List<DashboardInfoEntity> findByImageLink(@Param("imageLink") String imageLink, @Param("lmt") int lmt);

4
dao/src/main/java/org/thingsboard/server/dao/sql/widget/WidgetTypeInfoRepository.java

@ -198,13 +198,13 @@ public interface WidgetTypeInfoRepository extends JpaRepository<WidgetTypeInfoEn
@Query(nativeQuery = true,
value = "SELECT * FROM widget_type_info_view wti WHERE wti.id IN " +
"(select id from widget_type where tenant_id = :tenantId " +
"and (image = :imageLink or descriptor ILIKE CONCAT('%', :imageLink, '%')) limit :lmt)"
"and (image = :imageLink or descriptor ILIKE CONCAT('%\"', :imageLink, '\"%')) limit :lmt)"
)
List<WidgetTypeInfoEntity> findByTenantAndImageUrl(@Param("tenantId") UUID tenantId, @Param("imageLink") String imageLink, @Param("lmt") int lmt);
@Query(nativeQuery = true,
value = "SELECT * FROM widget_type_info_view wti WHERE wti.id IN " +
"(select id from widget_type where image = :imageLink or descriptor ILIKE CONCAT('%', :imageLink, '%') limit :lmt)"
"(select id from widget_type where image = :imageLink or descriptor ILIKE CONCAT('%\"', :imageLink, '\"%') limit :lmt)"
)
List<WidgetTypeInfoEntity> findByImageUrl(@Param("imageLink") String imageLink, @Param("lmt") int lmt);
}

Loading…
Cancel
Save