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
new file mode 100644
index 0000000000..85ee01ac5f
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java
@@ -0,0 +1,245 @@
+/**
+ * Copyright © 2016-2024 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.queue;
+
+import com.google.common.util.concurrent.ListeningExecutorService;
+import com.google.common.util.concurrent.MoreExecutors;
+import jakarta.annotation.PostConstruct;
+import jakarta.annotation.PreDestroy;
+import lombok.Data;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.context.ApplicationEventPublisher;
+import org.springframework.stereotype.Service;
+import org.thingsboard.common.util.ThingsBoardExecutors;
+import org.thingsboard.server.actors.ActorSystemContext;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.queue.QueueConfig;
+import org.thingsboard.server.common.msg.queue.ServiceType;
+import org.thingsboard.server.common.msg.queue.TbCallback;
+import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.common.TbProtoQueueMsg;
+import org.thingsboard.server.queue.discovery.PartitionService;
+import org.thingsboard.server.queue.discovery.QueueKey;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
+import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
+import org.thingsboard.server.queue.util.TbRuleEngineComponent;
+import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
+import org.thingsboard.server.service.cf.CalculatedFieldCache;
+import org.thingsboard.server.service.cf.CalculatedFieldExecutionService;
+import org.thingsboard.server.service.profile.TbAssetProfileCache;
+import org.thingsboard.server.service.profile.TbDeviceProfileCache;
+import org.thingsboard.server.service.queue.consumer.MainQueueConsumerManager;
+import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
+import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
+
+import java.util.List;
+import java.util.UUID;
+
+@Service
+@TbRuleEngineComponent
+@Slf4j
+public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerService implements TbCalculatedFieldConsumerService {
+
+ @Value("${queue.calculated_fields.poll_interval}")
+ private long pollInterval;
+ @Value("${queue.calculated_fields.pack_processing_timeout}")
+ private long packProcessingTimeout;
+ @Value("${queue.calculated_fields.consumer_per_partition:true}")
+ private boolean consumerPerPartition;
+ @Value("${queue.calculated_fields.pool_size:8}")
+ private int poolSize;
+
+ private final TbRuleEngineQueueFactory queueFactory;
+
+ private final CalculatedFieldExecutionService calculatedFieldExecutionService;
+
+ private MainQueueConsumerManager, CalculatedFieldQueueConfig> mainConsumer;
+
+ private volatile ListeningExecutorService calculatedFieldsExecutor;
+
+ public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory,
+ ActorSystemContext actorContext,
+ TbDeviceProfileCache deviceProfileCache,
+ TbAssetProfileCache assetProfileCache,
+ TbTenantProfileCache tenantProfileCache,
+ TbApiUsageStateService apiUsageStateService,
+ PartitionService partitionService,
+ ApplicationEventPublisher eventPublisher,
+ JwtSettingsService jwtSettingsService,
+ CalculatedFieldExecutionService calculatedFieldExecutionService,
+ CalculatedFieldCache calculatedFieldCache) {
+ super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService,
+ eventPublisher, jwtSettingsService);
+ this.queueFactory = tbQueueFactory;
+ this.calculatedFieldExecutionService = calculatedFieldExecutionService;
+ }
+
+ @PostConstruct
+ public void init() {
+ super.init("tb-cf");
+ this.calculatedFieldsExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(poolSize, "tb-cf-executor")); // TODO: multiple threads.
+
+ this.mainConsumer = MainQueueConsumerManager., CalculatedFieldQueueConfig>builder()
+ .queueKey(new QueueKey(ServiceType.TB_CORE))
+ .config(CalculatedFieldQueueConfig.of(consumerPerPartition, (int) pollInterval))
+ .msgPackProcessor(this::processMsgs)
+ .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer())
+ .consumerExecutor(consumersExecutor)
+ .scheduler(scheduler)
+ .taskExecutor(mgmtExecutor)
+ .build();
+ }
+
+ @PreDestroy
+ public void destroy() {
+ super.destroy();
+ if (calculatedFieldsExecutor != null) {
+ calculatedFieldsExecutor.shutdownNow();
+ }
+ }
+
+ @Override
+ protected void startConsumers() {
+ super.startConsumers();
+ }
+
+ @Override
+ protected void onTbApplicationEvent(PartitionChangeEvent event) {
+ log.debug("Subscribing to partitions: {}", event.getCalculatedFieldsPartitions());
+ mainConsumer.update(event.getCalculatedFieldsPartitions());
+ }
+
+ private void processMsgs(List> msgs, TbQueueConsumer> consumer, CalculatedFieldQueueConfig config) throws Exception {
+
+ }
+
+ @Override
+ protected ServiceType getServiceType() {
+ return ServiceType.TB_RULE_ENGINE;
+ }
+
+ @Override
+ protected long getNotificationPollDuration() {
+ return pollInterval;
+ }
+
+ @Override
+ protected long getNotificationPackProcessingTimeout() {
+ return packProcessingTimeout;
+ }
+
+ @Override
+ protected int getMgmtThreadPoolSize() {
+ return Math.max(Runtime.getRuntime().availableProcessors(), 4);
+ }
+
+ @Override
+ protected TbQueueConsumer> createNotificationsConsumer() {
+ return queueFactory.createToCalculatedFieldNotificationsMsgConsumer();
+ }
+
+ @Override
+ protected void handleNotification(UUID id, TbProtoQueueMsg msg, TbCallback callback) {
+ ToCalculatedFieldNotificationMsg notification = msg.getValue();
+
+ callback.onSuccess();
+ }
+
+// private void processEntityProfileUpdateMsg(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg) {
+// var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
+// var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
+// var oldProfile = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityProfileType(), new UUID(profileUpdateMsg.getOldProfileIdMSB(), profileUpdateMsg.getOldProfileIdLSB()));
+// var newProfile = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityProfileType(), new UUID(profileUpdateMsg.getNewProfileIdMSB(), profileUpdateMsg.getNewProfileIdLSB()));
+// calculatedFieldCache.getEntitiesByProfile(tenantId, oldProfile).remove(entityId);
+// calculatedFieldCache.getEntitiesByProfile(tenantId, newProfile).add(entityId);
+// }
+//
+// private void processProfileEntityMsg(TransportProtos.ProfileEntityMsgProto profileEntityMsg) {
+// var tenantId = toTenantId(profileEntityMsg.getTenantIdMSB(), profileEntityMsg.getTenantIdLSB());
+// var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsg.getEntityType(), new UUID(profileEntityMsg.getEntityIdMSB(), profileEntityMsg.getEntityIdLSB()));
+// var profileId = EntityIdFactory.getByTypeAndUuid(profileEntityMsg.getEntityProfileType(), new UUID(profileEntityMsg.getProfileIdMSB(), profileEntityMsg.getProfileIdLSB()));
+// boolean added = profileEntityMsg.getAdded();
+// Set entitiesByProfile = calculatedFieldCache.getEntitiesByProfile(tenantId, profileId);
+// if (added) {
+// entitiesByProfile.add(entityId);
+// } else {
+// entitiesByProfile.remove(entityId);
+// }
+// }
+//
+// private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldMsgProto calculatedFieldMsg, TbCallback callback) {
+// var tenantId = toTenantId(calculatedFieldMsg.getTenantIdMSB(), calculatedFieldMsg.getTenantIdLSB());
+// var calculatedFieldId = new CalculatedFieldId(new UUID(calculatedFieldMsg.getCalculatedFieldIdMSB(), calculatedFieldMsg.getCalculatedFieldIdLSB()));
+// ListenableFuture> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback));
+// DonAsynchron.withCallback(future,
+// __ -> callback.onSuccess(),
+// t -> {
+// log.warn("[{}] Failed to process calculated field message for calculated field [{}]", tenantId.getId(), calculatedFieldId.getId(), t);
+// callback.onFailure(t);
+// });
+// }
+//
+// private void forwardToCalculatedFieldService(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg, TbCallback callback) {
+// var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
+// var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
+// ListenableFuture> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback));
+// DonAsynchron.withCallback(future,
+// __ -> callback.onSuccess(),
+// t -> {
+// log.warn("[{}] Failed to process entity profile updated message for entity [{}]", tenantId.getId(), entityId.getId(), t);
+// callback.onFailure(t);
+// });
+// }
+//
+// private void forwardToCalculatedFieldService(TransportProtos.ProfileEntityMsgProto profileEntityMsgProto, TbCallback callback) {
+// var tenantId = toTenantId(profileEntityMsgProto.getTenantIdMSB(), profileEntityMsgProto.getTenantIdLSB());
+// var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsgProto.getEntityType(), new UUID(profileEntityMsgProto.getEntityIdMSB(), profileEntityMsgProto.getEntityIdLSB()));
+// ListenableFuture> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onProfileEntityMsg(profileEntityMsgProto, callback));
+// DonAsynchron.withCallback(future,
+// __ -> callback.onSuccess(),
+// t -> {
+// log.warn("[{}] Failed to process profile entity message for entityId [{}]", tenantId.getId(), entityId.getId(), t);
+// callback.onFailure(t);
+// });
+// }
+
+ private void throwNotHandled(Object msg, TbCallback callback) {
+ log.warn("Message not handled: {}", msg);
+ callback.onFailure(new RuntimeException("Message not handled!"));
+ }
+
+ private TenantId toTenantId(long tenantIdMSB, long tenantIdLSB) {
+ return TenantId.fromUUID(new UUID(tenantIdMSB, tenantIdLSB));
+ }
+
+ @Override
+ protected void stopConsumers() {
+ super.stopConsumers();
+ mainConsumer.stop();
+ mainConsumer.awaitStop();
+ }
+
+ @Data(staticConstructor = "of")
+ public static class CalculatedFieldQueueConfig implements QueueConfig {
+ private final boolean consumerPerPartition;
+ private final int pollInterval;
+ }
+
+}
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 7d4d975cb4..5e07e33a9e 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
@@ -63,6 +63,8 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
+import static org.thingsboard.server.queue.discovery.HashPartitionService.CALCULATED_FIELD_QUEUE_KEY;
+
@Service
@TbRuleEngineComponent
@Slf4j
@@ -107,6 +109,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
@Override
protected void onTbApplicationEvent(PartitionChangeEvent event) {
event.getPartitionsMap().forEach((queueKey, partitions) -> {
+ if (CALCULATED_FIELD_QUEUE_KEY.equals(queueKey)) {
+ return;
+ }
if (partitionService.isManagedByCurrentService(queueKey.getTenantId())) {
var consumer = getConsumer(queueKey).orElseGet(() -> {
Queue config = queueService.findQueueByTenantIdAndName(queueKey.getTenantId(), queueKey.getQueueName());
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 dcb72b8dd0..fdf3ed6c50 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
@@ -149,13 +149,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
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;
- }
- });
+// 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()) {
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 5151bc019b..94c1bef71c 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -1738,17 +1738,21 @@ 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}"
+ calculated_fields:
+ # Topic name for Calculated Field (CF) events from Rule Engine
+ event_topic: "${TB_QUEUE_CF_EVENT_TOPIC:tb_cf_event}"
+ # Topic name for Calculated Field (CF) compacted states
+ state_topic: "${TB_QUEUE_CF_STATE_TOPIC:tb_cf_state}"
# Interval in milliseconds to poll messages by CF (Rule Engine) microservices
- poll-interval: "${TB_QUEUE_CF_POLL_INTERVAL_MS:25}"
+ 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}"
+ 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}"
+ consumer_per_partition: "${TB_QUEUE_CF_CONSUMER_PER_PARTITION:true}"
+ # Thread pool size for processing of the incoming messages
+ pool_size: "${TB_QUEUE_CF_POOL_SIZE:8}"
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/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
index 073f47d59b..1b743316bd 100644
--- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
+++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
@@ -644,15 +644,15 @@ public class ProtoUtils {
return new BasicTsKvEntry(proto.getTs(), entry, proto.hasVersion() ? proto.getVersion() : null);
}
- public static KvEntry fromTelemetryProto(TransportProtos.TelemetryProto telemetryProto) {
- if (telemetryProto.hasAttrKv()) {
- return fromProto(telemetryProto.getAttrKv().getValue());
- } else if (telemetryProto.hasTsKv()) {
- return fromProto(telemetryProto.getTsKv());
- } else {
- throw new IllegalArgumentException("Unsupported TelemetryProto type: " + telemetryProto);
- }
- }
+// public static KvEntry fromTelemetryProto(TransportProtos.TelemetryProto telemetryProto) {
+// if (telemetryProto.hasAttrKv()) {
+// return fromProto(telemetryProto.getAttrKv().getValue());
+// } else if (telemetryProto.hasTsKv()) {
+// return fromProto(telemetryProto.getTsKv());
+// } else {
+// throw new IllegalArgumentException("Unsupported TelemetryProto type: " + telemetryProto);
+// }
+// }
public static TransportProtos.AttributeKey toAttributeKeyProto(String key, AttributeScope scope) {
TransportProtos.AttributeKey.Builder builder = TransportProtos.AttributeKey.newBuilder();
@@ -673,12 +673,12 @@ public class ProtoUtils {
return builder.build();
}
- public static TransportProtos.AttributeKvProto toAttributeKvProto(AttributeKvEntry attributeKvEntry, AttributeScope scope) {
- return TransportProtos.AttributeKvProto.newBuilder()
- .setKey(ProtoUtils.toAttributeKeyProto(attributeKvEntry.getKey(), scope))
- .setValue(ProtoUtils.toAttributeValueProto(attributeKvEntry))
- .build();
- }
+// public static TransportProtos.AttributeKvProto toAttributeKvProto(AttributeKvEntry attributeKvEntry, AttributeScope scope) {
+// return TransportProtos.AttributeKvProto.newBuilder()
+// .setKey(ProtoUtils.toAttributeKeyProto(attributeKvEntry.getKey(), scope))
+// .setValue(ProtoUtils.toAttributeValueProto(attributeKvEntry))
+// .build();
+// }
public static TransportProtos.AttributeValueProto toAttributeValueProto(AttributeKvEntry attributeKvEntry) {
TransportProtos.AttributeValueProto.Builder builder = TransportProtos.AttributeValueProto.newBuilder();
diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto
index 4e719a02df..288c923aaa 100644
--- a/common/proto/src/main/proto/queue.proto
+++ b/common/proto/src/main/proto/queue.proto
@@ -819,6 +819,35 @@ message CalculatedFieldIdProto {
int64 calculatedFieldIdLSB = 2;
}
+message SingleValueProto {
+ int64 ts = 1;
+ int64 version = 2;
+ KeyValueType type = 3;
+ bool has_v = 4;
+ bool bool_v = 5;
+ int64 long_v = 6;
+ double double_v = 7;
+ string string_v = 8;
+ string json_v = 9;
+}
+
+message SingleValueArgumentProto {
+ string argName = 1;
+ SingleValueProto value = 2;
+}
+
+message RollingArgumentProto {
+ string argName = 1;
+ repeated SingleValueProto values = 2;
+}
+
+message CalculatedFieldStateProto {
+ CalculatedFieldEntityCtxIdProto id = 1;
+ // int32 version = 2;
+ repeated SingleValueArgumentProto singleValueArguments = 3;
+ repeated RollingArgumentProto rollingValueArguments = 4;
+}
+
//Used to report session state to tb-Service and persist this state in the cache on the tb-Service level.
message SubscriptionInfoProto {
int64 lastActivityTime = 1;
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 53bdc78c93..7ac938f52f 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
@@ -61,9 +61,11 @@ 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}")
+ @Value("${queue.calculated_fields.event_topic}")
+ private String cfEventTopic;
+ @Value("${queue.calculated_fields.state_topic}")
+ private String cfStateTopic;
+ @Value("${queue.calculated_fields.partitions:10}")
private Integer cfPartitions;
@Value("${queue.vc.topic:tb_version_control}")
private String vcTopic;
@@ -76,6 +78,8 @@ public class HashPartitionService implements PartitionService {
@Value("${queue.partitions.hash_function_name:murmur3_128}")
private String hashFunctionName;
+ public static final QueueKey CALCULATED_FIELD_QUEUE_KEY = new QueueKey(ServiceType.TB_RULE_ENGINE).withQueueName(CF_QUEUE_NAME);
+
private final ApplicationEventPublisher applicationEventPublisher;
private final TbServiceInfoProvider serviceInfoProvider;
private final TenantRoutingInfoService tenantRoutingInfoService;
@@ -116,9 +120,8 @@ public class HashPartitionService implements PartitionService {
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);
+ partitionSizesMap.put(CALCULATED_FIELD_QUEUE_KEY, cfPartitions);
+ partitionTopicsMap.put(CALCULATED_FIELD_QUEUE_KEY, cfEventTopic);
QueueKey vcKey = new QueueKey(ServiceType.TB_VC_EXECUTOR);
partitionSizesMap.put(vcKey, vcPartitions);
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java
index 927c311a2d..8dec36c15c 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java
@@ -35,6 +35,7 @@ public class TopicService {
private final ConcurrentMap tbCoreNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentMap tbRuleEngineNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentMap tbEdgeNotificationTopics = new ConcurrentHashMap<>();
+ private final ConcurrentMap tbCalculatedFieldNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentReferenceHashMap tbEdgeEventsNotificationTopics = new ConcurrentReferenceHashMap<>();
/**
@@ -62,6 +63,11 @@ public class TopicService {
return buildTopicPartitionInfo("tb_edge.notifications." + serviceId, null, null, false);
}
+ public TopicPartitionInfo getCalculatedFieldNotificationsTopic(String serviceId) {
+ return tbCalculatedFieldNotificationTopics.computeIfAbsent(serviceId,
+ id -> buildNotificationsTopicPartitionInfo("calculated_field", serviceId));
+ }
+
public TopicPartitionInfo getEdgeEventNotificationsTopic(TenantId tenantId, EdgeId edgeId) {
return tbEdgeEventsNotificationTopics.computeIfAbsent(edgeId, id -> buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId));
}
@@ -71,7 +77,11 @@ public class TopicService {
}
private TopicPartitionInfo buildNotificationsTopicPartitionInfo(ServiceType serviceType, String serviceId) {
- return buildTopicPartitionInfo(serviceType.name().toLowerCase() + ".notifications." + serviceId, null, null, false);
+ return buildNotificationsTopicPartitionInfo(serviceType.name().toLowerCase(), serviceId);
+ }
+
+ private TopicPartitionInfo buildNotificationsTopicPartitionInfo(String serviceType, String serviceId) {
+ return buildTopicPartitionInfo(serviceType + ".notifications." + serviceId, null, null, false);
}
public TopicPartitionInfo buildTopicPartitionInfo(String topic, TenantId tenantId, Integer partition, boolean myPartition) {
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 3bb0c56f9a..57a4941981 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
@@ -23,10 +23,13 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.queue.discovery.QueueKey;
import java.io.Serial;
+import java.util.Collections;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
+import static org.thingsboard.server.queue.discovery.HashPartitionService.CALCULATED_FIELD_QUEUE_KEY;
+
@ToString(callSuper = true)
public class PartitionChangeEvent extends TbApplicationEvent {
@@ -53,7 +56,7 @@ public class PartitionChangeEvent extends TbApplicationEvent {
}
public Set getCalculatedFieldsPartitions() {
- return getPartitionsByServiceTypeAndQueueName(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME);
+ return partitionsMap.getOrDefault(CALCULATED_FIELD_QUEUE_KEY, Collections.emptySet());
}
private Set getPartitionsByServiceTypeAndQueueName(ServiceType serviceType, String queueName) {
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 d70cad159b..c26e2d15c9 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
@@ -33,6 +33,7 @@ import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.memory.InMemoryStorage;
import org.thingsboard.server.queue.memory.InMemoryTbQueueConsumer;
import org.thingsboard.server.queue.memory.InMemoryTbQueueProducer;
+import org.thingsboard.server.queue.settings.TbQueueCalculatedFieldSettings;
import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
@@ -53,6 +54,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
private final TbQueueTransportApiSettings transportApiSettings;
private final TbQueueTransportNotificationSettings transportNotificationSettings;
private final TbQueueEdgeSettings edgeSettings;
+ private final TbQueueCalculatedFieldSettings calculatedFieldSettings;
private final InMemoryStorage storage;
public InMemoryMonolithQueueFactory(TopicService topicService, TbQueueCoreSettings coreSettings,
@@ -62,6 +64,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
TbQueueTransportApiSettings transportApiSettings,
TbQueueTransportNotificationSettings transportNotificationSettings,
TbQueueEdgeSettings edgeSettings,
+ TbQueueCalculatedFieldSettings calculatedFieldSettings,
InMemoryStorage storage) {
this.topicService = topicService;
this.coreSettings = coreSettings;
@@ -71,6 +74,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
this.transportApiSettings = transportApiSettings;
this.transportNotificationSettings = transportNotificationSettings;
this.edgeSettings = edgeSettings;
+ this.calculatedFieldSettings = calculatedFieldSettings;
this.storage = storage;
}
@@ -139,6 +143,31 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
return null;
}
+ @Override
+ public TbQueueConsumer> createToCalculatedFieldMsgConsumer() {
+ return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic()));
+ }
+
+ @Override
+ public TbQueueProducer> createToCalculatedFieldMsgProducer() {
+ return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic()));
+ }
+
+ @Override
+ public TbQueueConsumer> createToCalculatedFieldNotificationsMsgConsumer() {
+ return new InMemoryTbQueueConsumer<>(storage, topicService.getCalculatedFieldNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
+ }
+
+ @Override
+ public TbQueueConsumer> createCalculatedFieldStateConsumer() {
+ return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic()));
+ }
+
+ @Override
+ public TbQueueProducer> createCalculatedFieldStateProducer() {
+ return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic()));
+ }
+
@Override
public TbQueueConsumer> createToUsageStatsServiceMsgConsumer() {
return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
@@ -209,6 +238,11 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
return null;
}
+ @Override
+ public TbQueueProducer> createToCalculatedFieldNotificationMsgProducer() {
+ return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic()));
+ }
+
@Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}")
private void printInMemoryStats() {
storage.printStats();
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java
index c4002f4d3e..0b3df5bccf 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java
@@ -1,12 +1,12 @@
/**
* Copyright © 2016-2024 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
- *
+ *
+ * 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.
@@ -18,6 +18,7 @@ package org.thingsboard.server.queue.provider;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.gen.js.JsInvokeProtos;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
@@ -159,4 +160,6 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous
return null;
}
+ TbQueueProducer> createToCalculatedFieldNotificationMsgProducer();
+
}
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 c406aeb311..76dad05393 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
@@ -1,12 +1,12 @@
/**
* Copyright © 2016-2024 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
- *
+ *
+ * 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.
@@ -17,6 +17,9 @@ package org.thingsboard.server.queue.provider;
import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.gen.js.JsInvokeProtos;
+import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
@@ -109,11 +112,22 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory
}
/**
- * Used to consume high priority messages by TB Core Service
+ * Used to consume high priority messages by TB Rule Engine Service
*
* @return
*/
TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer();
TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate();
+
+ TbQueueConsumer> createToCalculatedFieldMsgConsumer();
+
+ TbQueueProducer> createToCalculatedFieldMsgProducer();
+
+ TbQueueConsumer> createToCalculatedFieldNotificationsMsgConsumer();
+
+ TbQueueConsumer> createCalculatedFieldStateConsumer();
+
+ TbQueueProducer> createCalculatedFieldStateProducer();
+
}
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
new file mode 100644
index 0000000000..22bbd7e0f3
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCalculatedFieldSettings.java
@@ -0,0 +1,35 @@
+/**
+ * Copyright © 2016-2024 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.settings;
+
+import lombok.Data;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.context.annotation.Lazy;
+import org.springframework.stereotype.Component;
+
+@Lazy
+@Data
+@Component
+public class TbQueueCalculatedFieldSettings {
+
+ @Value("${queue.calculated_fields.event_topic}")
+ private String eventTopic;
+
+ @Value("${queue.calculated_fields.state_topic}")
+ private String stateTopic;
+
+
+}