From 6b9d374a5f2957d14a3722c2e6a6459da211db11 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 22 Jan 2025 12:23:42 +0200 Subject: [PATCH] Tmp commit for merge --- .../cf/CalculatedFieldExecutionService.java | 11 +++ .../TbCalculatedFieldConsumerService.java | 8 ++ .../DefaultTelemetrySubscriptionService.java | 18 +++-- .../src/main/resources/thingsboard.yml | 12 ++- .../server/common/data/DataConstants.java | 2 + .../server/common/msg/queue/ServiceType.java | 3 +- common/proto/src/main/proto/queue.proto | 73 +++++++------------ .../queue/discovery/HashPartitionService.java | 12 ++- .../discovery/event/PartitionChangeEvent.java | 4 + 9 files changed, 86 insertions(+), 57 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/queue/TbCalculatedFieldConsumerService.java 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 8ba1f6dfed..e18c8b4119 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 @@ -21,6 +21,17 @@ import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdat public interface CalculatedFieldExecutionService { + /** + * Push incoming telemetry to the CF processing queue for async processing. + * @param request - telemetry request; + * @param callback - callback to be executed when the message is ack by the queue. + */ + void pushRequestToQueue(CalculatedFieldTelemetryUpdateRequest request, TbCallback callback); + + void pushEntityUpdateMsg(TransportProtos.CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback); + + /* ===================================================== */ + void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbCalculatedFieldConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/TbCalculatedFieldConsumerService.java new file mode 100644 index 0000000000..387bdd7143 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbCalculatedFieldConsumerService.java @@ -0,0 +1,8 @@ +package org.thingsboard.server.service.queue; + +import org.springframework.context.ApplicationListener; +import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; + +public interface TbCalculatedFieldConsumerService extends ApplicationListener { + +} 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 8773564e5d..dcb72b8dd0 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.telemetry; +import com.google.common.util.concurrent.AsyncFunction; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -128,8 +129,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer KvUtils.validate(request.getEntries(), valueNoXssValidation); ListenableFuture future = saveTimeseriesInternal(request); if (!request.isOnlyLatest()) { - FutureCallback callback = getApiUsageCallback(tenantId, request.getCustomerId(), sysTenant, request.getCallback()); - Futures.addCallback(future, callback, tsCallBackExecutor); + Futures.addCallback(future, getApiUsageCallback(tenantId, request.getCustomerId(), sysTenant), tsCallBackExecutor); } } else { request.getCallback().onFailure(new RuntimeException("DB storage writes are disabled due to API limits!")); @@ -148,7 +148,14 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer } else { saveFuture = tsService.saveWithoutLatest(tenantId, entityId, request.getEntries(), request.getTtl()); } - + // We need to guarantee, that the message is successfully pushed to the calculated fields service before we execute any callbacks. + saveFuture = Futures.transformAsync(saveFuture, new AsyncFunction() { + @Override + public ListenableFuture apply(Integer input) throws Exception { + calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(request)); + return input; + } + }); addMainCallback(saveFuture, request.getCallback()); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries())); if (request.isSaveLatest() && !request.isOnlyLatest()) { @@ -326,19 +333,18 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer } } - private FutureCallback getApiUsageCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant, FutureCallback callback) { + private FutureCallback getApiUsageCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant) { return new FutureCallback<>() { @Override public void onSuccess(Integer result) { if (!sysTenant && result != null && result > 0) { apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, result); } - callback.onSuccess(null); } @Override public void onFailure(Throwable t) { - callback.onFailure(t); + } }; } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 3014e1448b..5151bc019b 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1692,7 +1692,6 @@ queue: enabled: "${TB_HOUSEKEEPER_STATS_ENABLED:true}" # Statistics printing interval for Housekeeper print-interval-ms: "${TB_HOUSEKEEPER_STATS_PRINT_INTERVAL_MS:60000}" - vc: # Default topic name topic: "${TB_QUEUE_VC_TOPIC:tb_version_control}" @@ -1739,6 +1738,17 @@ queue: topic-deletion-delay: "${TB_QUEUE_RULE_ENGINE_TOPIC_DELETION_DELAY_SEC:15}" # Size of the thread pool that handles such operations as partition changes, config updates, queue deletion management-thread-pool-size: "${TB_QUEUE_RULE_ENGINE_MGMT_THREAD_POOL_SIZE:12}" + calculated-fields: + # Topic name for Calculated Field (CF) tasks + topic: "${TB_QUEUE_CF_TOPIC:tb_calculated_fields}" + # Interval in milliseconds to poll messages by CF (Rule Engine) microservices + poll-interval: "${TB_QUEUE_CF_POLL_INTERVAL_MS:25}" + # 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}" + # Enable/disable a separate consumer per partition for CF queue + consumer-per-partition: "${TB_QUEUE_CF_CONSUMER_PER_PARTITION:true}" 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/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 56a5e135f8..77a7c4a781 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -145,4 +145,6 @@ public class DataConstants { public static final String EDGE_QUEUE_NAME = "Edge"; public static final String EDGE_EVENT_QUEUE_NAME = "EdgeEvent"; + public static final String CF_QUEUE_NAME = "CalculatedFields"; + } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceType.java index f3a0e47d09..022c46bcc6 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/ServiceType.java @@ -26,7 +26,8 @@ public enum ServiceType { TB_RULE_ENGINE("TB Rule Engine"), TB_TRANSPORT("TB Transport"), JS_EXECUTOR("JS Executor"), - TB_VC_EXECUTOR("TB VC Executor"); + TB_VC_EXECUTOR("TB VC Executor"), + TB_CF_ENGINE("TB Calculated Fields Engine"); private final String label; diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 1036d5ba67..4e719a02df 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -183,18 +183,6 @@ message TsKvListProto { repeated KeyValueProto kv = 2; } -message AttributeKvProto { - AttributeKey key = 1; - AttributeValueProto value = 2; -} - -message TelemetryProto { - oneof proto { - AttributeKvProto attrKv = 1; - TsKvProto tsKv = 2; - } -} - message DeviceInfoProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; @@ -785,17 +773,7 @@ message DeviceInactivityProto { int64 lastInactivityTime = 5; } -message CalculatedFieldMsgProto { - int64 tenantIdMSB = 1; - int64 tenantIdLSB = 2; - int64 calculatedFieldIdMSB = 3; - int64 calculatedFieldIdLSB = 4; - bool added = 5; - bool updated = 6; - bool deleted = 7; -} - -message EntityProfileUpdateMsgProto { +message CalculatedFieldEntityUpdateMsgProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; string entityType = 3; @@ -806,31 +784,26 @@ message EntityProfileUpdateMsgProto { int64 oldProfileIdLSB = 8; int64 newProfileIdMSB = 9; int64 newProfileIdLSB = 10; + bool added = 11; + bool updated = 12; + bool deleted = 13; } -message ProfileEntityMsgProto { +message CalculatedFieldTelemetryMsgProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; string entityType = 3; int64 entityIdMSB = 4; int64 entityIdLSB = 5; - string entityProfileType = 6; - int64 profileIdMSB = 7; - int64 profileIdLSB = 8; - bool added = 9; - bool deleted = 10; + repeated CalculatedFieldIdProto previousCalculatedFields = 7; + repeated TsKvProto tsData = 9; + AttributeScopeProto scope = 10; + repeated AttributeValueProto attrData = 11; } -message TelemetryUpdateMsgProto { - int64 tenantIdMSB = 1; - int64 tenantIdLSB = 2; - string entityType = 3; - int64 entityIdMSB = 4; - int64 entityIdLSB = 5; - repeated CalculatedFieldEntityCtxIdProto links = 6; - repeated CalculatedFieldIdProto previousCalculatedFields = 7; - string scope = 8; - repeated TelemetryProto updatedTelemetry = 9; +message CalculatedFieldLinkedTelemetryMsgProto { + CalculatedFieldTelemetryMsgProto msg = 1; + repeated CalculatedFieldEntityCtxIdProto links = 2; } message CalculatedFieldEntityCtxIdProto { @@ -1589,9 +1562,8 @@ message ToCoreMsg { DeviceConnectProto deviceConnectMsg = 50; DeviceDisconnectProto deviceDisconnectMsg = 51; DeviceInactivityProto deviceInactivityMsg = 52; - CalculatedFieldMsgProto calculatedFieldMsg = 53; - EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; - ProfileEntityMsgProto profileEntityMsg = 55; +// CalculatedFieldMsgProto calculatedFieldMsg = 53; +// EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; } /* High priority messages with low latency are handled by ThingsBoard Core Service separately */ @@ -1611,8 +1583,8 @@ message ToCoreNotificationMsg { FromEdgeSyncResponseMsgProto fromEdgeSyncResponse = 12 [deprecated = true]; ResourceCacheInvalidateMsg resourceCacheInvalidateMsg = 13; RestApiCallResponseMsgProto restApiCallResponseMsg = 50; - EntityProfileUpdateMsgProto entityProfileUpdateMsg = 51; - ProfileEntityMsgProto profileEntityMsg = 52; +// EntityProfileUpdateMsgProto entityProfileUpdateMsg = 51; +// ProfileEntityMsgProto profileEntityMsg = 52; } /* Messages to Edge queue that are handled by ThingsBoard Core Service */ @@ -1632,6 +1604,16 @@ message ToEdgeEventNotificationMsg { EdgeEventMsgProto edgeEventMsg = 1; } +message ToCalculatedFieldMsg { + CalculatedFieldTelemetryMsgProto telemetryMsg = 1; + CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsg = 2; +} + +message ToCalculatedFieldNotificationMsg { + ComponentLifecycleMsgProto componentLifecycle = 1; + CalculatedFieldEntityUpdateMsgProto entityUpdateMsg = 2; +} + /* Messages that are handled by ThingsBoard RuleEngine Service */ message ToRuleEngineMsg { int64 tenantIdMSB = 1; @@ -1639,9 +1621,6 @@ message ToRuleEngineMsg { bytes tbMsg = 3; repeated string relationTypes = 4; string failureMessage = 5; - TelemetryUpdateMsgProto cfTelemetryUpdateMsg = 6; - EntityProfileUpdateMsgProto entityProfileUpdateMsg = 7; - ProfileEntityMsgProto profileEntityMsg = 8; } message ToRuleEngineNotificationMsg { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 37e519e3f2..53bdc78c93 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -51,8 +51,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.stream.Collectors; -import static org.thingsboard.server.common.data.DataConstants.EDGE_QUEUE_NAME; -import static org.thingsboard.server.common.data.DataConstants.MAIN_QUEUE_NAME; +import static org.thingsboard.server.common.data.DataConstants.*; @Service @Slf4j @@ -62,6 +61,10 @@ public class HashPartitionService implements PartitionService { private String coreTopic; @Value("${queue.core.partitions:10}") private Integer corePartitions; + @Value("${queue.calculated-fields.topic}") + private String cfTopic; + @Value("${queue.calculated-fields.partitions:10}") + private Integer cfPartitions; @Value("${queue.vc.topic:tb_version_control}") private String vcTopic; @Value("${queue.vc.partitions:10}") @@ -108,10 +111,15 @@ public class HashPartitionService implements PartitionService { @PostConstruct public void init() { this.hashFunction = forName(hashFunctionName); + QueueKey coreKey = new QueueKey(ServiceType.TB_CORE); partitionSizesMap.put(coreKey, corePartitions); partitionTopicsMap.put(coreKey, coreTopic); + QueueKey cfKey = new QueueKey(ServiceType.TB_RULE_ENGINE).withQueueName(CF_QUEUE_NAME); + partitionSizesMap.put(cfKey, cfPartitions); + partitionTopicsMap.put(cfKey, cfTopic); + QueueKey vcKey = new QueueKey(ServiceType.TB_VC_EXECUTOR); partitionSizesMap.put(vcKey, vcPartitions); partitionTopicsMap.put(vcKey, vcTopic); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java index 88ceb4aa08..3bb0c56f9a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java @@ -52,6 +52,10 @@ public class PartitionChangeEvent extends TbApplicationEvent { return getPartitionsByServiceTypeAndQueueName(ServiceType.TB_CORE, DataConstants.EDGE_QUEUE_NAME); } + public Set getCalculatedFieldsPartitions() { + return getPartitionsByServiceTypeAndQueueName(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); + } + private Set getPartitionsByServiceTypeAndQueueName(ServiceType serviceType, String queueName) { return partitionsMap.entrySet() .stream()