diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 1ed919e922..26c82a33de 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -30,9 +30,9 @@ import org.springframework.context.annotation.Lazy; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.DeviceStateManager; import org.thingsboard.rule.engine.api.MailService; import org.thingsboard.rule.engine.api.NotificationCenter; -import org.thingsboard.rule.engine.api.DeviceStateManager; import org.thingsboard.rule.engine.api.SmsService; import org.thingsboard.rule.engine.api.notification.SlackService; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; @@ -106,11 +106,11 @@ import org.thingsboard.server.dao.widget.WidgetsBundleService; import org.thingsboard.server.queue.discovery.DiscoveryService; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.settings.TbQueueCalculatedFieldSettings; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldQueueService; import org.thingsboard.server.service.cf.CalculatedFieldStateService; -import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; @@ -460,6 +460,11 @@ public class ActorSystemContext { @Getter private DebugModeRateLimitsConfig debugModeRateLimitsConfig; + @Lazy + @Autowired(required = false) + @Getter + private TbQueueCalculatedFieldSettings calculatedFieldSettings; + /** * The following Service will be null if we operate in tb-core mode */ @@ -546,11 +551,6 @@ public class ActorSystemContext { @Getter private CalculatedFieldQueueService calculatedFieldQueueService; - @Lazy - @Autowired(required = false) - @Getter - private CalculatedFieldEntityProfileCache calculatedFieldEntityProfileCache; - @Value("${actors.session.max_concurrent_sessions_per_device:1}") @Getter private int maxConcurrentSessionsPerDevice; diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index 904547a2b6..a79a182fa1 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -113,6 +113,8 @@ public class AppActor extends ContextAwareActor { case SESSION_TIMEOUT_MSG: ctx.broadcastToChildrenByType(msg, EntityType.TENANT); break; + case CF_CACHE_INIT_MSG: + case CF_INIT_PROFILE_ENTITY_MSG: case CF_INIT_MSG: case CF_LINK_INIT_MSG: case CF_STATE_RESTORE_MSG: diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java index 70ed2849e8..9f59a80e67 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java @@ -22,8 +22,10 @@ import org.thingsboard.server.actors.TbActorException; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbActorStopReason; import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; +import org.thingsboard.server.common.msg.cf.CalculatedFieldCacheInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldInitMsg; +import org.thingsboard.server.common.msg.cf.CalculatedFieldInitProfileEntityMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldLinkInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; @@ -65,6 +67,12 @@ public class CalculatedFieldManagerActor extends AbstractCalculatedFieldActor { case CF_PARTITIONS_CHANGE_MSG: processor.onPartitionChange((CalculatedFieldPartitionChangeMsg) msg); break; + case CF_CACHE_INIT_MSG: + processor.onCacheInitMsg((CalculatedFieldCacheInitMsg) msg); + break; + case CF_INIT_PROFILE_ENTITY_MSG: + processor.onProfileEntityMsg((CalculatedFieldInitProfileEntityMsg) msg); + break; case CF_INIT_MSG: processor.onFieldInitMsg((CalculatedFieldInitMsg) msg); break; diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 31152c7c98..ba1ca71bda 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -24,6 +24,7 @@ import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.id.AssetId; @@ -31,17 +32,23 @@ import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageDataIterable; +import org.thingsboard.server.common.msg.cf.CalculatedFieldCacheInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldInitMsg; +import org.thingsboard.server.common.msg.cf.CalculatedFieldInitProfileEntityMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldLinkInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.cf.CalculatedFieldService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.queue.settings.TbQueueCalculatedFieldSettings; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldStateService; -import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache; +import org.thingsboard.server.service.cf.cache.TenantEntityProfileCache; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.profile.TbAssetProfileCache; @@ -68,22 +75,28 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware private final CalculatedFieldProcessingService cfExecService; private final CalculatedFieldStateService cfStateService; - private final CalculatedFieldEntityProfileCache cfEntityCache; private final CalculatedFieldService cfDaoService; + private final DeviceService deviceService; + private final AssetService assetService; private final TbAssetProfileCache assetProfileCache; private final TbDeviceProfileCache deviceProfileCache; + private final TenantEntityProfileCache entityProfileCache; + private final TbQueueCalculatedFieldSettings cfSettings; protected final TenantId tenantId; protected TbActorCtx ctx; CalculatedFieldManagerMessageProcessor(ActorSystemContext systemContext, TenantId tenantId) { super(systemContext); - this.cfEntityCache = systemContext.getCalculatedFieldEntityProfileCache(); this.cfExecService = systemContext.getCalculatedFieldProcessingService(); this.cfStateService = systemContext.getCalculatedFieldStateService(); this.cfDaoService = systemContext.getCalculatedFieldService(); + this.deviceService = systemContext.getDeviceService(); + this.assetService = systemContext.getAssetService(); this.assetProfileCache = systemContext.getAssetProfileCache(); this.deviceProfileCache = systemContext.getDeviceProfileCache(); + this.entityProfileCache = new TenantEntityProfileCache(); + this.cfSettings = systemContext.getCalculatedFieldSettings(); this.tenantId = tenantId; } @@ -100,6 +113,19 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware ctx.stop(ctx.getSelf()); } + public void onCacheInitMsg(CalculatedFieldCacheInitMsg msg) { + log.debug("[{}] Processing CF actor init message.", msg.getTenantId().getId()); + initEntityProfileCache(); + initCalculatedFields(); + msg.getCallback().onSuccess(); + } + + public void onProfileEntityMsg(CalculatedFieldInitProfileEntityMsg msg) { + log.debug("[{}] Processing profile entity message.", msg.getTenantId().getId()); + entityProfileCache.add(msg.getProfileEntityId(), msg.getEntityId()); + msg.getCallback().onSuccess(); + } + public void onFieldInitMsg(CalculatedFieldInitMsg msg) throws CalculatedFieldException { log.debug("[{}] Processing CF init message.", msg.getCf().getId()); var cf = msg.getCf(); @@ -180,16 +206,35 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } break; } + case DEVICE_PROFILE: + case ASSET_PROFILE: { + switch (event) { + case DELETED: + onProfileDeleted(msg.getData(), msg.getCallback()); + break; + default: + msg.getCallback().onSuccess(); + break; + } + break; + } default: { msg.getCallback().onSuccess(); } } } + private void onProfileDeleted(ComponentLifecycleMsg msg, TbCallback callback) { + entityProfileCache.removeProfileId(msg.getEntityId()); + callback.onSuccess(); + } + private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) { EntityId entityId = msg.getEntityId(); EntityId profileId = getProfileId(tenantId, entityId); - cfEntityCache.add(tenantId, profileId, entityId); + if (profileId != null) { + entityProfileCache.add(profileId, entityId); + } if (!isMyPartition(entityId, callback)) { return; } @@ -207,7 +252,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware private void onEntityUpdated(ComponentLifecycleMsg msg, TbCallback callback) { if (msg.getOldProfileId() != null && !msg.getOldProfileId().equals(msg.getProfileId())) { - cfEntityCache.update(tenantId, msg.getOldProfileId(), msg.getProfileId(), msg.getEntityId()); + entityProfileCache.update(msg.getOldProfileId(), msg.getProfileId(), msg.getEntityId()); if (!isMyPartition(msg.getEntityId(), callback)) { return; } @@ -226,7 +271,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } private void onEntityDeleted(ComponentLifecycleMsg msg, TbCallback callback) { - cfEntityCache.evict(tenantId, msg.getEntityId()); + entityProfileCache.removeEntityId(msg.getEntityId()); if (isMyPartition(msg.getEntityId(), callback)) { log.debug("Pushing entity lifecycle msg to specific actor [{}]", msg.getEntityId()); getOrCreateActor(msg.getEntityId()).tell(new CalculatedFieldEntityDeleteMsg(tenantId, msg.getEntityId(), callback)); @@ -322,11 +367,15 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware EntityId entityId = cfCtx.getEntityId(); EntityType entityType = cfCtx.getEntityId().getEntityType(); if (isProfileEntity(entityType)) { - var entityIds = cfEntityCache.getMyEntityIdsByProfileId(tenantId, entityId); + var entityIds = entityProfileCache.getEntityIdsByProfileId(entityId); if (!entityIds.isEmpty()) { //TODO: no need to do this if we cache all created actors and know which one belong to us; var multiCallback = new MultipleTbCallback(entityIds.size(), callback); - entityIds.forEach(id -> deleteCfForEntity(id, cfId, multiCallback)); + entityIds.forEach(id -> { + if (isMyPartition(id, multiCallback)) { + deleteCfForEntity(id, cfId, multiCallback); + } + }); } else { callback.onSuccess(); } @@ -366,10 +415,11 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware EntityId sourceEntityId = msg.getEntityId(); log.debug("Received linked telemetry msg from entity [{}]", sourceEntityId); var proto = msg.getProto(); + var callback = msg.getCallback(); var linksList = proto.getLinksList(); if (linksList.isEmpty()) { log.debug("[{}] No CF links to process new telemetry.", msg.getTenantId()); - msg.getCallback().onSuccess(); + callback.onSuccess(); } for (var linkProto : linksList) { var link = fromProto(linkProto); @@ -378,21 +428,23 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware var cf = calculatedFields.get(link.cfId()); if (EntityType.DEVICE_PROFILE.equals(targetEntityType) || EntityType.ASSET_PROFILE.equals(targetEntityType)) { // iterate over all entities that belong to profile and push the message for corresponding CF - var entityIds = cfEntityCache.getMyEntityIdsByProfileId(tenantId, targetEntityId); + var entityIds = entityProfileCache.getEntityIdsByProfileId(targetEntityId); if (!entityIds.isEmpty()) { - MultipleTbCallback callback = new MultipleTbCallback(entityIds.size(), msg.getCallback()); - var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, callback); + MultipleTbCallback multipleCallback = new MultipleTbCallback(entityIds.size(), callback); + var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, multipleCallback); entityIds.forEach(entityId -> { - log.debug("Pushing linked telemetry msg to specific actor [{}]", entityId); - getOrCreateActor(entityId).tell(newMsg); + if (isMyPartition(entityId, multipleCallback)) { + log.debug("Pushing linked telemetry msg to specific actor [{}]", entityId); + getOrCreateActor(entityId).tell(newMsg); + } }); } else { - msg.getCallback().onSuccess(); + callback.onSuccess(); } } else { - if (isMyPartition(targetEntityId, msg.getCallback())) { + if (isMyPartition(targetEntityId, callback)) { log.debug("Pushing linked telemetry msg to specific actor [{}]", targetEntityId); - var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, msg.getCallback()); + var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, callback); getOrCreateActor(targetEntityId).tell(newMsg); } } @@ -438,10 +490,14 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware EntityId entityId = cfCtx.getEntityId(); EntityType entityType = cfCtx.getEntityId().getEntityType(); if (isProfileEntity(entityType)) { - var entityIds = cfEntityCache.getMyEntityIdsByProfileId(tenantId, entityId); + var entityIds = entityProfileCache.getEntityIdsByProfileId(entityId); if (!entityIds.isEmpty()) { var multiCallback = new MultipleTbCallback(entityIds.size(), callback); - entityIds.forEach(id -> initCfForEntity(id, cfCtx, forceStateReinit, multiCallback)); + entityIds.forEach(id -> { + if (isMyPartition(id, multiCallback)) { + initCfForEntity(id, cfCtx, forceStateReinit, multiCallback); + } + }); } else { callback.onSuccess(); } @@ -504,4 +560,45 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware public void onPartitionChange(CalculatedFieldPartitionChangeMsg msg) { ctx.broadcastToChildren(msg, true); } + + public void initCalculatedFields() { + PageDataIterable cfs = new PageDataIterable<>(pageLink -> cfDaoService.findCalculatedFieldsByTenantId(tenantId, pageLink), cfSettings.getInitTenantFetchPackSize()); + cfs.forEach(cf -> { + log.trace("Processing calculated field record: {}", cf); + try { + onFieldInitMsg(new CalculatedFieldInitMsg(cf.getTenantId(), cf)); + } catch (CalculatedFieldException e) { + log.error("Failed to process calculated field record: {}", cf, e); + } + }); + calculatedFields.values().forEach(cf -> { + entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cf); + }); + PageDataIterable cfls = new PageDataIterable<>(pageLink -> cfDaoService.findAllCalculatedFieldLinksByTenantId(tenantId, pageLink), cfSettings.getInitTenantFetchPackSize()); + cfls.forEach(link -> { + onLinkInitMsg(new CalculatedFieldLinkInitMsg(link.getTenantId(), link)); + }); + } + + private void initEntityProfileCache() { + PageDataIterable deviceIdInfos = new PageDataIterable<>(pageLink -> deviceService.findProfileEntityIdInfosByTenantId(tenantId, pageLink), cfSettings.getInitTenantFetchPackSize()); + for (ProfileEntityIdInfo idInfo : deviceIdInfos) { + log.trace("Processing device record: {}", idInfo); + try { + entityProfileCache.add(idInfo.getProfileId(), idInfo.getEntityId()); + } catch (Exception e) { + log.error("Failed to process device record: {}", idInfo, e); + } + } + PageDataIterable assetIdInfos = new PageDataIterable<>(pageLink -> assetService.findProfileEntityIdInfosByTenantId(tenantId, pageLink), cfSettings.getInitTenantFetchPackSize()); + for (ProfileEntityIdInfo idInfo : assetIdInfos) { + log.trace("Processing asset record: {}", idInfo); + try { + entityProfileCache.add(idInfo.getProfileId(), idInfo.getEntityId()); + } catch (Exception e) { + log.error("Failed to process asset record: {}", idInfo, e); + } + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index 728f715af4..f0d330726d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -50,6 +50,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; +import org.thingsboard.server.common.msg.cf.CalculatedFieldCacheInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; @@ -176,6 +177,8 @@ public class TenantActor extends RuleChainManagerActor { case RULE_CHAIN_TO_RULE_CHAIN_MSG: onRuleChainMsg((RuleChainAwareMsg) msg); break; + case CF_CACHE_INIT_MSG: + case CF_INIT_PROFILE_ENTITY_MSG: case CF_INIT_MSG: case CF_LINK_INIT_MSG: case CF_STATE_RESTORE_MSG: @@ -199,6 +202,7 @@ public class TenantActor extends RuleChainManagerActor { } else { log.debug("[{}] CF Actor is not initialized. ToCalculatedFieldSystemMsg: [{}]", tenantId, msg); } + msg.getCallback().onSuccess(); return; } if (priority) { @@ -274,6 +278,7 @@ public class TenantActor extends RuleChainManagerActor { () -> DefaultActorService.CF_MANAGER_DISPATCHER_NAME, () -> new CalculatedFieldManagerActorCreator(systemContext, tenantId), () -> true); + cfActor.tellWithHighPriority(new CalculatedFieldCacheInitMsg(tenantId)); } catch (Exception e) { log.info("[{}] Failed to init CF Actor.", tenantId, e); } 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 219a261183..ed35d96cb7 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 @@ -64,7 +64,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { private final ConcurrentMap> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); - @Value("${calculatedField.initFetchPackSize:50000}") + @Value("${queue.calculated_fields.init_fetch_pack_size:50000}") @Getter private int initFetchPackSize; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldInitService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldInitService.java index cc3022fa29..79950f2e3f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldInitService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldInitService.java @@ -20,13 +20,14 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; +import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.page.PageDataIterable; +import org.thingsboard.server.common.msg.cf.CalculatedFieldInitProfileEntityMsg; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbRuleEngineComponent; -import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache; @Slf4j @Service @@ -34,11 +35,11 @@ import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache @RequiredArgsConstructor public class DefaultCalculatedFieldInitService implements CalculatedFieldInitService { - private final CalculatedFieldEntityProfileCache entityProfileCache; private final AssetService assetService; private final DeviceService deviceService; + private final ActorSystemContext actorSystemContext; - @Value("${calculated_fields.init_fetch_pack_size:50000}") + @Value("${queue.calculated_fields.init_fetch_pack_size:50000}") @Getter private int initFetchPackSize; @@ -48,7 +49,7 @@ public class DefaultCalculatedFieldInitService implements CalculatedFieldInitSer for (ProfileEntityIdInfo idInfo : deviceIdInfos) { log.trace("Processing device record: {}", idInfo); try { - entityProfileCache.add(idInfo.getTenantId(), idInfo.getProfileId(), idInfo.getEntityId()); + actorSystemContext.tell(new CalculatedFieldInitProfileEntityMsg(idInfo.getTenantId(), idInfo.getProfileId(), idInfo.getEntityId())); } catch (Exception e) { log.error("Failed to process device record: {}", idInfo, e); } @@ -57,7 +58,7 @@ public class DefaultCalculatedFieldInitService implements CalculatedFieldInitSer for (ProfileEntityIdInfo idInfo : assetIdInfos) { log.trace("Processing asset record: {}", idInfo); try { - entityProfileCache.add(idInfo.getTenantId(), idInfo.getProfileId(), idInfo.getEntityId()); + actorSystemContext.tell(new CalculatedFieldInitProfileEntityMsg(idInfo.getTenantId(), idInfo.getProfileId(), idInfo.getEntityId())); } catch (Exception e) { log.error("Failed to process asset record: {}", idInfo, e); } 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 deleted file mode 100644 index 2f5772ae50..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/cf/cache/DefaultCalculatedFieldEntityProfileCache.java +++ /dev/null @@ -1,95 +0,0 @@ -/** - * Copyright © 2016-2025 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.service.cf.cache; - -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Service; -import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.id.EntityId; -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; - -@TbRuleEngineComponent -@Service -@Slf4j -@RequiredArgsConstructor -//TODO ashvayka: remove and use TenantEntityProfileCache in each CalculatedFieldManagerMessageProcessor; -public class DefaultCalculatedFieldEntityProfileCache extends TbApplicationEventListener implements CalculatedFieldEntityProfileCache { - - 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() - .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)); - } - - @Override - public void add(TenantId tenantId, EntityId profileId, EntityId entityId) { - var tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, entityId); - var partition = tpi.getPartition().orElse(UNKNOWN); - tenantCache.computeIfAbsent(tenantId, id -> new TenantEntityProfileCache()) - .add(profileId, entityId, partition, tpi.isMyPartition()); - } - - @Override - public void update(TenantId tenantId, EntityId oldProfileId, EntityId newProfileId, EntityId entityId) { - var tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, entityId); - var partition = tpi.getPartition().orElse(UNKNOWN); - var cache = tenantCache.computeIfAbsent(tenantId, id -> new TenantEntityProfileCache()); - //TODO: make this method atomic; - cache.remove(oldProfileId, entityId); - cache.add(newProfileId, entityId, partition, tpi.isMyPartition()); - } - - @Override - public void evict(TenantId tenantId, EntityId entityId) { - var cache = tenantCache.computeIfAbsent(tenantId, id -> new TenantEntityProfileCache()); - cache.removeEntityId(entityId); - } - - @Override - public Collection getMyEntityIdsByProfileId(TenantId tenantId, EntityId profileId) { - return tenantCache.computeIfAbsent(tenantId, id -> new TenantEntityProfileCache()).getMyEntityIdsByProfileId(profileId); - } - - @Override - public int getEntityIdPartition(TenantId tenantId, EntityId entityId) { - var tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, entityId); - return tpi.getPartition().orElse(UNKNOWN); - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/cache/TenantEntityProfileCache.java b/application/src/main/java/org/thingsboard/server/service/cf/cache/TenantEntityProfileCache.java index 1a17b9b8be..6435ffb14c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/cache/TenantEntityProfileCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/cache/TenantEntityProfileCache.java @@ -32,31 +32,13 @@ import java.util.concurrent.locks.ReentrantReadWriteLock; public class TenantEntityProfileCache { private final ReadWriteLock lock = new ReentrantReadWriteLock(); - private final Map>> allEntities = new HashMap<>(); - private final Map> myEntities = new HashMap<>(); - - public void setMyPartitions(List myPartitions) { - lock.writeLock().lock(); - try { - myEntities.clear(); - myPartitions.forEach(partitionId -> { - var map = allEntities.get(partitionId); - if (map != null) { - map.forEach((profileId, entityIds) -> myEntities.computeIfAbsent(profileId, k -> new HashSet<>()).addAll(entityIds)); - } - }); - } finally { - lock.writeLock().unlock(); - } - } + private final Map> allEntities = new HashMap<>(); public void removeProfileId(EntityId profileId) { lock.writeLock().lock(); try { // Remove from allEntities - allEntities.values().forEach(map -> map.remove(profileId)); - // Remove from myEntities - myEntities.remove(profileId); + allEntities.remove(profileId); } finally { lock.writeLock().unlock(); } @@ -66,9 +48,7 @@ public class TenantEntityProfileCache { lock.writeLock().lock(); try { // Remove from allEntities - allEntities.values().forEach(map -> map.values().forEach(set -> set.remove(entityId))); - // Remove from myEntities - myEntities.values().forEach(set -> set.remove(entityId)); + allEntities.values().forEach(set -> set.remove(entityId)); } finally { lock.writeLock().unlock(); } @@ -78,33 +58,33 @@ public class TenantEntityProfileCache { lock.writeLock().lock(); try { // Remove from allEntities - allEntities.values().forEach(map -> removeSafely(map, profileId, entityId)); - // Remove from myEntities - removeSafely(myEntities, profileId, entityId); + removeSafely(allEntities, profileId, entityId); } finally { lock.writeLock().unlock(); } } - public void add(EntityId profileId, EntityId entityId, Integer partition, boolean mine) { + public void add(EntityId profileId, EntityId entityId) { lock.writeLock().lock(); try { - if(EntityType.DEVICE.equals(profileId.getEntityType())){ - throw new RuntimeException("WTF?"); - } - if (mine) { - myEntities.computeIfAbsent(profileId, k -> new HashSet<>()).add(entityId); + if (EntityType.DEVICE.equals(profileId.getEntityType()) || EntityType.ASSET.equals(profileId.getEntityType())) { + throw new RuntimeException("Entity type '" + profileId.getEntityType() + "' is not a profileId."); } - allEntities.computeIfAbsent(partition, k -> new HashMap<>()).computeIfAbsent(profileId, p -> new HashSet<>()).add(entityId); + allEntities.computeIfAbsent(profileId, k -> new HashSet<>()).add(entityId); } finally { lock.writeLock().unlock(); } } - public Collection getMyEntityIdsByProfileId(EntityId profileId) { + public void update(EntityId oldProfileId, EntityId newProfileId, EntityId entityId) { + remove(oldProfileId, entityId); + add(newProfileId, entityId); + } + + public Collection getEntityIdsByProfileId(EntityId profileId) { lock.readLock().lock(); try { - var entities = myEntities.getOrDefault(profileId, Collections.emptySet()); + var entities = allEntities.getOrDefault(profileId, Collections.emptySet()); List result = new ArrayList<>(entities.size()); result.addAll(entities); return result; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java index 3620ad3639..533b487d38 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java @@ -91,7 +91,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } } }) - .consumerCreator((config, partitionId) -> queueFactory.createCalculatedFieldStateConsumer()) + .consumerCreator((queueConfig, tpi) -> queueFactory.createCalculatedFieldStateConsumer()) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(eventConsumer.getConsumerExecutor()) .scheduler(eventConsumer.getScheduler()) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java index 0eaa506dfd..9dc6139ca5 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java @@ -22,11 +22,11 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.queue.common.state.DefaultQueueStateService; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.common.state.DefaultQueueStateService; import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; import org.thingsboard.server.service.cf.CfRocksDb; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java index 076e8eb8f6..37687f14a9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java @@ -141,6 +141,7 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edge.getId()).getTopic(); TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); kafkaAdmin.deleteTopic(topic); + kafkaAdmin.deleteConsumerGroup(topic); } } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java index 5da5fa2be0..5946b84d09 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java @@ -106,7 +106,7 @@ public class HousekeeperStatsService { this.failedProcessingCounter = register("failedProcessing", statsFactory); this.reprocessedCounter = register("reprocessed", statsFactory); this.failedReprocessingCounter = register("failedReprocessing", statsFactory); - this.processingTimer = statsFactory.createTimer(StatsType.HOUSEKEEPER, "processingTime", "taskType", taskType.name()); + this.processingTimer = statsFactory.createStatsTimer(StatsType.HOUSEKEEPER.getName(), "processingTime", "taskType", taskType.name()); } private StatsCounter register(String statsName, StatsFactory statsFactory) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java index 125a06299d..f8c398b016 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java @@ -16,7 +16,6 @@ package org.thingsboard.server.service.queue; import jakarta.annotation.PreDestroy; -import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.CollectionUtils; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; @@ -70,7 +69,6 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent -@Slf4j public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBasedConsumerService implements TbCalculatedFieldConsumerService { @Value("${queue.calculated_fields.poll_interval:25}") @@ -101,12 +99,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa @Override protected void onStartUp() { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); - PartitionedQueueConsumerManager> eventConsumer = PartitionedQueueConsumerManager.>create() + var eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) .topic(partitionService.getTopic(queueKey)) .pollInterval(pollInterval) .msgPackProcessor(this::processMsgs) - .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer()) + .consumerCreator((queueConfig, tpi) -> queueFactory.createToCalculatedFieldMsgConsumer(tpi)) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) @@ -129,7 +127,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa protected void onPartitionChangeEvent(PartitionChangeEvent event) { try { event.getNewPartitions().forEach((queueKey, partitions) -> { - if (queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { + if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName())) { stateService.restore(queueKey, partitions); } }); @@ -177,9 +175,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa packSubmitFuture.cancel(true); log.info("Timeout to process message: {}", pendingMsgHolder.getMsg()); } - if (log.isDebugEnabled()) { - ctx.getAckMap().forEach((id, msg) -> log.debug("[{}] Timeout to process message: {}", id, msg.getValue())); - } + ctx.getAckMap().forEach((id, msg) -> log.warn("[{}] Timeout to process message: {}", id, msg.getValue())); ctx.getFailedMap().forEach((id, msg) -> log.warn("[{}] Failed to process message: {}", id, msg.getValue())); } consumer.commit(); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index a3003ba6ff..47a0c473e9 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -117,8 +117,8 @@ import java.util.function.Function; import java.util.stream.Collectors; @Service -@TbCoreComponent @Slf4j +@TbCoreComponent public class DefaultTbCoreConsumerService extends AbstractConsumerService implements TbCoreConsumerService { @Value("${queue.core.poll-interval}") @@ -206,7 +206,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService queueFactory.createToCoreMsgConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createToCoreMsgConsumer()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java index 2e4c15f8e8..a9ae18a11f 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java @@ -19,8 +19,6 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import lombok.Data; -import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.jetbrains.annotations.NotNull; import org.springframework.beans.factory.annotation.Value; @@ -45,12 +43,12 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.consumer.MainQueueConsumerManager; import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; -import org.thingsboard.server.queue.common.consumer.MainQueueConsumerManager; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.IdMsgPair; @@ -66,7 +64,6 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; -@Slf4j @Service @TbCoreComponent public class DefaultTbEdgeConsumerService extends AbstractConsumerService implements TbEdgeConsumerService { @@ -105,7 +102,7 @@ public class DefaultTbEdgeConsumerService extends AbstractConsumerService queueFactory.createEdgeMsgConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdgeMsgConsumer()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index 98f857750a..f51678310e 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.queue; -import lombok.extern.slf4j.Slf4j; import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.event.EventListener; import org.springframework.scheduling.annotation.Scheduled; @@ -65,7 +64,6 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent -@Slf4j public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedConsumerService implements TbRuleEngineConsumerService { private final TbRuleEngineConsumerContext ctx; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index d12bd896ff..26b689df9d 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -17,7 +17,8 @@ package org.thingsboard.server.service.queue.processing; import jakarta.annotation.PreDestroy; import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.context.ApplicationEventPublisher; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; @@ -62,10 +63,11 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; -@Slf4j @RequiredArgsConstructor public abstract class AbstractConsumerService extends TbApplicationEventListener { + protected final Logger log = LoggerFactory.getLogger(getClass()); + protected final ActorSystemContext actorContext; protected final TbTenantProfileCache tenantProfileCache; protected final TbDeviceProfileCache deviceProfileCache; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java index 97aa81d41c..f42908bcb5 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java @@ -28,6 +28,8 @@ import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -35,7 +37,7 @@ public abstract class AbstractPartitionBasedConsumerService pendingEvents = new ArrayList<>(); public AbstractPartitionBasedConsumerService(ActorSystemContext actorContext, TbTenantProfileCache tenantProfileCache, @@ -61,8 +63,16 @@ public abstract class AbstractPartitionBasedConsumerService { + Integer partitionId = tpi != null ? tpi.getPartition().orElse(-1) : null; + return ctx.getQueueFactory().createToRuleEngineMsgConsumer(queueConfig, partitionId); + }, + consumerExecutor, scheduler, taskExecutor, null); this.ctx = ctx; this.stats = new TbRuleEngineConsumerStats(queueKey, ctx.getStatsFactory()); } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java index 73712542fd..6d06c85585 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java @@ -113,7 +113,7 @@ public class KafkaEdgeTopicsCleanUpService extends AbstractCleanUpService { .ifPresentOrElse(lastConnectTime -> { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); if (kafkaAdmin.isTopicEmpty(topic)) { - kafkaAdmin.deleteTopic(topic); + deleteTopicAndConsumerGroup(topic); log.info("[{}] Removed outdated topic {} for edge {} older than {}", tenantId, topic, edgeId, Date.from(Instant.ofEpochMilli(currentTimeMillis - ttlMillis))); } @@ -121,7 +121,7 @@ public class KafkaEdgeTopicsCleanUpService extends AbstractCleanUpService { Edge edge = edgeService.findEdgeById(tenantId, edgeId); if (edge == null) { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - kafkaAdmin.deleteTopic(topic); + deleteTopicAndConsumerGroup(topic); log.info("[{}] Removed topic {} for deleted edge {}", tenantId, topic, edgeId); } }); @@ -132,12 +132,17 @@ public class KafkaEdgeTopicsCleanUpService extends AbstractCleanUpService { } else { for (EdgeId edgeId : edgeIds) { String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - kafkaAdmin.deleteTopic(topic); + deleteTopicAndConsumerGroup(topic); } log.info("[{}] Removed topics for not existing tenant and edges {}", tenantId, edgeIds); } } + private void deleteTopicAndConsumerGroup(String topic) { + kafkaAdmin.deleteTopic(topic); + kafkaAdmin.deleteConsumerGroup(topic); + } + private boolean isTopicExpired(long lastConnectTime, long ttlMillis, long currentTimeMillis) { return lastConnectTime + ttlMillis < currentTimeMillis; } @@ -146,7 +151,7 @@ public class KafkaEdgeTopicsCleanUpService extends AbstractCleanUpService { Map> tenantEdgeMap = new HashMap<>(); for (String topic : topics) { try { - String remaining = topic.substring(prefix.length()); + String remaining = topic.substring(prefix.length() + 1); String[] parts = remaining.split("\\."); TenantId tenantId = TenantId.fromUUID(UUID.fromString(parts[0])); EdgeId edgeId = new EdgeId(UUID.fromString(parts[1])); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 716797c2d2..abafd81e7c 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1773,7 +1773,7 @@ queue: # EDQS responses topic responses_topic: "${TB_EDQS_RESPONSES_TOPIC:edqs.responses}" # Poll interval for EDQS topics - poll_interval: "${TB_EDQS_POLL_INTERVAL_MS:125}" + poll_interval: "${TB_EDQS_POLL_INTERVAL_MS:25}" # Maximum amount of pending requests to EDQS max_pending_requests: "${TB_EDQS_MAX_PENDING_REQUESTS:10000}" # Maximum timeout for requests to EDQS @@ -1848,6 +1848,10 @@ queue: pool_size: "${TB_QUEUE_CF_POOL_SIZE:8}" # RocksDB path for storing CF states rocks_db_path: "${TB_QUEUE_CF_ROCKS_DB_PATH:${user.home}/.rocksdb/cf_states}" + # The fetch size specifies how many rows will be fetched from the database per request for initial fetching + init_fetch_pack_size: "${TB_QUEUE_CF_FETCH_PACK_SIZE:50000}" + # The fetch size specifies how many rows will be fetched from the database per request for per-tenant fetching + init_tenant_fetch_pack_size: "${TB_QUEUE_CF_TENANT_FETCH_PACK_SIZE:1000}" transport: # For high-priority notifications that require minimum latency and processing time notifications_topic: "${TB_QUEUE_TRANSPORT_NOTIFICATIONS_TOPIC:tb_transport.notifications}" diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbEdgeQueueAdmin.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbEdgeQueueAdmin.java new file mode 100644 index 0000000000..9be50bb145 --- /dev/null +++ b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbEdgeQueueAdmin.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.queue; + +public interface TbEdgeQueueAdmin extends TbQueueAdmin { + void syncEdgeNotificationsOffsets(String fatGroupId, String newGroupId); + + void deleteConsumerGroup(String consumerGroupId); +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/asset/AssetService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/asset/AssetService.java index 09bc8f1f93..c22c9c4140 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/asset/AssetService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/asset/AssetService.java @@ -66,6 +66,8 @@ public interface AssetService extends EntityDaoService { PageData findProfileEntityIdInfos(PageLink pageLink); + PageData findProfileEntityIdInfosByTenantId(TenantId tenantId, PageLink pageLink); + PageData findAssetIdsByTenantIdAndAssetProfileId(TenantId tenantId, AssetProfileId assetProfileId, PageLink pageLink); ListenableFuture> findAssetsByTenantIdAndIdsAsync(TenantId tenantId, List assetIds); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java index 3d91790790..5101d6d57e 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java @@ -39,6 +39,8 @@ public interface CalculatedFieldService extends EntityDaoService { PageData findAllCalculatedFields(PageLink pageLink); + PageData findCalculatedFieldsByTenantId(TenantId tenantId, PageLink pageLink); + PageData findAllCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink); void deleteCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId); @@ -53,6 +55,8 @@ public interface CalculatedFieldService extends EntityDaoService { List findAllCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + PageData findAllCalculatedFieldLinksByTenantId(TenantId tenantId, PageLink pageLink); + PageData findAllCalculatedFieldLinks(PageLink pageLink); boolean referencedInAnyCalculatedField(TenantId tenantId, EntityId referencedEntityId); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java index 4ef653855d..9eb258f182 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java @@ -76,6 +76,8 @@ public interface DeviceService extends EntityDaoService { PageData findProfileEntityIdInfos(PageLink pageLink); + PageData findProfileEntityIdInfosByTenantId(TenantId tenantId, PageLink pageLink); + PageData findDevicesByTenantIdAndType(TenantId tenantId, String type, PageLink pageLink); PageData findDeviceIdsByTenantIdAndDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId, PageLink pageLink); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/BaseEntityData.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/BaseEntityData.java index 33f32b9781..21b9ff0c4f 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/BaseEntityData.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/BaseEntityData.java @@ -125,7 +125,7 @@ public abstract class BaseEntityData implements EntityDa return switch (key) { case "createdTime" -> new LongDataPoint(System.currentTimeMillis(), fields.getCreatedTime()); case "edgeTemplate" -> new BoolDataPoint(System.currentTimeMillis(), fields.isEdgeTemplate()); - case "parentId" -> new StringDataPoint(System.currentTimeMillis(), getRelatedParentId(ctx)); + case "parentId" -> new StringDataPoint(System.currentTimeMillis(), getRelatedParentId(ctx), false); default -> new StringDataPoint(System.currentTimeMillis(), getField(key), false); }; } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedJsonDataPoint.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedJsonDataPoint.java index bce9d86875..c05a724e79 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedJsonDataPoint.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedJsonDataPoint.java @@ -17,10 +17,12 @@ package org.thingsboard.server.edqs.data.dp; import org.thingsboard.server.common.data.kv.DataType; +import java.util.function.Function; + public class CompressedJsonDataPoint extends CompressedStringDataPoint { - public CompressedJsonDataPoint(long ts, byte[] compressedValue) { - super(ts, compressedValue); + public CompressedJsonDataPoint(long ts, byte[] compressedValue, Function uncompressor) { + super(ts, compressedValue, uncompressor); } @Override diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedStringDataPoint.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedStringDataPoint.java index cf4267e443..45db6bf72a 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedStringDataPoint.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/data/dp/CompressedStringDataPoint.java @@ -19,17 +19,21 @@ import lombok.Getter; import lombok.SneakyThrows; import org.thingsboard.common.util.TbBytePool; import org.thingsboard.server.common.data.kv.DataType; -import org.xerial.snappy.Snappy; + +import java.util.function.Function; public class CompressedStringDataPoint extends AbstractDataPoint { @Getter private final byte[] compressedValue; + protected final Function uncompressor; + @SneakyThrows - public CompressedStringDataPoint(long ts, byte[] compressedValue) { + public CompressedStringDataPoint(long ts, byte[] compressedValue, Function uncompressor) { super(ts); this.compressedValue = TbBytePool.intern(compressedValue); + this.uncompressor = uncompressor; } @Override @@ -40,7 +44,7 @@ public class CompressedStringDataPoint extends AbstractDataPoint { @SneakyThrows @Override public String getStr() { - return Snappy.uncompressString(compressedValue); + return uncompressor.apply(compressedValue); } @Override diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index 07575220eb..78d17ef368 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -139,7 +139,7 @@ public class EdqsProcessor implements TbQueueHandler, } consumer.commit(); }) - .consumerCreator((config, partitionId) -> queueFactory.createEdqsEventsConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdqsEventsConsumer()) .queueAdmin(queueFactory.getEdqsQueueAdmin()) .consumerExecutor(consumersExecutor) .taskExecutor(taskExecutor) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/DefaultEdqsRepository.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/DefaultEdqsRepository.java index 215c64194f..88e4cc63d2 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/DefaultEdqsRepository.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/DefaultEdqsRepository.java @@ -16,6 +16,7 @@ package org.thingsboard.server.edqs.repo; import lombok.AllArgsConstructor; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.ObjectType; @@ -40,6 +41,7 @@ import java.util.function.Predicate; @Slf4j public class DefaultEdqsRepository implements EdqsRepository { + @Getter private final static ConcurrentMap repos = new ConcurrentHashMap<>(); private final EdqsStatsService statsService; @@ -52,6 +54,7 @@ public class DefaultEdqsRepository implements EdqsRepository { if (event.getEventType() == EdqsEventType.DELETED && event.getObjectType() == ObjectType.TENANT) { log.info("Tenant {} deleted", event.getTenantId()); repos.remove(event.getTenantId()); + statsService.reportRemoved(ObjectType.TENANT); } else { get(event.getTenantId()).processEvent(event); } @@ -61,8 +64,7 @@ public class DefaultEdqsRepository implements EdqsRepository { public long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query, boolean ignorePermissionCheck) { long startNs = System.nanoTime(); long result = get(tenantId).countEntitiesByQuery(customerId, query, ignorePermissionCheck); - double timingMs = (double) (System.nanoTime() - startNs) / 1000_000; - log.info("countEntitiesByQuery done in {} ms", timingMs); + statsService.reportEdqsCountQuery(tenantId, query, System.nanoTime() - startNs); return result; } @@ -71,8 +73,7 @@ public class DefaultEdqsRepository implements EdqsRepository { EntityDataQuery query, boolean ignorePermissionCheck) { long startNs = System.nanoTime(); var result = get(tenantId).findEntityDataByQuery(customerId, query, ignorePermissionCheck); - double timingMs = (double) (System.nanoTime() - startNs) / 1000_000; - log.info("findEntityDataByQuery done in {} ms", timingMs); + statsService.reportEdqsDataQuery(tenantId, query, System.nanoTime() - startNs); return result; } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java index a47559d6d8..ffba7583c9 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/repo/TenantRepo.java @@ -332,7 +332,6 @@ public class TenantRepo { public PageData findEntityDataByQuery(CustomerId customerId, EntityDataQuery oldQuery, boolean ignorePermissionCheck) { EdqsDataQuery query = RepositoryUtils.toNewQuery(oldQuery); - log.info("[{}][{}] findEntityDataByQuery: {}", tenantId, customerId, query); QueryContext ctx = buildContext(customerId, query.getEntityFilter(), ignorePermissionCheck); EntityQueryProcessor queryProcessor = EntityQueryProcessorFactory.create(this, ctx, query); return sortAndConvert(query, queryProcessor.processQuery(), ctx); @@ -340,7 +339,6 @@ public class TenantRepo { public long countEntitiesByQuery(CustomerId customerId, EntityCountQuery oldQuery, boolean ignorePermissionCheck) { EdqsQuery query = RepositoryUtils.toNewQuery(oldQuery); - log.info("[{}][{}] countEntitiesByQuery: {}", tenantId, customerId, query); QueryContext ctx = buildContext(customerId, query.getEntityFilter(), ignorePermissionCheck); EntityQueryProcessor queryProcessor = EntityQueryProcessorFactory.create(this, ctx, query); return queryProcessor.count(); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index efdb1ead1c..0efe6e7d3b 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -59,7 +59,8 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final KafkaEdqsQueueFactory queueFactory; - @Autowired @Lazy + @Autowired + @Lazy private EdqsProcessor edqsProcessor; private PartitionedQueueConsumerManager> stateConsumer; @@ -93,7 +94,7 @@ public class KafkaEdqsStateService implements EdqsStateService { } consumer.commit(); }) - .consumerCreator((config, partitionId) -> queueFactory.createEdqsStateConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdqsStateConsumer()) .queueAdmin(queueAdmin) .consumerExecutor(eventConsumer.getConsumerExecutor()) .taskExecutor(eventConsumer.getTaskExecutor()) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/stats/DefaultEdqsStatsService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/stats/DefaultEdqsStatsService.java index 3767d0f60b..435c7f0884 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/stats/DefaultEdqsStatsService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/stats/DefaultEdqsStatsService.java @@ -15,78 +15,117 @@ */ package org.thingsboard.server.edqs.stats; +import jakarta.annotation.PostConstruct; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; +import org.thingsboard.common.util.TbBytePool; +import org.thingsboard.common.util.TbStringPool; import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.stats.EdqsStatsService; +import org.thingsboard.server.common.stats.StatsCounter; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsTimer; -import org.thingsboard.server.common.stats.StatsType; +import org.thingsboard.server.edqs.repo.DefaultEdqsRepository; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @Service @Slf4j +@RequiredArgsConstructor @ConditionalOnExpression("'${queue.edqs.api.supported:true}' == 'true' && '${queue.edqs.stats.enabled:true}' == 'true'") public class DefaultEdqsStatsService implements EdqsStatsService { private final StatsFactory statsFactory; - @Value("${queue.edqs.stats.slow_query_threshold:3000}") + @Value("${queue.edqs.stats.slow_query_threshold}") private int slowQueryThreshold; - private final ConcurrentHashMap objectCounters = new ConcurrentHashMap<>(); - private final StatsTimer dataQueryTimer; - private final StatsTimer countQueryTimer; + private final ConcurrentMap objectCounters = new ConcurrentHashMap<>(); + private final ConcurrentMap timers = new ConcurrentHashMap<>(); + private final ConcurrentMap counters = new ConcurrentHashMap<>(); - private DefaultEdqsStatsService(StatsFactory statsFactory) { - this.statsFactory = statsFactory; - dataQueryTimer = statsFactory.createTimer(StatsType.EDQS, "entityDataQueryTimer"); - countQueryTimer = statsFactory.createTimer(StatsType.EDQS, "entityCountQueryTimer"); + @PostConstruct + private void init() { + statsFactory.createGauge("edqsMapGauges", "stringPoolSize", TbStringPool.getPool(), Map::size); + statsFactory.createGauge("edqsMapGauges", "bytePoolSize", TbBytePool.getPool(), Map::size); + statsFactory.createGauge("edqsMapGauges", "tenantReposSize", DefaultEdqsRepository.getRepos(), Map::size); } @Override public void reportAdded(ObjectType objectType) { - getObjectCounter(objectType).incrementAndGet(); + getObjectGauge(objectType).incrementAndGet(); } @Override public void reportRemoved(ObjectType objectType) { - getObjectCounter(objectType).decrementAndGet(); + getObjectGauge(objectType).decrementAndGet(); } @Override - public void reportDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) { - double timingMs = timingNanos / 1000_000.0; - if (timingMs < slowQueryThreshold) { - log.debug("[{}] Executed data query in {} ms: {}", tenantId, timingMs, query); - } else { - log.warn("[{}] Executed slow data query in {} ms: {}", tenantId, timingMs, query); - } - dataQueryTimer.record(timingNanos, TimeUnit.NANOSECONDS); + public void reportEntityDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) { + checkTiming(tenantId, query, timingNanos); + getTimer("entityDataQueryTimer").record(timingNanos, TimeUnit.NANOSECONDS); + } + + @Override + public void reportEntityCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) { + checkTiming(tenantId, query, timingNanos); + getTimer("entityCountQueryTimer").record(timingNanos, TimeUnit.NANOSECONDS); + } + + @Override + public void reportEdqsDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) { + checkTiming(tenantId, query, timingNanos); + getTimer("edqsDataQueryTimer").record(timingNanos, TimeUnit.NANOSECONDS); + } + + @Override + public void reportEdqsCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) { + checkTiming(tenantId, query, timingNanos); + getTimer("edqsCountQueryTimer").record(timingNanos, TimeUnit.NANOSECONDS); + } + + @Override + public void reportStringCompressed() { + getCounter("stringsCompressed").increment(); } @Override - public void reportCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) { + public void reportStringUncompressed() { + getCounter("stringsUncompressed").increment(); + } + + private void checkTiming(TenantId tenantId, EntityCountQuery query, long timingNanos) { double timingMs = timingNanos / 1000_000.0; + String queryType = query instanceof EntityDataQuery ? "data" : "count"; if (timingMs < slowQueryThreshold) { - log.debug("[{}] Executed count query in {} ms: {}", tenantId, timingMs, query); + log.debug("[{}] Executed " + queryType + " query in {} ms: {}", tenantId, timingMs, query); } else { - log.warn("[{}] Executed slow count query in {} ms: {}", tenantId, timingMs, query); + log.warn("[{}] Executed slow " + queryType + " query in {} ms: {}", tenantId, timingMs, query); } - countQueryTimer.record(timingNanos, TimeUnit.NANOSECONDS); } - private AtomicInteger getObjectCounter(ObjectType objectType) { + private StatsTimer getTimer(String name) { + return timers.computeIfAbsent(name, __ -> statsFactory.createStatsTimer("edqsTimers", name)); + } + + private StatsCounter getCounter(String name) { + return counters.computeIfAbsent(name, __ -> statsFactory.createStatsCounter("edqsCounters", name)); + } + + private AtomicInteger getObjectGauge(ObjectType objectType) { return objectCounters.computeIfAbsent(objectType, type -> - statsFactory.createGauge("edqsObjectsCount", new AtomicInteger(), "objectType", type.name())); + statsFactory.createGauge("edqsGauges", "objectsCount", new AtomicInteger(), "objectType", type.name())); } } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java index 167037b889..9a84436f36 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/EdqsConverter.java @@ -43,6 +43,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.stats.EdqsStatsService; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.edqs.data.dp.BoolDataPoint; import org.thingsboard.server.edqs.data.dp.CompressedJsonDataPoint; @@ -56,14 +57,18 @@ import org.thingsboard.server.gen.transport.TransportProtos.DataPointProto; import org.xerial.snappy.Snappy; import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.util.HashMap; import java.util.Map; import java.util.UUID; @Service +@RequiredArgsConstructor @Slf4j public class EdqsConverter { + private final EdqsStatsService edqsStatsService; + @Value("${queue.edqs.string_compression_length_threshold:512}") private int stringCompressionLengthThreshold; @@ -171,26 +176,44 @@ public class EdqsConverter { } else if (proto.hasDoubleV()) { return new DoubleDataPoint(ts, proto.getDoubleV()); } else if (proto.hasStringV()) { - return new StringDataPoint(ts, proto.getStringV()); + String stringV = proto.getStringV(); + if (stringV.length() < stringCompressionLengthThreshold) { + return new StringDataPoint(ts, stringV); + } else { + return new CompressedStringDataPoint(ts, compress(stringV), this::uncompress); + } } else if (proto.hasCompressedStringV()) { - return new CompressedStringDataPoint(ts, proto.getCompressedStringV().toByteArray()); + return new CompressedStringDataPoint(ts, proto.getCompressedStringV().toByteArray(), this::uncompress); } else if (proto.hasJsonV()) { - return new JsonDataPoint(ts, proto.getJsonV()); + String jsonV = proto.getJsonV(); + if (jsonV.length() < stringCompressionLengthThreshold) { + return new JsonDataPoint(ts, jsonV); + } else { + return new CompressedJsonDataPoint(ts, compress(jsonV), this::uncompress); + } } else if (proto.hasCompressedJsonV()) { - return new CompressedJsonDataPoint(ts, proto.getCompressedJsonV().toByteArray()); + return new CompressedJsonDataPoint(ts, proto.getCompressedJsonV().toByteArray(), this::uncompress); } else { throw new IllegalArgumentException("Unsupported data point proto: " + proto); } } @SneakyThrows - private static byte[] compress(String value) { - byte[] compressed = Snappy.compress(value); - // TODO: limit the size - log.debug("Compressed {} bytes to {} bytes", value.length(), compressed.length); + private byte[] compress(String value) { + byte[] compressed = Snappy.compress(value, StandardCharsets.UTF_8); + log.debug("Compressed {} chars to {} bytes", value.length(), compressed.length); + edqsStatsService.reportStringCompressed(); return compressed; } + @SneakyThrows + private String uncompress(byte[] compressed) { + String value = Snappy.uncompressString(compressed, StandardCharsets.UTF_8); + log.debug("Uncompressed {} bytes to {} chars", compressed.length, value.length()); + edqsStatsService.reportStringUncompressed(); + return value; + } + public static Entity toEntity(EntityType entityType, Object entity) { Entity edqsEntity = new Entity(); edqsEntity.setType(entityType); diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java index 178caf7961..f1c404ce16 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java @@ -136,6 +136,8 @@ public enum MsgType { EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG, + CF_CACHE_INIT_MSG, // Sent to init caches for CF actor; + CF_INIT_PROFILE_ENTITY_MSG, // Sent to init profile entities cache; CF_INIT_MSG, // Sent to init particular calculated field; CF_LINK_INIT_MSG, // Sent to init particular calculated field; CF_STATE_RESTORE_MSG, // Sent to restore particular calculated field entity state; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldCacheInitMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldCacheInitMsg.java new file mode 100644 index 0000000000..bf2054dfcf --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldCacheInitMsg.java @@ -0,0 +1,33 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.msg.cf; + +import lombok.Data; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; + +@Data +public class CalculatedFieldCacheInitMsg implements ToCalculatedFieldSystemMsg { + + private final TenantId tenantId; + + @Override + public MsgType getMsgType() { + return MsgType.CF_CACHE_INIT_MSG; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/cache/CalculatedFieldEntityProfileCache.java b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldInitProfileEntityMsg.java similarity index 51% rename from application/src/main/java/org/thingsboard/server/service/cf/cache/CalculatedFieldEntityProfileCache.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldInitProfileEntityMsg.java index bb5ef91974..66cd2ac441 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/cache/CalculatedFieldEntityProfileCache.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldInitProfileEntityMsg.java @@ -13,24 +13,24 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.cf.cache; +package org.thingsboard.server.common.msg.cf; -import org.springframework.context.ApplicationListener; +import lombok.Data; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; +import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; -import java.util.Collection; +@Data +public class CalculatedFieldInitProfileEntityMsg implements ToCalculatedFieldSystemMsg { -public interface CalculatedFieldEntityProfileCache extends ApplicationListener { + private final TenantId tenantId; + private final EntityId profileEntityId; + private final EntityId entityId; - void add(TenantId tenantId, EntityId profileId, EntityId entityId); + @Override + public MsgType getMsgType() { + return MsgType.CF_INIT_PROFILE_ENTITY_MSG; + } - void update(TenantId tenantId, EntityId oldProfileId, EntityId newProfileId, EntityId entityId); - - void evict(TenantId tenantId, EntityId entityId); - - Collection getMyEntityIdsByProfileId(TenantId tenantId, EntityId profileId); - - int getEntityIdPartition(TenantId tenantId, EntityId entityId); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java index ef9728344c..2233855a37 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java @@ -54,7 +54,7 @@ public class MainQueueConsumerManager msgPackProcessor; - protected final BiFunction> consumerCreator; + protected final BiFunction> consumerCreator; @Getter protected final ExecutorService consumerExecutor; @Getter @@ -74,7 +74,7 @@ public class MainQueueConsumerManager msgPackProcessor, - BiFunction> consumerCreator, + BiFunction> consumerCreator, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, @@ -313,7 +313,7 @@ public class MainQueueConsumerManager onStop.accept(tpi) : null; TbQueueConsumerTask consumer = new TbQueueConsumerTask<>(key, () -> { - TbQueueConsumer queueConsumer = consumerCreator.apply(config, partitionId); + TbQueueConsumer queueConsumer = consumerCreator.apply(config, tpi); if (startOffsetProvider != null && queueConsumer instanceof TbKafkaConsumerTemplate kafkaConsumer) { kafkaConsumer.setStartOffsetProvider(startOffsetProvider); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java index 0de1e53753..1b19fcab0e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java @@ -45,7 +45,7 @@ public class PartitionedQueueConsumerManager extends MainQ @Builder(builderMethodName = "create") // not to conflict with super.builder() public PartitionedQueueConsumerManager(QueueKey queueKey, String topic, long pollInterval, MsgPackProcessor msgPackProcessor, - BiFunction> consumerCreator, TbQueueAdmin queueAdmin, + BiFunction> consumerCreator, TbQueueAdmin queueAdmin, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, Consumer uncaughtErrorHandler) { super(queueKey, QueueConfig.of(true, pollInterval), msgPackProcessor, consumerCreator, consumerExecutor, scheduler, taskExecutor, uncaughtErrorHandler); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java index 401b451f59..3c927f135b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java @@ -38,7 +38,7 @@ public class EdqsConfig { private String requestsTopic; @Value("${queue.edqs.responses_topic:edqs.responses}") private String responsesTopic; - @Value("${queue.edqs.poll_interval:125}") + @Value("${queue.edqs.poll_interval:25}") private long pollInterval; @Value("${queue.edqs.max_pending_requests:10000}") private int maxPendingRequests; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index 37f881b49c..1e0064a5c8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -26,6 +26,7 @@ import org.apache.kafka.clients.admin.TopicDescription; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.TopicExistsException; +import org.thingsboard.server.queue.TbEdgeQueueAdmin; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.util.PropertyUtils; @@ -43,7 +44,7 @@ import java.util.stream.Collectors; * Created by ashvayka on 24.09.18. */ @Slf4j -public class TbKafkaAdmin implements TbQueueAdmin { +public class TbKafkaAdmin implements TbQueueAdmin, TbEdgeQueueAdmin { private final TbKafkaSettings settings; private final Map topicConfigs; @@ -149,17 +150,38 @@ public class TbKafkaAdmin implements TbQueueAdmin { * */ public void syncOffsets(String fatGroupId, String newGroupId, Integer partitionId) { try { - syncOffsetsUnsafe(fatGroupId, newGroupId, partitionId); + log.info("syncOffsets [{}][{}][{}]", fatGroupId, newGroupId, partitionId); + if (partitionId == null) { + return; + } + syncOffsetsUnsafe(fatGroupId, newGroupId, "." + partitionId); } catch (Exception e) { log.warn("Failed to syncOffsets from {} to {} partitionId {}", fatGroupId, newGroupId, partitionId, e); } } - void syncOffsetsUnsafe(String fatGroupId, String newGroupId, Integer partitionId) throws ExecutionException, InterruptedException, TimeoutException { - log.info("syncOffsets [{}][{}][{}]", fatGroupId, newGroupId, partitionId); - if (partitionId == null) { - return; + /** + * Sync edge notifications offsets from a fat group to a single group per edge + * */ + public void syncEdgeNotificationsOffsets(String fatGroupId, String newGroupId) { + try { + log.info("syncEdgeNotificationsOffsets [{}][{}]", fatGroupId, newGroupId); + syncOffsetsUnsafe(fatGroupId, newGroupId, newGroupId); + } catch (Exception e) { + log.warn("Failed to syncEdgeNotificationsOffsets from {} to {}", fatGroupId, newGroupId, e); + } + } + + @Override + public void deleteConsumerGroup(String consumerGroupId) { + try { + settings.getAdminClient().deleteConsumerGroups(Collections.singletonList(consumerGroupId)); + } catch (Exception e) { + log.warn("Failed to delete consumer group {}", consumerGroupId, e); } + } + + void syncOffsetsUnsafe(String fatGroupId, String newGroupId, String topicSuffix) throws ExecutionException, InterruptedException, TimeoutException { Map oldOffsets = getConsumerGroupOffsets(fatGroupId); if (oldOffsets.isEmpty()) { return; @@ -167,7 +189,7 @@ public class TbKafkaAdmin implements TbQueueAdmin { for (var consumerOffset : oldOffsets.entrySet()) { var tp = consumerOffset.getKey(); - if (!tp.topic().endsWith("." + partitionId)) { + if (!tp.topic().endsWith(topicSuffix)) { continue; } var om = consumerOffset.getValue(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java index 89d83af826..085d04f28c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java @@ -22,6 +22,7 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; @@ -133,7 +134,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer() { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index dcebe085b3..2269004f90 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -20,10 +20,12 @@ import jakarta.annotation.PreDestroy; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; @@ -44,6 +46,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg; +import org.thingsboard.server.queue.TbEdgeQueueAdmin; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; @@ -101,7 +104,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi private final TbQueueAdmin vcAdmin; private final TbQueueAdmin housekeeperAdmin; private final TbQueueAdmin housekeeperReprocessingAdmin; - private final TbQueueAdmin edgeAdmin; + private final TbEdgeQueueAdmin edgeAdmin; private final TbQueueAdmin edgeEventAdmin; private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfStateAdmin; @@ -493,9 +496,13 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi public TbQueueConsumer> createEdgeEventMsgConsumer(TenantId tenantId, EdgeId edgeId) { TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); - consumerBuilder.topic(topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic()); + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + + edgeAdmin.syncEdgeNotificationsOffsets(topicService.buildTopicName("monolith-edge-event-consumer"), topic); + + consumerBuilder.topic(topic); consumerBuilder.clientId("monolith-to-edge-event-consumer-" + serviceInfoProvider.getServiceId() + "-" + edgeConsumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("monolith-edge-event-consumer")); + consumerBuilder.groupId(topic); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeEventNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(edgeEventAdmin); consumerBuilder.statsService(consumerStatsService); @@ -513,15 +520,24 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer() { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { + String queueName = DataConstants.CF_QUEUE_NAME; + if (tpi == null) { + throw new IllegalArgumentException("TopicPartitionInfo is required."); + } + TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); + Integer partitionId = tpi.getPartition().orElseThrow(() -> new IllegalArgumentException("PartitionId is required.")); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); + TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); - consumerBuilder.clientId("monolith-calculated-field-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("monolith-calculated-field-consumer")); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.groupId(groupId); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(cfAdmin); consumerBuilder.statsService(consumerStatsService); + return consumerBuilder.build(); } @@ -571,7 +587,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi .readFromBeginning(true) .stopWhenRead(true) .clientId("monolith-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) - .groupId(topicService.buildTopicName("monolith-calculated-field-state-consumer")) + .groupId(null) // not using consumer group management .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())) .admin(cfStateAdmin) .statsService(consumerStatsService) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java index e9c42b0022..ea7c56f0aa 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java @@ -42,6 +42,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg; +import org.thingsboard.server.queue.TbEdgeQueueAdmin; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; @@ -99,7 +100,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { private final TbQueueAdmin vcAdmin; private final TbQueueAdmin housekeeperAdmin; private final TbQueueAdmin housekeeperReprocessingAdmin; - private final TbQueueAdmin edgeAdmin; + private final TbEdgeQueueAdmin edgeAdmin; private final TbQueueAdmin edgeEventAdmin; private final TbQueueAdmin cfAdmin; private final TbQueueAdmin edqsEventsAdmin; @@ -439,9 +440,13 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { public TbQueueConsumer> createEdgeEventMsgConsumer(TenantId tenantId, EdgeId edgeId) { TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); - consumerBuilder.topic(topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic()); + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + + edgeAdmin.syncEdgeNotificationsOffsets(topicService.buildTopicName("tb-core-edge-event-consumer"), topic); + + consumerBuilder.topic(topic); consumerBuilder.clientId("tb-core-edge-event-consumer-" + serviceInfoProvider.getServiceId() + "-" + edgeConsumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("tb-core-edge-event-consumer")); + consumerBuilder.groupId(topic); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeEventNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(edgeEventAdmin); consumerBuilder.statsService(consumerStatsService); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index d43ef5c9ac..b4884ae72c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -20,8 +20,11 @@ import jakarta.annotation.PreDestroy; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; @@ -313,12 +316,20 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer() { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { + String queueName = DataConstants.CF_QUEUE_NAME; + if (tpi == null) { + throw new IllegalArgumentException("TopicPartitionInfo is required."); + } + TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); + Integer partitionId = tpi.getPartition().orElseThrow(() -> new IllegalArgumentException("PartitionId is required.")); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); + TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); - consumerBuilder.clientId("tb-rule-engine-calculated-field-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("tb-rule-engine-calculated-field-consumer")); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.groupId(groupId); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(cfAdmin); consumerBuilder.statsService(consumerStatsService); @@ -371,7 +382,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { .readFromBeginning(true) .stopWhenRead(true) .clientId("tb-rule-engine-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) - .groupId(topicService.buildTopicName("tb-rule-engine-calculated-field-state-consumer")) + .groupId(null) // not using consumer group management .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())) .admin(cfStateAdmin) .statsService(consumerStatsService) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java index 329e3e346a..18bb6db14a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java @@ -16,6 +16,7 @@ package org.thingsboard.server.queue.provider; import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; @@ -95,8 +96,8 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory /** * Used to consume messages by TB Rule Engine Service * - * @return * @param configuration + * @return */ TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration); @@ -104,9 +105,9 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory * Used to consume messages by TB Rule Engine Service * Intended usage for consumer per partition strategy * - * @return TbQueueConsumer * @param configuration - * @param partitionId as a suffix for consumer name + * @param partitionId as a suffix for consumer name + * @return TbQueueConsumer */ default TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { return createToRuleEngineMsgConsumer(configuration); @@ -121,7 +122,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate(); - TbQueueConsumer> createToCalculatedFieldMsgConsumer(); + TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi); TbQueueAdmin getCalculatedFieldQueueAdmin(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCalculatedFieldSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCalculatedFieldSettings.java index c2de8eff4e..9a87e88416 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCalculatedFieldSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCalculatedFieldSettings.java @@ -31,5 +31,7 @@ public class TbQueueCalculatedFieldSettings { @Value("${queue.calculated_fields.state_topic}") private String stateTopic; + @Value("${queue.calculated_fields.init_tenant_fetch_pack_size:1000}") + private int initTenantFetchPackSize; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java index 46c29b867b..6a3aaa4a7f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java @@ -36,10 +36,11 @@ public @interface AfterStartUp { int STARTUP_SERVICE = 8; int ACTOR_SYSTEM = 9; - int REGULAR_SERVICE = 10; int CF_READ_PROFILE_ENTITIES_SERVICE = 10; - int CF_READ_CF_SERVICE = 11; + int CF_READ_CF_SERVICE = 10; + + int REGULAR_SERVICE = 11; int BEFORE_TRANSPORT_SERVICE = Integer.MAX_VALUE - 1001; int TRANSPORT_SERVICE = Integer.MAX_VALUE - 1000; diff --git a/common/stats/src/main/java/org/thingsboard/server/common/stats/DefaultStatsFactory.java b/common/stats/src/main/java/org/thingsboard/server/common/stats/DefaultStatsFactory.java index 8e2291270d..4d858be4bb 100644 --- a/common/stats/src/main/java/org/thingsboard/server/common/stats/DefaultStatsFactory.java +++ b/common/stats/src/main/java/org/thingsboard/server/common/stats/DefaultStatsFactory.java @@ -27,6 +27,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.StringUtils; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.ToDoubleFunction; @Service public class DefaultStatsFactory implements StatsFactory { @@ -86,6 +87,16 @@ public class DefaultStatsFactory implements StatsFactory { return meterRegistry.gauge(key, Tags.of(tags), number); } + @Override + public T createGauge(String type, String name, T number, String... tags) { + return createGauge(type, number, getTags(name, tags)); + } + + @Override + public void createGauge(String type, String name, S stateObject, ToDoubleFunction numberProvider, String... tags) { + meterRegistry.gauge(type, Tags.of(getTags(name, tags)), stateObject, numberProvider); + } + @Override public MessagesStats createMessagesStats(String key) { StatsCounter totalCounter = createStatsCounter(key, TOTAL_MSGS); @@ -106,8 +117,8 @@ public class DefaultStatsFactory implements StatsFactory { } @Override - public StatsTimer createTimer(StatsType type, String name, String... tags) { - return new StatsTimer(name, Timer.builder(type.getName()) + public StatsTimer createStatsTimer(String type, String name, String... tags) { + return new StatsTimer(name, Timer.builder(type) .tags(getTags(name, tags)) .register(meterRegistry)); } diff --git a/common/stats/src/main/java/org/thingsboard/server/common/stats/DummyEdqsStatsService.java b/common/stats/src/main/java/org/thingsboard/server/common/stats/DummyEdqsStatsService.java index df78e5fc89..3db1c6c790 100644 --- a/common/stats/src/main/java/org/thingsboard/server/common/stats/DummyEdqsStatsService.java +++ b/common/stats/src/main/java/org/thingsboard/server/common/stats/DummyEdqsStatsService.java @@ -33,9 +33,21 @@ public class DummyEdqsStatsService implements EdqsStatsService { public void reportRemoved(ObjectType objectType) {} @Override - public void reportDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) {} + public void reportEntityDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) {} @Override - public void reportCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) {} + public void reportEntityCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) {} + + @Override + public void reportEdqsDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos) {} + + @Override + public void reportEdqsCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos) {} + + @Override + public void reportStringCompressed() {} + + @Override + public void reportStringUncompressed() {} } diff --git a/common/stats/src/main/java/org/thingsboard/server/common/stats/EdqsStatsService.java b/common/stats/src/main/java/org/thingsboard/server/common/stats/EdqsStatsService.java index 106e43e913..3a6eeb26b8 100644 --- a/common/stats/src/main/java/org/thingsboard/server/common/stats/EdqsStatsService.java +++ b/common/stats/src/main/java/org/thingsboard/server/common/stats/EdqsStatsService.java @@ -26,8 +26,16 @@ public interface EdqsStatsService { void reportRemoved(ObjectType objectType); - void reportDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos); + void reportEntityDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos); - void reportCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos); + void reportEntityCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos); + + void reportEdqsDataQuery(TenantId tenantId, EntityDataQuery query, long timingNanos); + + void reportEdqsCountQuery(TenantId tenantId, EntityCountQuery query, long timingNanos); + + void reportStringCompressed(); + + void reportStringUncompressed(); } diff --git a/common/stats/src/main/java/org/thingsboard/server/common/stats/StatsFactory.java b/common/stats/src/main/java/org/thingsboard/server/common/stats/StatsFactory.java index bd46c09285..291b95e25a 100644 --- a/common/stats/src/main/java/org/thingsboard/server/common/stats/StatsFactory.java +++ b/common/stats/src/main/java/org/thingsboard/server/common/stats/StatsFactory.java @@ -17,6 +17,8 @@ package org.thingsboard.server.common.stats; import io.micrometer.core.instrument.Timer; +import java.util.function.ToDoubleFunction; + public interface StatsFactory { StatsCounter createStatsCounter(String key, String statsName, String... otherTags); @@ -25,10 +27,14 @@ public interface StatsFactory { T createGauge(String key, T number, String... tags); + T createGauge(String type, String name, T number, String... tags); + + void createGauge(String type, String name, S stateObject, ToDoubleFunction numberProvider, String... tags); + MessagesStats createMessagesStats(String key); Timer createTimer(String key, String... tags); - StatsTimer createTimer(StatsType type, String name, String... tags); + StatsTimer createStatsTimer(String type, String name, String... tags); } diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbBytePool.java b/common/util/src/main/java/org/thingsboard/common/util/TbBytePool.java index fe16a14e7c..1aadb49129 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/TbBytePool.java +++ b/common/util/src/main/java/org/thingsboard/common/util/TbBytePool.java @@ -16,12 +16,14 @@ package org.thingsboard.common.util; import com.google.common.hash.Hashing; +import lombok.Getter; import org.springframework.util.ConcurrentReferenceHashMap; import java.util.concurrent.ConcurrentMap; public class TbBytePool { + @Getter private static final ConcurrentMap pool = new ConcurrentReferenceHashMap<>(); public static byte[] intern(byte[] data) { @@ -32,8 +34,4 @@ public class TbBytePool { return pool.computeIfAbsent(checksum, c -> data); } - public static int size(){ - return pool.size(); - } - } diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStringPool.java b/common/util/src/main/java/org/thingsboard/common/util/TbStringPool.java index 38c010fbd3..167bf1bbed 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/TbStringPool.java +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStringPool.java @@ -15,12 +15,14 @@ */ package org.thingsboard.common.util; +import lombok.Getter; import org.springframework.util.ConcurrentReferenceHashMap; import java.util.concurrent.ConcurrentMap; public class TbStringPool { + @Getter private static final ConcurrentMap pool = new ConcurrentReferenceHashMap<>(); public static String intern(String data) { @@ -30,8 +32,4 @@ public class TbStringPool { return pool.computeIfAbsent(data, str -> str); } - public static int size(){ - return pool.size(); - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java index 36700ff59f..098dc4e83d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java @@ -240,4 +240,6 @@ public interface AssetDao extends Dao, TenantEntityDao, Exportable PageData findProfileEntityIdInfos(PageLink pageLink); + PageData findProfileEntityIdInfosByTenantId(UUID tenantId, PageLink pageLink); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index 21d5ea7f4e..fed201d403 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -293,6 +293,14 @@ public class BaseAssetService extends AbstractCachedEntityService findProfileEntityIdInfosByTenantId(TenantId tenantId, PageLink pageLink) { + log.trace("Executing findProfileEntityIdInfosByTenantId, tenantId[{}], pageLink [{}]", tenantId, pageLink); + validateId(tenantId, id -> INCORRECT_TENANT_ID + id); + validatePageLink(pageLink); + return assetDao.findProfileEntityIdInfosByTenantId(tenantId.getId(), pageLink); + } + @Override public PageData findAssetIdsByTenantIdAndAssetProfileId(TenantId tenantId, AssetProfileId assetProfileId, PageLink pageLink) { log.trace("Executing findAssetIdsByTenantIdAndAssetProfileId, tenantId [{}], assetProfileId [{}]", tenantId, assetProfileId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index dceef73aba..0c5df18e80 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java @@ -104,6 +104,14 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements return calculatedFieldDao.findAll(pageLink); } + @Override + public PageData findCalculatedFieldsByTenantId(TenantId tenantId, PageLink pageLink) { + log.trace("Executing findAllByTenantId, tenantId [{}], pageLink [{}]", tenantId, pageLink); + validateId(tenantId, id -> INCORRECT_TENANT_ID + id); + validatePageLink(pageLink); + return calculatedFieldDao.findAllByTenantId(tenantId, pageLink); + } + @Override public PageData findAllCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) { log.trace("Executing findAllByEntityId, entityId [{}], pageLink [{}]", entityId, pageLink); @@ -174,6 +182,14 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements return calculatedFieldLinkDao.findCalculatedFieldLinksByEntityId(tenantId, entityId); } + @Override + public PageData findAllCalculatedFieldLinksByTenantId(TenantId tenantId, PageLink pageLink) { + log.trace("Executing findAllCalculatedFieldLinksByTenantId, tenantId[{}] pageLink [{}]", tenantId, pageLink); + validateId(tenantId, id -> INCORRECT_TENANT_ID + id); + validatePageLink(pageLink); + return calculatedFieldLinkDao.findAllByTenantId(tenantId, pageLink); + } + @Override public PageData findAllCalculatedFieldLinks(PageLink pageLink) { log.trace("Executing findAllCalculatedFieldLinks, pageLink [{}]", pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java index a966977968..aadae93893 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java @@ -37,6 +37,8 @@ public interface CalculatedFieldDao extends Dao { PageData findAll(PageLink pageLink); + PageData findAllByTenantId(TenantId tenantId, PageLink pageLink); + PageData findAllByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink); List removeAllByEntityId(TenantId tenantId, EntityId entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java index 8b4a5e7086..dd184289ed 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java @@ -31,8 +31,12 @@ public interface CalculatedFieldLinkDao extends Dao { List findCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + List findCalculatedFieldLinksByTenantId(TenantId tenantId); + List findAll(); PageData findAll(PageLink pageLink); + PageData findAllByTenantId(TenantId tenantId, PageLink pageLink); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java index efc57119eb..6bc8903d20 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java @@ -23,9 +23,7 @@ import org.thingsboard.server.common.data.DeviceInfoFilter; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.ProfileEntityIdInfo; -import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.page.PageData; @@ -233,6 +231,8 @@ public interface DeviceDao extends Dao, TenantEntityDao, Exporta PageData findProfileEntityIdInfos(PageLink pageLink); + PageData findProfileEntityIdInfosByTenantId(UUID tenantId, PageLink pageLink); + PageData findDeviceInfosByFilter(DeviceInfoFilter filter, PageLink pageLink); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 2c890082a0..6d993f3e3d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -394,6 +394,14 @@ public class DeviceServiceImpl extends CachedVersionedEntityService findProfileEntityIdInfosByTenantId(TenantId tenantId, PageLink pageLink) { + log.trace("Executing findProfileEntityIdInfosByTenantId, tenantId[{}], pageLink [{}]", tenantId, pageLink); + validateId(tenantId, id -> INCORRECT_TENANT_ID + id); + validatePageLink(pageLink); + return deviceDao.findProfileEntityIdInfosByTenantId(tenantId.getId(), pageLink); + } + @Override public PageData findDevicesByTenantIdAndType(TenantId tenantId, String type, PageLink pageLink) { log.trace("Executing findDevicesByTenantIdAndType, tenantId [{}], type [{}], pageLink [{}]", tenantId, type, pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java index 4df33779ee..be7cbf7f84 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java @@ -20,7 +20,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.springframework.util.CollectionUtils; -import org.thingsboard.common.util.TbStopWatch; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.HasCustomerId; import org.thingsboard.server.common.data.HasEmail; @@ -43,6 +42,7 @@ import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityFilterType; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityListFilter; +import org.thingsboard.server.common.data.query.EntityNameFilter; import org.thingsboard.server.common.data.query.EntityTypeFilter; import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.RelationsQueryFilter; @@ -61,6 +61,7 @@ import java.util.function.Function; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.id.EntityId.NULL_UUID; +import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_NAME; import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_TYPE; import static org.thingsboard.server.dao.service.Validator.validateEntityDataPageLink; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -101,7 +102,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe validateId(customerId, id -> INCORRECT_CUSTOMER_ID + id); validateEntityCountQuery(query); - TbStopWatch stopWatch = TbStopWatch.create(); + long startNs = System.nanoTime(); Long result; if (edqsApiService.isEnabled() && validForEdqs(query) && !tenantId.isSysTenantId()) { EdqsRequest request = EdqsRequest.builder() @@ -112,7 +113,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } else { result = entityQueryDao.countEntitiesByQuery(tenantId, customerId, query); } - edqsStatsService.reportCountQuery(tenantId, query, stopWatch.stopAndGetTotalTimeNanos()); + edqsStatsService.reportEntityCountQuery(tenantId, query, System.nanoTime() - startNs); return result; } @@ -123,7 +124,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe validateId(customerId, id -> INCORRECT_CUSTOMER_ID + id); validateEntityDataQuery(query); - TbStopWatch stopWatch = TbStopWatch.create(); + long startNs = System.nanoTime(); PageData result; if (edqsApiService.isEnabled() && validForEdqs(query)) { EdqsRequest request = EdqsRequest.builder() @@ -146,7 +147,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } } } - edqsStatsService.reportDataQuery(tenantId, query, stopWatch.stopAndGetTotalTimeNanos()); + edqsStatsService.reportEntityDataQuery(tenantId, query, System.nanoTime() - startNs); return result; } @@ -250,6 +251,8 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe validateRelationQuery((RelationsQueryFilter) query.getEntityFilter()); } else if (query.getEntityFilter().getType().equals(ENTITY_TYPE)) { validateEntityTypeQuery((EntityTypeFilter) query.getEntityFilter()); + } else if (query.getEntityFilter().getType().equals(ENTITY_NAME)) { + validateEntityNameQuery((EntityNameFilter) query.getEntityFilter()); } } @@ -264,6 +267,12 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } } + private static void validateEntityNameQuery(EntityNameFilter filter) { + if (filter.getEntityType() == null) { + throw new IncorrectParameterException("Entity type is required"); + } + } + private static void validateRelationQuery(RelationsQueryFilter queryFilter) { if (queryFilter.isMultiRoot() && queryFilter.getMultiRootEntitiesType() == null) { throw new IncorrectParameterException("Multi-root relation query filter should contain 'multiRootEntitiesType'"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java index c61d894de5..4e99fb57e4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java @@ -39,14 +39,12 @@ import org.thingsboard.server.dao.model.sql.AssetEntity; import org.thingsboard.server.dao.model.sql.AssetInfoEntity; import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.sql.device.NativeAssetRepository; -import org.thingsboard.server.dao.sql.device.NativeDeviceRepository; import org.thingsboard.server.dao.util.SqlDao; import java.util.Arrays; import java.util.List; import java.util.Optional; import java.util.UUID; -import java.util.stream.Collectors; import static org.thingsboard.server.dao.DaoUtil.convertTenantEntityInfosToDto; @@ -262,10 +260,16 @@ public class JpaAssetDao extends JpaAbstractDao implements A @Override public PageData findProfileEntityIdInfos(PageLink pageLink) { - log.debug("Find profile device id infos by pageLink [{}]", pageLink); + log.debug("Find profile asset id infos by pageLink [{}]", pageLink); return nativeAssetRepository.findProfileEntityIdInfos(DaoUtil.toPageable(pageLink)); } + @Override + public PageData findProfileEntityIdInfosByTenantId(UUID tenantId, PageLink pageLink) { + log.debug("Find profile asset id infos by pageLink [{}]", pageLink); + return nativeAssetRepository.findProfileEntityIdInfosByTenantId(tenantId, DaoUtil.toPageable(pageLink)); + } + @Override public Long countByTenantId(TenantId tenantId) { return assetRepository.countByTenantId(tenantId.getId()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java index 584a3b5199..6f6a0775f3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.dao.sql.cf; +import org.springframework.data.domain.Page; +import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.JpaRepository; import org.thingsboard.server.dao.model.sql.CalculatedFieldLinkEntity; @@ -27,4 +29,8 @@ public interface CalculatedFieldLinkRepository extends JpaRepository findAllByTenantIdAndEntityId(UUID tenantId, UUID entityId); + List findAllByTenantId(UUID tenantId); + + Page findAllByTenantId(UUID tenantId, Pageable pageable); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java index 0f48f3b00d..be122816ba 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java @@ -32,6 +32,8 @@ public interface CalculatedFieldRepository extends JpaRepository findAllByTenantIdAndEntityId(UUID tenantId, UUID entityId); + Page findAllByTenantId(UUID tenantId, Pageable pageable); + Page findAllByTenantIdAndEntityId(UUID tenantId, UUID entityId, Pageable pageable); List findAllByTenantId(UUID tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java index 8922eaca4e..4bb52c29db 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java @@ -71,6 +71,12 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao findAllByTenantId(TenantId tenantId, PageLink pageLink) { + log.debug("Try to find calculated fields by tenantId[{}] and pageLink [{}]", tenantId, pageLink); + return DaoUtil.toPageData(calculatedFieldRepository.findAllByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink))); + } + @Override public PageData findAllByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) { log.debug("Try to find calculated fields by entityId[{}] and pageLink [{}]", entityId, pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java index dbb2fd87da..38871f2fb8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java @@ -54,6 +54,11 @@ public class JpaCalculatedFieldLinkDao extends JpaAbstractDao findCalculatedFieldLinksByTenantId(TenantId tenantId) { + return DaoUtil.convertDataList(calculatedFieldLinkRepository.findAllByTenantId(tenantId.getId())); + } + @Override public List findAll() { return DaoUtil.convertDataList(calculatedFieldLinkRepository.findAll()); @@ -65,6 +70,12 @@ public class JpaCalculatedFieldLinkDao extends JpaAbstractDao findAllByTenantId(TenantId tenantId, PageLink pageLink) { + log.debug("Try to find calculated field links by tenantId [{}], pageLink [{}]", tenantId, pageLink); + return DaoUtil.toPageData(calculatedFieldLinkRepository.findAllByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink))); + } + @Override protected Class getEntityClass() { return CalculatedFieldLinkEntity.class; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeAssetRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeAssetRepository.java index 43d66e2ff0..ec47da3499 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeAssetRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeAssetRepository.java @@ -20,12 +20,9 @@ import org.springframework.data.domain.Pageable; import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; import org.springframework.stereotype.Repository; import org.springframework.transaction.support.TransactionTemplate; -import org.thingsboard.server.common.data.DeviceIdInfo; import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -43,8 +40,8 @@ public class DefaultNativeAssetRepository extends AbstractNativeRepository imple @Override public PageData findProfileEntityIdInfos(Pageable pageable) { - String PROFILE_DEVICE_ID_INFO_QUERY = "SELECT tenant_id as tenantId, asset_profile_id as profileId, id as id FROM asset ORDER BY created_time ASC LIMIT %s OFFSET %s"; - return find(COUNT_QUERY, PROFILE_DEVICE_ID_INFO_QUERY, pageable, row -> { + String PROFILE_ASSET_ID_INFO_QUERY = "SELECT tenant_id as tenantId, asset_profile_id as profileId, id as id FROM asset ORDER BY created_time ASC LIMIT %s OFFSET %s"; + return find(COUNT_QUERY, PROFILE_ASSET_ID_INFO_QUERY, pageable, row -> { AssetId id = new AssetId((UUID) row.get("id")); AssetProfileId profileId = new AssetProfileId((UUID) row.get("profileId")); var tenantIdObj = row.get("tenantId"); @@ -52,4 +49,14 @@ public class DefaultNativeAssetRepository extends AbstractNativeRepository imple }); } + @Override + public PageData findProfileEntityIdInfosByTenantId(UUID tenantId, Pageable pageable) { + String PROFILE_ASSET_ID_INFO_QUERY = String.format("SELECT tenant_id as tenantId, asset_profile_id as profileId, id as id FROM asset WHERE tenant_id = '%s' ORDER BY created_time ASC LIMIT %%s OFFSET %%s", tenantId); + return find(COUNT_QUERY, PROFILE_ASSET_ID_INFO_QUERY, pageable, row -> { + AssetId id = new AssetId((UUID) row.get("id")); + AssetProfileId profileId = new AssetProfileId((UUID) row.get("profileId")); + var tenantIdObj = row.get("tenantId"); + return ProfileEntityIdInfo.create(tenantIdObj != null ? (UUID) tenantIdObj : TenantId.SYS_TENANT_ID.getId(), profileId, id); + }); + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeDeviceRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeDeviceRepository.java index 776dedc2d5..78ee2795b0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeDeviceRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DefaultNativeDeviceRepository.java @@ -61,4 +61,15 @@ public class DefaultNativeDeviceRepository extends AbstractNativeRepository impl }); } + @Override + public PageData findProfileEntityIdInfosByTenantId(UUID tenantId, Pageable pageable) { + String PROFILE_DEVICE_ID_INFO_QUERY = String.format("SELECT tenant_id as tenantId, device_profile_id as profileId, id as id FROM device WHERE tenant_id = '%s' ORDER BY created_time ASC LIMIT %%s OFFSET %%s", tenantId); + return find(COUNT_QUERY, PROFILE_DEVICE_ID_INFO_QUERY, pageable, row -> { + DeviceId id = new DeviceId((UUID) row.get("id")); + DeviceProfileId profileId = new DeviceProfileId((UUID) row.get("profileId")); + var tenantIdObj = row.get("tenantId"); + return ProfileEntityIdInfo.create(tenantIdObj != null ? (UUID) tenantIdObj : TenantId.SYS_TENANT_ID.getId(), profileId, id); + }); + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java index 34835f52f1..9c79637e39 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java @@ -179,11 +179,11 @@ public class JpaDeviceDao extends JpaAbstractDao implement @Override public PageData findDeviceIdsByTenantIdAndDeviceProfileId(UUID tenantId, UUID deviceProfileId, PageLink pageLink) { return DaoUtil.pageToPageData( - deviceRepository.findIdsByTenantIdAndDeviceProfileId( - tenantId, - deviceProfileId, - pageLink.getTextSearch(), - DaoUtil.toPageable(pageLink))) + deviceRepository.findIdsByTenantIdAndDeviceProfileId( + tenantId, + deviceProfileId, + pageLink.getTextSearch(), + DaoUtil.toPageable(pageLink))) .mapData(DeviceId::new); } @@ -281,6 +281,12 @@ public class JpaDeviceDao extends JpaAbstractDao implement return nativeDeviceRepository.findProfileEntityIdInfos(DaoUtil.toPageable(pageLink)); } + @Override + public PageData findProfileEntityIdInfosByTenantId(UUID tenantId, PageLink pageLink) { + log.debug("Find profile device id infos by tenantId[{}], pageLink [{}]", tenantId, pageLink); + return nativeDeviceRepository.findProfileEntityIdInfosByTenantId(tenantId, DaoUtil.toPageable(pageLink)); + } + @Override public Device findByTenantIdAndExternalId(UUID tenantId, UUID externalId) { return DaoUtil.getData(deviceRepository.findByTenantIdAndExternalId(tenantId, externalId)); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/NativeProfileEntityRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/NativeProfileEntityRepository.java index 750f0c8787..3568f01851 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/NativeProfileEntityRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/NativeProfileEntityRepository.java @@ -19,8 +19,12 @@ import org.springframework.data.domain.Pageable; import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.page.PageData; +import java.util.UUID; + public interface NativeProfileEntityRepository { PageData findProfileEntityIdInfos(Pageable pageable); + PageData findProfileEntityIdInfosByTenantId(UUID tenantId, Pageable pageable); + } diff --git a/edqs/src/main/resources/edqs.yml b/edqs/src/main/resources/edqs.yml index 2c0e0b8c5d..1cc32a4230 100644 --- a/edqs/src/main/resources/edqs.yml +++ b/edqs/src/main/resources/edqs.yml @@ -66,7 +66,7 @@ queue: # EDQS responses topic responses_topic: "${TB_EDQS_RESPONSES_TOPIC:edqs.responses}" # Poll interval for EDQS topics - poll_interval: "${TB_EDQS_POLL_INTERVAL_MS:125}" + poll_interval: "${TB_EDQS_POLL_INTERVAL_MS:25}" # Maximum amount of pending requests to EDQS max_pending_requests: "${TB_EDQS_MAX_PENDING_REQUESTS:10000}" # Maximum timeout for requests to EDQS @@ -76,6 +76,8 @@ queue: stats: # Enable/disable statistics for EDQS enabled: "${TB_EDQS_STATS_ENABLED:true}" + # Threshold for slow queries to log, in milliseconds + slow_query_threshold: "${TB_EDQS_SLOW_QUERY_THRESHOLD_MS:200}" kafka: # Kafka Bootstrap nodes in "host:port" format diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts index 9490db52a4..60d05d4c86 100644 --- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts +++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts @@ -1118,6 +1118,7 @@ export class DashboardPageComponent extends PageComponent implements IDashboardC this.dashboardCtx.aliasController.dashboardStateChanged(); this.isRightLayoutOpened = openRightLayout ? true : false; this.updateLayouts(layoutsData); + this.cd.markForCheck(); } setTimeout(() => { this.mobileService.onDashboardLoaded(this.layouts.right.show, this.isRightLayoutOpened); diff --git a/ui-ngx/src/app/modules/home/components/widget/action/manage-widget-actions.component.scss b/ui-ngx/src/app/modules/home/components/widget/action/manage-widget-actions.component.scss index 66cbe675b2..e122794bb0 100644 --- a/ui-ngx/src/app/modules/home/components/widget/action/manage-widget-actions.component.scss +++ b/ui-ngx/src/app/modules/home/components/widget/action/manage-widget-actions.component.scss @@ -33,6 +33,7 @@ .table-container { overflow: auto; + z-index: 1; } } } diff --git a/ui-ngx/src/app/modules/home/components/widget/config/timewindow-config-panel.component.html b/ui-ngx/src/app/modules/home/components/widget/config/timewindow-config-panel.component.html index ff0acd5fca..138dde324e 100644 --- a/ui-ngx/src/app/modules/home/components/widget/config/timewindow-config-panel.component.html +++ b/ui-ngx/src/app/modules/home/components/widget/config/timewindow-config-panel.component.html @@ -33,6 +33,7 @@ noMargin isEdit="true" alwaysDisplayTypePrefix + timezone="true" [historyOnly]="onlyHistoryTimewindow" quickIntervalOnly="{{ widgetType === widgetTypes.latest }}" aggregation="{{ widgetType === widgetTypes.timeseries }}" diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/scada/scada-symbol.models.ts b/ui-ngx/src/app/modules/home/components/widget/lib/scada/scada-symbol.models.ts index e52c6f9f54..abe570fa3c 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/scada/scada-symbol.models.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/scada/scada-symbol.models.ts @@ -271,7 +271,7 @@ export const updateScadaSymbolMetadataInContent = (svgContent: string, metadata: const svgDoc = new DOMParser().parseFromString(svgContent, 'image/svg+xml'); const parsererror = svgDoc.getElementsByTagName('parsererror'); if (parsererror?.length) { - return parsererror[0].outerHTML; + throw Error(parsererror[0].textContent) } updateScadaSymbolMetadataInDom(svgDoc, metadata); return svgDoc.documentElement.outerHTML; diff --git a/ui-ngx/src/app/modules/home/models/widget-component.models.ts b/ui-ngx/src/app/modules/home/models/widget-component.models.ts index 00afa2b68d..aee890a5b2 100644 --- a/ui-ngx/src/app/modules/home/models/widget-component.models.ts +++ b/ui-ngx/src/app/modules/home/models/widget-component.models.ts @@ -115,6 +115,7 @@ import { DataKeySettingsFunction } from '@home/components/widget/lib/settings/co import { UtilsService } from '@core/services/utils.service'; import { CompiledTbFunction } from '@shared/models/js-function.models'; import { FormProperty } from '@shared/models/dynamic-form.models'; +import { ExportableEntity } from '@shared/models/base-data'; export interface IWidgetAction { name: string; @@ -573,7 +574,7 @@ export interface IDynamicWidgetComponent { [key: string]: any; } -export interface WidgetInfo extends WidgetTypeDescriptor, WidgetControllerDescriptor { +export interface WidgetInfo extends WidgetTypeDescriptor, WidgetControllerDescriptor, ExportableEntity { widgetName: string; fullFqn: string; deprecated: boolean; @@ -681,6 +682,7 @@ export const toWidgetInfo = (widgetTypeEntity: WidgetType): WidgetInfo => ({ fullFqn: fullWidgetTypeFqn(widgetTypeEntity), deprecated: widgetTypeEntity.deprecated, scada: widgetTypeEntity.scada, + externalId: widgetTypeEntity.externalId, type: widgetTypeEntity.descriptor.type, sizeX: widgetTypeEntity.descriptor.sizeX, sizeY: widgetTypeEntity.descriptor.sizeY, @@ -736,6 +738,7 @@ export const toWidgetType = (widgetInfo: WidgetInfo, id: WidgetTypeId, tenantId: name: widgetInfo.widgetName, deprecated: widgetInfo.deprecated, scada: widgetInfo.scada, + externalId: widgetInfo.externalId, descriptor }; }; diff --git a/ui-ngx/src/app/modules/home/pages/admin/resource/js-resource.component.html b/ui-ngx/src/app/modules/home/pages/admin/resource/js-resource.component.html index 6287d27f2c..62036a8ea1 100644 --- a/ui-ngx/src/app/modules/home/pages/admin/resource/js-resource.component.html +++ b/ui-ngx/src/app/modules/home/pages/admin/resource/js-resource.component.html @@ -80,6 +80,7 @@ (fileNameChanged)="entityForm?.get('fileName').patchValue($event)"> - {{ 'javascript.module-script' | translate }} - + {{ 'javascript.module-script' | translate }} { - imageInfo.title = metadata.title; - return this.imageService.updateImageInfo(imageInfo); - }) - ); - } - imageInfoObservable.pipe( - switchMap(imageInfo => this.imageService.getImageString( + try { + const scadaSymbolContent = this.prepareScadaSymbolContent(metadata); + const file = createFileFromContent(scadaSymbolContent, this.symbolData.imageResource.fileName, + this.symbolData.imageResource.descriptor.mediaType); + const type = imageResourceType(this.symbolData.imageResource); + let imageInfoObservable = + this.imageService.updateImage(type, this.symbolData.imageResource.resourceKey, file); + if (metadata.title !== this.symbolData.imageResource.title) { + imageInfoObservable = imageInfoObservable.pipe( + switchMap(imageInfo => { + imageInfo.title = metadata.title; + return this.imageService.updateImageInfo(imageInfo); + }) + ); + } + imageInfoObservable.pipe( + switchMap(imageInfo => this.imageService.getImageString( `${IMAGES_URL_PREFIX}/${type}/${encodeURIComponent(imageInfo.resourceKey)}`).pipe( - map(content => ({ - imageResource: imageInfo, - scadaSymbolContent: content - })) + map(content => ({ + imageResource: imageInfo, + scadaSymbolContent: content + })) )) - ).subscribe(data => { - this.init(data); - this.updateBreadcrumbs.emit(); - }); + ).subscribe(data => { + this.init(data); + this.updateBreadcrumbs.emit(); + }); + } catch (e) { + this.store.dispatch(new ActionNotificationShow({ message: e.message, type: 'error' })); + } } } @@ -248,43 +253,47 @@ export class ScadaSymbolComponent extends PageComponent enterPreviewMode() { this.previewMetadata = this.scadaSymbolFormGroup.get('metadata').value; - this.symbolData.scadaSymbolContent = this.prepareScadaSymbolContent(this.previewMetadata); - this.previewScadaSymbolObjectSettings = { - behavior: {}, - properties: {} - }; - this.scadaPreviewFormGroup.patchValue({ - scadaSymbolObjectSettings: this.previewScadaSymbolObjectSettings - }, {emitEvent: false}); - this.scadaPreviewFormGroup.markAsPristine(); - const settings: ScadaSymbolWidgetSettings = {...scadaSymbolWidgetDefaultSettings, - ...{ + try { + this.symbolData.scadaSymbolContent = this.prepareScadaSymbolContent(this.previewMetadata); + this.previewScadaSymbolObjectSettings = { + behavior: {}, + properties: {} + }; + this.scadaPreviewFormGroup.patchValue({ + scadaSymbolObjectSettings: this.previewScadaSymbolObjectSettings + }, {emitEvent: false}); + this.scadaPreviewFormGroup.markAsPristine(); + const settings: ScadaSymbolWidgetSettings = {...scadaSymbolWidgetDefaultSettings, + ...{ simulated: true, scadaSymbolUrl: null, scadaSymbolContent: this.symbolData.scadaSymbolContent, scadaSymbolObjectSettings: this.previewScadaSymbolObjectSettings, padding: '0', background: colorBackground('rgba(0,0,0,0)') - } - }; - this.previewWidget = { - typeFullFqn: 'system.scada_symbol', - type: widgetType.rpc, - sizeX: this.previewMetadata.widgetSizeX || 3, - sizeY: this.previewMetadata.widgetSizeY || 3, - row: 0, - col: 0, - config: { - settings, - showTitle: false, - dropShadow: false, - padding: '0', - margin: '0', - backgroundColor: 'rgba(0,0,0,0)' - } - }; - this.previewWidgets = [this.previewWidget]; - this.previewMode = true; + } + }; + this.previewWidget = { + typeFullFqn: 'system.scada_symbol', + type: widgetType.rpc, + sizeX: this.previewMetadata.widgetSizeX || 3, + sizeY: this.previewMetadata.widgetSizeY || 3, + row: 0, + col: 0, + config: { + settings, + showTitle: false, + dropShadow: false, + padding: '0', + margin: '0', + backgroundColor: 'rgba(0,0,0,0)' + } + }; + this.previewWidgets = [this.previewWidget]; + this.previewMode = true; + } catch (e) { + this.store.dispatch(new ActionNotificationShow({ message: e.message, type: 'error' })); + } } exitPreviewMode() { @@ -374,19 +383,23 @@ export class ScadaSymbolComponent extends PageComponent metadata = parseScadaSymbolMetadataFromContent(this.origSymbolData.scadaSymbolContent); } const linkElement = document.createElement('a'); - const scadaSymbolContent = this.prepareScadaSymbolContent(metadata); - const blob = new Blob([scadaSymbolContent], { type: this.symbolData.imageResource.descriptor.mediaType }); - const url = URL.createObjectURL(blob); - linkElement.setAttribute('href', url); - linkElement.setAttribute('download', this.symbolData.imageResource.fileName); - const clickEvent = new MouseEvent('click', - { - view: window, - bubbles: true, - cancelable: false - } - ); - linkElement.dispatchEvent(clickEvent); + try { + const scadaSymbolContent = this.prepareScadaSymbolContent(metadata); + const blob = new Blob([scadaSymbolContent], { type: this.symbolData.imageResource.descriptor.mediaType }); + const url = URL.createObjectURL(blob); + linkElement.setAttribute('href', url); + linkElement.setAttribute('download', this.symbolData.imageResource.fileName); + const clickEvent = new MouseEvent('click', + { + view: window, + bubbles: true, + cancelable: false + } + ); + linkElement.dispatchEvent(clickEvent); + } catch (e) { + this.store.dispatch(new ActionNotificationShow({ message: e.message, type: 'error' })); + } } createWidget() { diff --git a/ui-ngx/src/app/shared/components/image/upload-image-dialog.component.html b/ui-ngx/src/app/shared/components/image/upload-image-dialog.component.html index bbaf7d88ae..47950e53b7 100644 --- a/ui-ngx/src/app/shared/components/image/upload-image-dialog.component.html +++ b/ui-ngx/src/app/shared/components/image/upload-image-dialog.component.html @@ -28,8 +28,8 @@
-
-
+
+
{ - this.dialogRef.close({image: Object.assign(imageInfo, {base64})}); - }); - } else { - if (this.isScada) { - blobToText(file).subscribe(scadaSymbolContent => { - this.dialogRef.close({scadaSymbolContent}); - }); - } else { - const image = this.data.image; forkJoin([ - this.imageService.updateImage(imageResourceType(image), image.resourceKey, file), + this.imageService.uploadImage(file, title, this.data.imageSubType), blobToBase64(file) ]).subscribe(([imageInfo, base64]) => { - this.dialogRef.close({image:Object.assign(imageInfo, {base64})}); + this.dialogRef.close({image: Object.assign(imageInfo, {base64})}); }); + } else { + if (this.isScada) { + blobToText(file).subscribe(scadaSymbolContent => { + this.dialogRef.close({scadaSymbolContent}); + }); + } else { + const image = this.data.image; + forkJoin([ + this.imageService.updateImage(imageResourceType(image), image.resourceKey, file), + blobToBase64(file) + ]).subscribe(([imageInfo, base64]) => { + this.dialogRef.close({image:Object.assign(imageInfo, {base64})}); + }); + } } + } catch (e) { + this.store.dispatch(new ActionNotificationShow({ + message: e.message, + type: 'error', + verticalPosition: 'top', + horizontalPosition: 'right', + target: 'uploadRoot' + })); } } } diff --git a/ui-ngx/src/app/shared/components/js-func.component.scss b/ui-ngx/src/app/shared/components/js-func.component.scss index 44aab1abbf..7c43eaccf2 100644 --- a/ui-ngx/src/app/shared/components/js-func.component.scss +++ b/ui-ngx/src/app/shared/components/js-func.component.scss @@ -85,4 +85,8 @@ color: rgb(49, 132, 149); } } + + label.tb-title.tb-required::after { + content: "*"; + } } diff --git a/ui-ngx/src/app/shared/components/markdown.component.scss b/ui-ngx/src/app/shared/components/markdown.component.scss index 4989a77021..65ad64b189 100644 --- a/ui-ngx/src/app/shared/components/markdown.component.scss +++ b/ui-ngx/src/app/shared/components/markdown.component.scss @@ -270,7 +270,7 @@ outline: none; position: absolute; width: 206px; - height: 42px; + height: 32px; top: 0; right: 32px; background: 0 0; @@ -281,11 +281,11 @@ user-select: none; &.multiline { - right: 44px; + right: 38px; } p { - padding: 8px; + padding: 8px 8px 0; top: 1px; transition: .2s; color: #2a7dec; @@ -301,10 +301,10 @@ background-color: #fff; position: absolute; width: 38px; - height: 38px; + height: 28px; top: 3px; right: 3px; - padding: 10px; + padding: 10px 10px 0; img { position: initial; diff --git a/ui-ngx/src/app/shared/models/widget.models.ts b/ui-ngx/src/app/shared/models/widget.models.ts index ec55d4e8a3..c55dcdee37 100644 --- a/ui-ngx/src/app/shared/models/widget.models.ts +++ b/ui-ngx/src/app/shared/models/widget.models.ts @@ -206,7 +206,7 @@ export interface WidgetControllerDescriptor { actionSources?: {[actionSourceId: string]: WidgetActionSource}; } -export interface BaseWidgetType extends BaseData, HasTenantId, HasVersion { +export interface BaseWidgetType extends BaseData, HasTenantId, HasVersion, ExportableEntity { tenantId: TenantId; fqn: string; name: string; diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 984d364912..72afcfc0b9 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -5258,7 +5258,7 @@ "advanced-mode": "Advanced", "save-time-series": { "processing-settings": "Processing settings", - "processing-settings-hint": "Define how incoming messages are processed. In Basic mode, select a preconfigured processing strategy or enable only WebSocket updates. Advanced mode allows you to select individual processing strategies for each action.", + "processing-settings-hint": "Define how incoming messages are processed. Basic processing settings allow you to select preconfigured strategies, while Advanced settings allow you to select individual processing strategies for each action.", "advanced-settings-hint": "Be cautious when configuring processing strategies. Certain combinations can lead to unexpected behavior.", "strategy": "Strategy", "deduplication-interval": "Deduplication interval", @@ -5277,7 +5277,7 @@ }, "save-attribute": { "processing-settings": "Processing settings", - "processing-settings-hint": "Define how incoming messages are processed. In Basic mode, select a preconfigured processing strategy or enable only WebSocket updates. Advanced mode allows you to select individual processing strategies for each action.", + "processing-settings-hint": "Define how incoming messages are processed. Basic processing settings allow you to select preconfigured strategies, while Advanced settings allow you to select individual processing strategies for each action.", "advanced-settings-hint": "Be cautious when configuring processing strategies. Certain combinations can lead to unexpected behavior.", "strategy": "Strategy", "deduplication-interval": "Deduplication interval",