From 85119d02475a5cbfb6688bc609414f2d3d5bf5e0 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 24 Jan 2025 11:37:32 +0200 Subject: [PATCH] WIP: cluster mode implementation --- .../cf/CalculatedFieldExecutionService.java | 4 +- ...efaultCalculatedFieldExecutionService.java | 11 +++-- ...faultTbCalculatedFieldConsumerService.java | 6 +-- .../queue/DefaultTbClusterService.java | 42 +++++++------------ .../DefaultTelemetrySubscriptionService.java | 3 +- .../src/main/resources/thingsboard.yml | 2 +- 6 files changed, 29 insertions(+), 39 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index 1ad0377e27..93a0ecc75f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -23,6 +23,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntit import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; +import java.util.List; + public interface CalculatedFieldExecutionService { /** @@ -32,7 +34,7 @@ public interface CalculatedFieldExecutionService { */ void pushRequestToQueue(TimeseriesSaveRequest request, TimeseriesSaveResult result); - void pushRequestToQueue(AttributesSaveRequest request); + void pushRequestToQueue(AttributesSaveRequest request, List result); // void pushEntityUpdateMsg(TransportProtos.CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index cbd5445cdf..d46a7515ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -197,14 +197,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } @Override - public void pushRequestToQueue(AttributesSaveRequest request) { + public void pushRequestToQueue(AttributesSaveRequest request, List result) { var tenantId = request.getTenantId(); var entityId = request.getEntityId(); checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries(), request.getScope()), cf -> cf.linkMatches(entityId, request.getEntries(), request.getScope()), - () -> toCalculatedFieldTelemetryMsgProto(request), request.getCallback()); + () -> toCalculatedFieldTelemetryMsgProto(request, result), request.getCallback()); } - private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, Predicate mainEntityFilter, Predicate linkedEntityFilter, Supplier msg, FutureCallback callback) { + private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, + Predicate mainEntityFilter, Predicate linkedEntityFilter, + Supplier msg, FutureCallback callback) { boolean send = checkEntityForCalculatedFields(tenantId, entityId, mainEntityFilter, linkedEntityFilter); if (send) { clusterService.pushMsgToCalculatedFields(tenantId, entityId, msg.get(), wrap(callback)); @@ -827,7 +829,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return msg.build(); } - private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request) { + private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request, List result) { + //TODO: IM Use result in both methods to update the versions of telemetry/attributes. ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder(); CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds()); 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 4eb3f0cd89..19d6e295b5 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 @@ -62,9 +62,9 @@ import java.util.UUID; @Slf4j public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerService implements TbCalculatedFieldConsumerService { - @Value("${queue.calculated_fields.poll_interval}") + @Value("${queue.calculated_fields.poll_interval:25}") private long pollInterval; - @Value("${queue.calculated_fields.pack_processing_timeout}") + @Value("${queue.calculated_fields.pack_processing_timeout:60000}") private long packProcessingTimeout; @Value("${queue.calculated_fields.consumer_per_partition:true}") private boolean consumerPerPartition; @@ -102,7 +102,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer this.calculatedFieldsExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(poolSize, "tb-cf-executor")); // TODO: multiple threads. this.mainConsumer = MainQueueConsumerManager., CalculatedFieldQueueConfig>builder() - .queueKey(new QueueKey(ServiceType.TB_CORE)) + .queueKey(new QueueKey(ServiceType.TB_RULE_ENGINE)) .config(CalculatedFieldQueueConfig.of(consumerPerPartition, (int) pollInterval)) .msgPackProcessor(this::processMsgs) .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer()) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index b510b813d6..3064277bf4 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -102,6 +102,7 @@ import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; +import java.util.HashSet; import java.util.List; import java.util.Objects; import java.util.Optional; @@ -570,7 +571,9 @@ public class DefaultTbClusterService implements TbClusterService { private void broadcast(ComponentLifecycleMsg msg) { ComponentLifecycleMsgProto componentLifecycleMsgProto = toProto(msg); TbQueueProducer> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer(); + TbQueueProducer> toCalculatedFieldProducer = producerProvider.getCalculatedFieldsNotificationsMsgProducer(); Set tbRuleEngineServices = partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE); + Set tbCalculatedFieldServices = new HashSet<>(tbRuleEngineServices); EntityType entityType = msg.getEntityId().getEntityType(); if (entityType.equals(EntityType.TENANT) || entityType.equals(EntityType.TENANT_PROFILE) @@ -581,7 +584,6 @@ public class DefaultTbClusterService implements TbClusterService { || (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED) || entityType.equals(EntityType.ENTITY_VIEW) || entityType.equals(EntityType.NOTIFICATION_RULE) - || entityType.equals(EntityType.CALCULATED_FIELD) ) { TbQueueProducer> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); Set tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); @@ -594,6 +596,14 @@ public class DefaultTbClusterService implements TbClusterService { // No need to push notifications twice tbRuleEngineServices.removeAll(tbCoreServices); } + if (entityType.equals(EntityType.CALCULATED_FIELD)) { + for (String serviceId : tbCalculatedFieldServices) { + TopicPartitionInfo tpi = topicService.getCalculatedFieldNotificationsTopic(serviceId); + ToCalculatedFieldNotificationMsg toCfNotificationMsg = ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build(); + toCalculatedFieldProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toCfNotificationMsg), null); + toRuleEngineNfs.incrementAndGet(); // TODO: add separate counter when we will have new ServiceType.CALCULATED_FIELDS + } + } for (String serviceId : tbRuleEngineServices) { TopicPartitionInfo tpi = topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId); ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build(); @@ -641,11 +651,11 @@ public class DefaultTbClusterService implements TbClusterService { if (deviceNameChanged) { gatewayNotificationsService.onDeviceUpdated(device, old); } - boolean deviceTypeChanged = !device.getType().equals(old.getType()); - if (deviceTypeChanged) { + boolean deviceProfileChanged = !device.getDeviceProfileId().equals(old.getDeviceProfileId()); + if (deviceProfileChanged) { handleEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); } - if (deviceNameChanged || deviceTypeChanged) { + if (deviceNameChanged || deviceProfileChanged) { pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); } } else { @@ -802,37 +812,13 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void onCalculatedFieldUpdated(CalculatedField calculatedField, CalculatedField oldCalculatedField) { var created = oldCalculatedField == null; - broadcastEntityChangeToTransport(calculatedField.getTenantId(), calculatedField.getId(), calculatedField, null); broadcastEntityStateChangeEvent(calculatedField.getTenantId(), calculatedField.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - sendCalculatedFieldEvent(calculatedField.getTenantId(), calculatedField.getId(), created, !created, false); } @Override public void onCalculatedFieldDeleted(TenantId tenantId, CalculatedField calculatedField, TbQueueCallback callback) { CalculatedFieldId calculatedFieldId = calculatedField.getId(); - broadcastEntityDeleteToTransport(tenantId, calculatedFieldId, calculatedField.getName(), callback); broadcastEntityStateChangeEvent(tenantId, calculatedFieldId, ComponentLifecycleEvent.DELETED); - sendCalculatedFieldEvent(tenantId, calculatedFieldId, false, false, true); - } - - private void sendCalculatedFieldEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, boolean added, boolean updated, boolean deleted) { - ComponentLifecycleMsgProto.Builder builder = ComponentLifecycleMsgProto.newBuilder(); - builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); - builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); - builder.setEntityType(TransportProtos.EntityTypeProto.CALCULATED_FIELD); - builder.setEntityIdMSB(calculatedFieldId.getId().getMostSignificantBits()); - builder.setEntityIdLSB(calculatedFieldId.getId().getLeastSignificantBits()); - TransportProtos.ComponentLifecycleEvent event; - if (added) { - event = TransportProtos.ComponentLifecycleEvent.CREATED; - } else if (updated) { - event = TransportProtos.ComponentLifecycleEvent.UPDATED; - } else { - event = TransportProtos.ComponentLifecycleEvent.DELETED; - } - builder.setEvent(event); - - pushNotificationToCalculatedFields(tenantId, calculatedFieldId, ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycle(builder).build(), null); } private void handleEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 0f13dbb2f9..7f889d431f 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -167,9 +167,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer public void saveAttributesInternal(AttributesSaveRequest request) { log.trace("Executing saveInternal [{}]", request); ListenableFuture> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries()); - addMainCallback(saveFuture, request.getCallback()); DonAsynchron.withCallback(saveFuture, result -> { - calculatedFieldExecutionService.pushRequestToQueue(request); + calculatedFieldExecutionService.pushRequestToQueue(request, result); }, safeCallback(request.getCallback()), tsCallBackExecutor); addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice())); addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index fdd9923292..d62aac5e6c 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1750,7 +1750,7 @@ queue: # Amount of partitions used by CF microservices partitions: "${TB_QUEUE_CF_PARTITIONS:10}" # Timeout for processing a message pack by CF microservices - pack_processing_timeout: "${TB_QUEUE_CF_PACK_PROCESSING_TIMEOUT_MS:2000}" + pack_processing_timeout: "${TB_QUEUE_CF_PACK_PROCESSING_TIMEOUT_MS:60000}" # Enable/disable a separate consumer per partition for CF queue consumer_per_partition: "${TB_QUEUE_CF_CONSUMER_PER_PARTITION:true}" # Thread pool size for processing of the incoming messages