diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java index 7c7e72b511..f669df5305 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java @@ -15,10 +15,10 @@ */ package org.thingsboard.server.service.housekeeper; -import com.datastax.oss.driver.api.core.uuid.Uuids; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -40,6 +40,8 @@ import org.thingsboard.server.service.housekeeper.processor.HousekeeperTaskProce import javax.annotation.PreDestroy; import java.util.List; import java.util.Map; +import java.util.Optional; +import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.stream.Collectors; @@ -49,7 +51,7 @@ import java.util.stream.Collectors; @Slf4j public class DefaultHousekeeperService implements HousekeeperService { - private final Map taskProcessors; + private final Map> taskProcessors; private final TbQueueConsumer> consumer; private final TbQueueProducer> producer; @@ -65,7 +67,8 @@ public class DefaultHousekeeperService implements HousekeeperService { public DefaultHousekeeperService(HousekeeperReprocessingService reprocessingService, TbCoreQueueFactory queueFactory, TbQueueProducerProvider producerProvider, - DataDecodingEncodingService dataDecodingEncodingService, List taskProcessors) { + DataDecodingEncodingService dataDecodingEncodingService, + @Lazy List> taskProcessors) { this.consumer = queueFactory.createHousekeeperMsgConsumer(); this.producer = producerProvider.getHousekeeperMsgProducer(); this.reprocessingService = reprocessingService; @@ -86,7 +89,7 @@ public class DefaultHousekeeperService implements HousekeeperService { for (TbProtoQueueMsg msg : msgs) { try { - processTask(msg); + processTask(msg.getValue()); } catch (Exception e) { log.error("Message processing failed", e); } @@ -107,19 +110,18 @@ public class DefaultHousekeeperService implements HousekeeperService { log.info("Started Housekeeper service"); } - protected void processTask(TbProtoQueueMsg msg) { - HousekeeperTask task = dataDecodingEncodingService.decode(msg.getValue().getTask().getValue().toByteArray()).get(); - HousekeeperTaskProcessor taskProcessor = taskProcessors.get(task.getTaskType()); - if (taskProcessor == null) { - log.error("Unsupported task type {}: {}", task.getTaskType(), task); - return; - } + @SuppressWarnings("unchecked") + protected void processTask(ToHousekeeperServiceMsg msg) { + HousekeeperTask task = dataDecodingEncodingService.decode(msg.getTask().getValue().toByteArray()).get(); + HousekeeperTaskProcessor taskProcessor = getTaskProcessor(task.getTaskType()); - log.info("[{}] Processing task: {}", task.getTenantId(), task); + log.info("[{}][{}][{}] Processing task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType()); try { - taskProcessor.process(task); + taskProcessor.process((T) task); } catch (Exception e) { - log.error("[{}] Task processing failed: {}", task.getTenantId(), task, e); + log.error("[{}][{}][{}] {} task processing failed, submitting for reprocessing (attempt {}): {}", + task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), + task.getTaskType(), msg.getTask().getAttempt(), task, e); reprocessingService.submitForReprocessing(msg); } } @@ -127,11 +129,20 @@ public class DefaultHousekeeperService implements HousekeeperService { @Override public void submitTask(HousekeeperTask task) { TopicPartitionInfo tpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); - producer.send(tpi, new TbProtoQueueMsg<>(Uuids.timeBased(), ToHousekeeperServiceMsg.newBuilder() + producer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), ToHousekeeperServiceMsg.newBuilder() .setTask(HousekeeperTaskProto.newBuilder() .setValue(ByteString.copyFrom(dataDecodingEncodingService.encode(task))) + .setTs(task.getTs()) + .setAttempt(0) .build()) .build()), null); + log.trace("[{}][{}][{}] Submitted task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType()); + } + + @SuppressWarnings("unchecked") + private HousekeeperTaskProcessor getTaskProcessor(HousekeeperTaskType taskType) { + return Optional.ofNullable((HousekeeperTaskProcessor) taskProcessors.get(taskType)) + .orElseThrow(() -> new IllegalArgumentException("Unsupported task type " + taskType)); } @PreDestroy @@ -141,4 +152,5 @@ public class DefaultHousekeeperService implements HousekeeperService { consumer.unsubscribe(); consumerExecutor.shutdownNow(); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java index 234cd07c36..4fca14e44f 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java @@ -15,16 +15,15 @@ */ package org.thingsboard.server.service.housekeeper; -import com.datastax.oss.driver.api.core.uuid.Uuids; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.gen.transport.TransportProtos.HousekeeperTaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.queue.TbQueueConsumer; -import org.thingsboard.server.queue.TbQueueMsgHeaders; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; @@ -33,13 +32,10 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import javax.annotation.PreDestroy; import java.util.List; -import java.util.Optional; +import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToLong; -import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.longToBytes; - @TbCoreComponent @Service @Slf4j @@ -77,15 +73,21 @@ public class HousekeeperReprocessingService { return; } - for (TbProtoQueueMsg msg : msgs) { - long msgTs = Uuids.unixTimestamp(msg.getKey()); - if (msgTs >= startTs) { - stop(); - return; // fixme: we should commit already reprocessed messages + if (msgs.stream().anyMatch(msg -> { + boolean newMsg = msg.getValue().getTask().getTs() >= startTs; + if (newMsg) { + log.info("Stopping reprocessing due to msg is new {}", msg); } - + return newMsg; + })) { + stop(); // fixme: we should commit already reprocessed messages; maybe submit for reprocessing again and commit? + // msg batch size should be 1. otherwise some tasks won't be reprocessed + return; + } + for (TbProtoQueueMsg msg : msgs) { try { - reprocessTask(msg); + housekeeperService.processTask(msg.getValue());// fixme: or should we submit to queue? + Thread.sleep(1000); } catch (Exception e) { log.error("Message processing failed", e); } @@ -106,21 +108,19 @@ public class HousekeeperReprocessingService { log.info("Started Housekeeper tasks reprocessing"); } - private void reprocessTask(TbProtoQueueMsg msg) { - housekeeperService.processTask(msg);// fixme: or should we submit to queue? - } - - public void submitForReprocessing(TbProtoQueueMsg msg) { - TbQueueMsgHeaders msgHeaders = msg.getHeaders(); - long reprocessingAttempts = Optional.ofNullable(msgHeaders.get("reprocessingAttempts")) - .map(header -> bytesToLong(header)) - .orElse(0L); - reprocessingAttempts++; - msgHeaders.put("reprocessingAttempts", longToBytes(reprocessingAttempts)); + public void submitForReprocessing(ToHousekeeperServiceMsg msg) { + HousekeeperTaskProto task = msg.getTask(); + int attempt = task.getAttempt() + 1; + msg = msg.toBuilder() + .setTask(task.toBuilder() + .setAttempt(attempt) + .setTs(System.currentTimeMillis()) + .build()) + .build(); var producer = producerProvider.getHousekeeperDelayedMsgProducer(); TopicPartitionInfo tpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); - producer.send(tpi, new TbProtoQueueMsg<>(Uuids.timeBased(), msg.getValue(), msgHeaders), null); + producer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); } @PreDestroy diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java index 54bf7989af..247466af65 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java @@ -16,21 +16,27 @@ package org.thingsboard.server.service.housekeeper.processor; import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.dao.housekeeper.data.HousekeeperTaskType; import org.thingsboard.server.dao.housekeeper.data.AlarmsUnassignHousekeeperTask; import org.thingsboard.server.service.entitiy.alarm.TbAlarmService; +import java.util.List; + @Component @RequiredArgsConstructor +@Slf4j public class AlarmsUnassignTaskProcessor implements HousekeeperTaskProcessor { private final TbAlarmService alarmService; @Override public void process(AlarmsUnassignHousekeeperTask task) throws Exception { - alarmService.unassignDeletedUserAlarms(task.getTenantId(), (UserId) task.getEntityId(), task.getUserTitle(), task.getTs()); + List alarms = alarmService.unassignDeletedUserAlarms(task.getTenantId(), (UserId) task.getEntityId(), task.getUserTitle(), task.getTs()); + log.trace("[{}][{}] Unassigned {} alarms", task.getTenantId(), task.getEntityId(), alarms.size()); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java index 2a491148ea..e925c9f20c 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java @@ -35,7 +35,7 @@ public class TelemetryDeletionTaskProcessor implements HousekeeperTaskProcessor< @Override public HousekeeperTaskType getTaskType() { - return HousekeeperTaskType.DELETE_ATTRIBUTES; + return HousekeeperTaskType.DELETE_TELEMETRY; } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 17b4c85cfb..6a96c0a3db 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1565,6 +1565,7 @@ queue: print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:60000}" housekeeper: topic: "tb_housekeeper" + reprocessing-topic: "tb_housekeeper.reprocessing" poll-interval-ms: "1000" vc: # Default topic name for Kafka, RabbitMQ, etc. diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 0dbd96a53d..c4e8e64e42 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -1398,4 +1398,6 @@ message ToHousekeeperServiceMsg { message HousekeeperTaskProto { bytes value = 1; + int64 ts = 2; + int32 attempt = 3; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueTemplate.java index 9340cafa67..f11eaaef48 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueTemplate.java @@ -46,13 +46,13 @@ public class AbstractTbQueueTemplate { return new String(data, StandardCharsets.UTF_8); } - public static byte[] longToBytes(long x) { + protected static byte[] longToBytes(long x) { ByteBuffer longBuffer = ByteBuffer.allocate(Long.BYTES); longBuffer.putLong(0, x); return longBuffer.array(); } - public static long bytesToLong(byte[] bytes) { + protected static long bytesToLong(byte[] bytes) { return ByteBuffer.wrap(bytes).getLong(); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCoreSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCoreSettings.java index 4edb688235..0b66c2fa75 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCoreSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCoreSettings.java @@ -37,7 +37,7 @@ public class TbQueueCoreSettings { @Value("${queue.core.housekeeper.topic:tb_housekeeper}") private String housekeeperTopic; - @Value("${queue.core.housekeeper.topic:tb_housekeeper.delayed}") + @Value("${queue.core.housekeeper.reprocessing-topic:tb_housekeeper.reprocessing}") private String housekeeperDelayedTopic; @Value("${queue.core.partitions}") diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java index 7bfe76e753..6b88f19699 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java @@ -20,11 +20,8 @@ import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.annotation.Lazy; -import org.springframework.transaction.event.TransactionalEventListener; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.StringUtils; -import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -33,10 +30,8 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; -import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.exception.DataValidationException; -import org.thingsboard.server.dao.housekeeper.HousekeeperService; -import org.thingsboard.server.dao.housekeeper.data.HousekeeperTask; +import org.thingsboard.server.dao.housekeeper.CleanUpService; import org.thingsboard.server.dao.relation.RelationService; import java.util.Collections; @@ -70,28 +65,8 @@ public abstract class AbstractEntityService { protected EdgeService edgeService; @Autowired - protected HousekeeperService housekeeperService; - - @TransactionalEventListener(fallbackExecution = true) // todo: consider moving this to HousekeeperService - public void onEntityDeleted(DeleteEntityEvent event) { - TenantId tenantId = event.getTenantId(); - EntityId entityId = event.getEntityId(); - log.trace("[{}] DeleteEntityEvent handler: {}", tenantId, event); - - cleanUpRelatedData(tenantId, entityId); - if (entityId.getEntityType() == EntityType.USER) { - housekeeperService.submitTask(HousekeeperTask.unassignAlarms((User) event.getEntity())); - } - } - - protected void cleanUpRelatedData(TenantId tenantId, EntityId entityId) { - // todo: skipped entities list - relationService.deleteEntityRelations(tenantId, entityId); - housekeeperService.submitTask(HousekeeperTask.deleteAttributes(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteTelemetry(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteEvents(tenantId, entityId)); - housekeeperService.submitTask(HousekeeperTask.deleteEntityAlarms(tenantId, entityId)); - } + @Lazy + protected CleanUpService cleanUpService; protected void createRelation(TenantId tenantId, EntityRelation relation) { log.debug("Creating relation: {}", relation); diff --git a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java new file mode 100644 index 0000000000..0086e17926 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java @@ -0,0 +1,60 @@ +/** + * 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.dao.housekeeper; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.springframework.transaction.event.TransactionalEventListener; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; +import org.thingsboard.server.dao.housekeeper.data.HousekeeperTask; +import org.thingsboard.server.dao.relation.RelationService; + +@Component +@RequiredArgsConstructor +@Slf4j +public class CleanUpService { + + private final HousekeeperService housekeeperService; + private final RelationService relationService; + + @TransactionalEventListener(fallbackExecution = true) // todo: consider moving this to HousekeeperService + public void handleEntityDeletionEvent(DeleteEntityEvent event) { + TenantId tenantId = event.getTenantId(); + EntityId entityId = event.getEntityId(); + log.trace("[{}] DeleteEntityEvent handler: {}", tenantId, event); + + log.info("[{}][{}][{}] Handling DeleteEntityEvent", tenantId, entityId.getEntityType(), entityId.getId()); + cleanUpRelatedData(tenantId, entityId); + if (entityId.getEntityType() == EntityType.USER) { + housekeeperService.submitTask(HousekeeperTask.unassignAlarms((User) event.getEntity())); + } + } + + public void cleanUpRelatedData(TenantId tenantId, EntityId entityId) { + // todo: skipped entities list + relationService.deleteEntityRelations(tenantId, entityId); + housekeeperService.submitTask(HousekeeperTask.deleteAttributes(tenantId, entityId)); + housekeeperService.submitTask(HousekeeperTask.deleteTelemetry(tenantId, entityId)); + housekeeperService.submitTask(HousekeeperTask.deleteEvents(tenantId, entityId)); + housekeeperService.submitTask(HousekeeperTask.deleteEntityAlarms(tenantId, entityId)); + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/AlarmsUnassignHousekeeperTask.java b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/AlarmsUnassignHousekeeperTask.java index 876a8f736b..aa1858f59f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/AlarmsUnassignHousekeeperTask.java +++ b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/AlarmsUnassignHousekeeperTask.java @@ -22,12 +22,10 @@ import org.thingsboard.server.common.data.User; public class AlarmsUnassignHousekeeperTask extends HousekeeperTask { private final String userTitle; - private final long ts; protected AlarmsUnassignHousekeeperTask(User user) { super(user.getTenantId(), user.getId(), HousekeeperTaskType.UNASSIGN_ALARMS); this.userTitle = user.getTitle(); - this.ts = System.currentTimeMillis(); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java index 726e07a790..9fcbbfe43d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java +++ b/dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java @@ -15,27 +15,26 @@ */ package org.thingsboard.server.dao.housekeeper.data; -import lombok.Getter; +import lombok.Data; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import java.io.Serializable; -/* - * on start, read the retry queue and put the messages back to main queue (save offset) - * */ -@Getter +@Data public class HousekeeperTask implements Serializable { private final TenantId tenantId; private final EntityId entityId; private final HousekeeperTaskType taskType; + private final long ts; protected HousekeeperTask(TenantId tenantId, EntityId entityId, HousekeeperTaskType taskType) { this.tenantId = tenantId; this.entityId = entityId; this.taskType = taskType; + this.ts = System.currentTimeMillis(); } public static HousekeeperTask deleteAttributes(TenantId tenantId, EntityId entityId) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 88fa6f1351..ebbec8969f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -177,7 +177,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC List updatedRuleNodes = new ArrayList<>(); List existingRuleNodes = getRuleChainNodes(tenantId, ruleChainMetaData.getRuleChainId()); for (RuleNode existingNode : existingRuleNodes) { - cleanUpRelatedData(tenantId, existingNode.getId()); // fixme: for sure? + cleanUpService.cleanUpRelatedData(tenantId, existingNode.getId()); // fixme: for sure? Integer index = ruleNodeIndexMap.get(existingNode.getId()); RuleNode newRuleNode = null; if (index != null) { @@ -771,7 +771,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC private void deleteRuleNodes(TenantId tenantId, List ruleNodes) { List ruleNodeIds = ruleNodes.stream().map(RuleNode::getId).collect(Collectors.toList()); for (var node : ruleNodes) { - cleanUpRelatedData(tenantId, node.getId()); + cleanUpService.cleanUpRelatedData(tenantId, node.getId()); } ruleNodeDao.deleteByIdIn(ruleNodeIds); } @@ -820,7 +820,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC private void deleteRuleNode(TenantId tenantId, EntityId entityId) { ruleNodeDao.removeById(tenantId, entityId.getId()); - cleanUpRelatedData(tenantId, entityId); + cleanUpService.cleanUpRelatedData(tenantId, entityId); } private final PaginatedRemover tenantRuleChainsRemover = diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index cb51a7d8ad..d8cf6aab32 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java @@ -20,9 +20,11 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.tuple.Pair; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; @@ -31,6 +33,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.attributes.AttributesDao; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey; import org.thingsboard.server.dao.model.sql.AttributeKvEntity; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; @@ -132,10 +135,10 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl @Override public List findAll(TenantId tenantId, EntityId entityId, String attributeType) { return DaoUtil.convertDataList(Lists.newArrayList( - attributeKvRepository.findAllByEntityTypeAndEntityIdAndAttributeType( - entityId.getEntityType(), - entityId.getId(), - attributeType))); + attributeKvRepository.findAllByEntityTypeAndEntityIdAndAttributeType( + entityId.getEntityType(), + entityId.getId(), + attributeType))); } @Override @@ -188,6 +191,16 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl return futuresList; } + @Transactional + @Override + public List> removeAllByEntityId(TenantId tenantId, EntityId entityId) { + return jdbcTemplate.queryForList("DELETE FROM attribute_kv WHERE entity_type = ? and entity_id = ? " + + "RETURNING attribute_type, attribute_key", entityId.getEntityType().name(), entityId.getId()).stream() + .map(deleted -> Pair.of((String) deleted.get(ModelConstants.ATTRIBUTE_TYPE_COLUMN), + (String) deleted.get(ModelConstants.ATTRIBUTE_KEY_COLUMN))) + .collect(Collectors.toList()); + } + private AttributeKvCompositeKey getAttributeKvCompositeKey(EntityId entityId, String attributeType, String attributeKey) { return new AttributeKvCompositeKey( entityId.getEntityType(),