Browse Source

Housekeeper tasks reprocessing

pull/10201/head
ViacheslavKlimov 3 years ago
parent
commit
7e07771acc
  1. 42
      application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java
  2. 50
      application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java
  3. 8
      application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java
  5. 1
      application/src/main/resources/thingsboard.yml
  6. 2
      common/proto/src/main/proto/queue.proto
  7. 4
      common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueTemplate.java
  8. 2
      common/queue/src/main/java/org/thingsboard/server/queue/settings/TbQueueCoreSettings.java
  9. 31
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  10. 60
      dao/src/main/java/org/thingsboard/server/dao/housekeeper/CleanUpService.java
  11. 2
      dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/AlarmsUnassignHousekeeperTask.java
  12. 9
      dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java
  13. 6
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  14. 21
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java

42
application/src/main/java/org/thingsboard/server/service/housekeeper/DefaultHousekeeperService.java

@ -15,10 +15,10 @@
*/ */
package org.thingsboard.server.service.housekeeper; package org.thingsboard.server.service.housekeeper;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.protobuf.ByteString; import com.google.protobuf.ByteString;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; 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 javax.annotation.PreDestroy;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -49,7 +51,7 @@ import java.util.stream.Collectors;
@Slf4j @Slf4j
public class DefaultHousekeeperService implements HousekeeperService { public class DefaultHousekeeperService implements HousekeeperService {
private final Map<HousekeeperTaskType, HousekeeperTaskProcessor> taskProcessors; private final Map<HousekeeperTaskType, HousekeeperTaskProcessor<?>> taskProcessors;
private final TbQueueConsumer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> consumer; private final TbQueueConsumer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> consumer;
private final TbQueueProducer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> producer; private final TbQueueProducer<TbProtoQueueMsg<ToHousekeeperServiceMsg>> producer;
@ -65,7 +67,8 @@ public class DefaultHousekeeperService implements HousekeeperService {
public DefaultHousekeeperService(HousekeeperReprocessingService reprocessingService, public DefaultHousekeeperService(HousekeeperReprocessingService reprocessingService,
TbCoreQueueFactory queueFactory, TbCoreQueueFactory queueFactory,
TbQueueProducerProvider producerProvider, TbQueueProducerProvider producerProvider,
DataDecodingEncodingService dataDecodingEncodingService, List<HousekeeperTaskProcessor> taskProcessors) { DataDecodingEncodingService dataDecodingEncodingService,
@Lazy List<HousekeeperTaskProcessor<?>> taskProcessors) {
this.consumer = queueFactory.createHousekeeperMsgConsumer(); this.consumer = queueFactory.createHousekeeperMsgConsumer();
this.producer = producerProvider.getHousekeeperMsgProducer(); this.producer = producerProvider.getHousekeeperMsgProducer();
this.reprocessingService = reprocessingService; this.reprocessingService = reprocessingService;
@ -86,7 +89,7 @@ public class DefaultHousekeeperService implements HousekeeperService {
for (TbProtoQueueMsg<ToHousekeeperServiceMsg> msg : msgs) { for (TbProtoQueueMsg<ToHousekeeperServiceMsg> msg : msgs) {
try { try {
processTask(msg); processTask(msg.getValue());
} catch (Exception e) { } catch (Exception e) {
log.error("Message processing failed", e); log.error("Message processing failed", e);
} }
@ -107,19 +110,18 @@ public class DefaultHousekeeperService implements HousekeeperService {
log.info("Started Housekeeper service"); log.info("Started Housekeeper service");
} }
protected void processTask(TbProtoQueueMsg<ToHousekeeperServiceMsg> msg) { @SuppressWarnings("unchecked")
HousekeeperTask task = dataDecodingEncodingService.<HousekeeperTask>decode(msg.getValue().getTask().getValue().toByteArray()).get(); protected <T extends HousekeeperTask> void processTask(ToHousekeeperServiceMsg msg) {
HousekeeperTaskProcessor taskProcessor = taskProcessors.get(task.getTaskType()); HousekeeperTask task = dataDecodingEncodingService.<HousekeeperTask>decode(msg.getTask().getValue().toByteArray()).get();
if (taskProcessor == null) { HousekeeperTaskProcessor<T> taskProcessor = getTaskProcessor(task.getTaskType());
log.error("Unsupported task type {}: {}", task.getTaskType(), task);
return;
}
log.info("[{}] Processing task: {}", task.getTenantId(), task); log.info("[{}][{}][{}] Processing task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType());
try { try {
taskProcessor.process(task); taskProcessor.process((T) task);
} catch (Exception e) { } 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); reprocessingService.submitForReprocessing(msg);
} }
} }
@ -127,11 +129,20 @@ public class DefaultHousekeeperService implements HousekeeperService {
@Override @Override
public void submitTask(HousekeeperTask task) { public void submitTask(HousekeeperTask task) {
TopicPartitionInfo tpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); 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() .setTask(HousekeeperTaskProto.newBuilder()
.setValue(ByteString.copyFrom(dataDecodingEncodingService.encode(task))) .setValue(ByteString.copyFrom(dataDecodingEncodingService.encode(task)))
.setTs(task.getTs())
.setAttempt(0)
.build()) .build())
.build()), null); .build()), null);
log.trace("[{}][{}][{}] Submitted task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType());
}
@SuppressWarnings("unchecked")
private <T extends HousekeeperTask> HousekeeperTaskProcessor<T> getTaskProcessor(HousekeeperTaskType taskType) {
return Optional.ofNullable((HousekeeperTaskProcessor<T>) taskProcessors.get(taskType))
.orElseThrow(() -> new IllegalArgumentException("Unsupported task type " + taskType));
} }
@PreDestroy @PreDestroy
@ -141,4 +152,5 @@ public class DefaultHousekeeperService implements HousekeeperService {
consumer.unsubscribe(); consumer.unsubscribe();
consumerExecutor.shutdownNow(); consumerExecutor.shutdownNow();
} }
} }

50
application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java

@ -15,16 +15,15 @@
*/ */
package org.thingsboard.server.service.housekeeper; package org.thingsboard.server.service.housekeeper;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; 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.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueMsgHeaders;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
@ -33,13 +32,10 @@ import org.thingsboard.server.queue.util.TbCoreComponent;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.UUID;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToLong;
import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.longToBytes;
@TbCoreComponent @TbCoreComponent
@Service @Service
@Slf4j @Slf4j
@ -77,15 +73,21 @@ public class HousekeeperReprocessingService {
return; return;
} }
for (TbProtoQueueMsg<ToHousekeeperServiceMsg> msg : msgs) { if (msgs.stream().anyMatch(msg -> {
long msgTs = Uuids.unixTimestamp(msg.getKey()); boolean newMsg = msg.getValue().getTask().getTs() >= startTs;
if (msgTs >= startTs) { if (newMsg) {
stop(); log.info("Stopping reprocessing due to msg is new {}", msg);
return; // fixme: we should commit already reprocessed messages
} }
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<ToHousekeeperServiceMsg> msg : msgs) {
try { try {
reprocessTask(msg); housekeeperService.processTask(msg.getValue());// fixme: or should we submit to queue?
Thread.sleep(1000);
} catch (Exception e) { } catch (Exception e) {
log.error("Message processing failed", e); log.error("Message processing failed", e);
} }
@ -106,21 +108,19 @@ public class HousekeeperReprocessingService {
log.info("Started Housekeeper tasks reprocessing"); log.info("Started Housekeeper tasks reprocessing");
} }
private void reprocessTask(TbProtoQueueMsg<ToHousekeeperServiceMsg> msg) { public void submitForReprocessing(ToHousekeeperServiceMsg msg) {
housekeeperService.processTask(msg);// fixme: or should we submit to queue? HousekeeperTaskProto task = msg.getTask();
} int attempt = task.getAttempt() + 1;
msg = msg.toBuilder()
public void submitForReprocessing(TbProtoQueueMsg<ToHousekeeperServiceMsg> msg) { .setTask(task.toBuilder()
TbQueueMsgHeaders msgHeaders = msg.getHeaders(); .setAttempt(attempt)
long reprocessingAttempts = Optional.ofNullable(msgHeaders.get("reprocessingAttempts")) .setTs(System.currentTimeMillis())
.map(header -> bytesToLong(header)) .build())
.orElse(0L); .build();
reprocessingAttempts++;
msgHeaders.put("reprocessingAttempts", longToBytes(reprocessingAttempts));
var producer = producerProvider.getHousekeeperDelayedMsgProducer(); var producer = producerProvider.getHousekeeperDelayedMsgProducer();
TopicPartitionInfo tpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); 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 @PreDestroy

8
application/src/main/java/org/thingsboard/server/service/housekeeper/processor/AlarmsUnassignTaskProcessor.java

@ -16,21 +16,27 @@
package org.thingsboard.server.service.housekeeper.processor; package org.thingsboard.server.service.housekeeper.processor;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component; 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.common.data.id.UserId;
import org.thingsboard.server.dao.housekeeper.data.HousekeeperTaskType; import org.thingsboard.server.dao.housekeeper.data.HousekeeperTaskType;
import org.thingsboard.server.dao.housekeeper.data.AlarmsUnassignHousekeeperTask; import org.thingsboard.server.dao.housekeeper.data.AlarmsUnassignHousekeeperTask;
import org.thingsboard.server.service.entitiy.alarm.TbAlarmService; import org.thingsboard.server.service.entitiy.alarm.TbAlarmService;
import java.util.List;
@Component @Component
@RequiredArgsConstructor @RequiredArgsConstructor
@Slf4j
public class AlarmsUnassignTaskProcessor implements HousekeeperTaskProcessor<AlarmsUnassignHousekeeperTask> { public class AlarmsUnassignTaskProcessor implements HousekeeperTaskProcessor<AlarmsUnassignHousekeeperTask> {
private final TbAlarmService alarmService; private final TbAlarmService alarmService;
@Override @Override
public void process(AlarmsUnassignHousekeeperTask task) throws Exception { public void process(AlarmsUnassignHousekeeperTask task) throws Exception {
alarmService.unassignDeletedUserAlarms(task.getTenantId(), (UserId) task.getEntityId(), task.getUserTitle(), task.getTs()); List<AlarmId> alarms = alarmService.unassignDeletedUserAlarms(task.getTenantId(), (UserId) task.getEntityId(), task.getUserTitle(), task.getTs());
log.trace("[{}][{}] Unassigned {} alarms", task.getTenantId(), task.getEntityId(), alarms.size());
} }
@Override @Override

2
application/src/main/java/org/thingsboard/server/service/housekeeper/processor/TelemetryDeletionTaskProcessor.java

@ -35,7 +35,7 @@ public class TelemetryDeletionTaskProcessor implements HousekeeperTaskProcessor<
@Override @Override
public HousekeeperTaskType getTaskType() { public HousekeeperTaskType getTaskType() {
return HousekeeperTaskType.DELETE_ATTRIBUTES; return HousekeeperTaskType.DELETE_TELEMETRY;
} }
} }

1
application/src/main/resources/thingsboard.yml

@ -1565,6 +1565,7 @@ queue:
print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:60000}" print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:60000}"
housekeeper: housekeeper:
topic: "tb_housekeeper" topic: "tb_housekeeper"
reprocessing-topic: "tb_housekeeper.reprocessing"
poll-interval-ms: "1000" poll-interval-ms: "1000"
vc: vc:
# Default topic name for Kafka, RabbitMQ, etc. # Default topic name for Kafka, RabbitMQ, etc.

2
common/proto/src/main/proto/queue.proto

@ -1398,4 +1398,6 @@ message ToHousekeeperServiceMsg {
message HousekeeperTaskProto { message HousekeeperTaskProto {
bytes value = 1; bytes value = 1;
int64 ts = 2;
int32 attempt = 3;
} }

4
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); 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); ByteBuffer longBuffer = ByteBuffer.allocate(Long.BYTES);
longBuffer.putLong(0, x); longBuffer.putLong(0, x);
return longBuffer.array(); return longBuffer.array();
} }
public static long bytesToLong(byte[] bytes) { protected static long bytesToLong(byte[] bytes) {
return ByteBuffer.wrap(bytes).getLong(); return ByteBuffer.wrap(bytes).getLong();
} }
} }

2
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}") @Value("${queue.core.housekeeper.topic:tb_housekeeper}")
private String housekeeperTopic; private String housekeeperTopic;
@Value("${queue.core.housekeeper.topic:tb_housekeeper.delayed}") @Value("${queue.core.housekeeper.reprocessing-topic:tb_housekeeper.reprocessing}")
private String housekeeperDelayedTopic; private String housekeeperDelayedTopic;
@Value("${queue.core.partitions}") @Value("${queue.core.partitions}")

31
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.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.annotation.Lazy; 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.EntityView;
import org.thingsboard.server.common.data.StringUtils; 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.EdgeId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; 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.alarm.AlarmService;
import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.dao.entityview.EntityViewService; 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.exception.DataValidationException;
import org.thingsboard.server.dao.housekeeper.HousekeeperService; import org.thingsboard.server.dao.housekeeper.CleanUpService;
import org.thingsboard.server.dao.housekeeper.data.HousekeeperTask;
import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.relation.RelationService;
import java.util.Collections; import java.util.Collections;
@ -70,28 +65,8 @@ public abstract class AbstractEntityService {
protected EdgeService edgeService; protected EdgeService edgeService;
@Autowired @Autowired
protected HousekeeperService housekeeperService; @Lazy
protected CleanUpService cleanUpService;
@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));
}
protected void createRelation(TenantId tenantId, EntityRelation relation) { protected void createRelation(TenantId tenantId, EntityRelation relation) {
log.debug("Creating relation: {}", relation); log.debug("Creating relation: {}", relation);

60
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));
}
}

2
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 { public class AlarmsUnassignHousekeeperTask extends HousekeeperTask {
private final String userTitle; private final String userTitle;
private final long ts;
protected AlarmsUnassignHousekeeperTask(User user) { protected AlarmsUnassignHousekeeperTask(User user) {
super(user.getTenantId(), user.getId(), HousekeeperTaskType.UNASSIGN_ALARMS); super(user.getTenantId(), user.getId(), HousekeeperTaskType.UNASSIGN_ALARMS);
this.userTitle = user.getTitle(); this.userTitle = user.getTitle();
this.ts = System.currentTimeMillis();
} }
} }

9
dao/src/main/java/org/thingsboard/server/dao/housekeeper/data/HousekeeperTask.java

@ -15,27 +15,26 @@
*/ */
package org.thingsboard.server.dao.housekeeper.data; 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.User;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import java.io.Serializable; import java.io.Serializable;
/* @Data
* on start, read the retry queue and put the messages back to main queue (save offset)
* */
@Getter
public class HousekeeperTask implements Serializable { public class HousekeeperTask implements Serializable {
private final TenantId tenantId; private final TenantId tenantId;
private final EntityId entityId; private final EntityId entityId;
private final HousekeeperTaskType taskType; private final HousekeeperTaskType taskType;
private final long ts;
protected HousekeeperTask(TenantId tenantId, EntityId entityId, HousekeeperTaskType taskType) { protected HousekeeperTask(TenantId tenantId, EntityId entityId, HousekeeperTaskType taskType) {
this.tenantId = tenantId; this.tenantId = tenantId;
this.entityId = entityId; this.entityId = entityId;
this.taskType = taskType; this.taskType = taskType;
this.ts = System.currentTimeMillis();
} }
public static HousekeeperTask deleteAttributes(TenantId tenantId, EntityId entityId) { public static HousekeeperTask deleteAttributes(TenantId tenantId, EntityId entityId) {

6
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -177,7 +177,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
List<RuleNodeUpdateResult> updatedRuleNodes = new ArrayList<>(); List<RuleNodeUpdateResult> updatedRuleNodes = new ArrayList<>();
List<RuleNode> existingRuleNodes = getRuleChainNodes(tenantId, ruleChainMetaData.getRuleChainId()); List<RuleNode> existingRuleNodes = getRuleChainNodes(tenantId, ruleChainMetaData.getRuleChainId());
for (RuleNode existingNode : existingRuleNodes) { 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()); Integer index = ruleNodeIndexMap.get(existingNode.getId());
RuleNode newRuleNode = null; RuleNode newRuleNode = null;
if (index != null) { if (index != null) {
@ -771,7 +771,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
private void deleteRuleNodes(TenantId tenantId, List<RuleNode> ruleNodes) { private void deleteRuleNodes(TenantId tenantId, List<RuleNode> ruleNodes) {
List<RuleNodeId> ruleNodeIds = ruleNodes.stream().map(RuleNode::getId).collect(Collectors.toList()); List<RuleNodeId> ruleNodeIds = ruleNodes.stream().map(RuleNode::getId).collect(Collectors.toList());
for (var node : ruleNodes) { for (var node : ruleNodes) {
cleanUpRelatedData(tenantId, node.getId()); cleanUpService.cleanUpRelatedData(tenantId, node.getId());
} }
ruleNodeDao.deleteByIdIn(ruleNodeIds); ruleNodeDao.deleteByIdIn(ruleNodeIds);
} }
@ -820,7 +820,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
private void deleteRuleNode(TenantId tenantId, EntityId entityId) { private void deleteRuleNode(TenantId tenantId, EntityId entityId) {
ruleNodeDao.removeById(tenantId, entityId.getId()); ruleNodeDao.removeById(tenantId, entityId.getId());
cleanUpRelatedData(tenantId, entityId); cleanUpService.cleanUpRelatedData(tenantId, entityId);
} }
private final PaginatedRemover<TenantId, RuleChain> tenantRuleChainsRemover = private final PaginatedRemover<TenantId, RuleChain> tenantRuleChainsRemover =

21
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.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; 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.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; 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.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.attributes.AttributesDao; 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.AttributeKvCompositeKey;
import org.thingsboard.server.dao.model.sql.AttributeKvEntity; import org.thingsboard.server.dao.model.sql.AttributeKvEntity;
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService;
@ -132,10 +135,10 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
@Override @Override
public List<AttributeKvEntry> findAll(TenantId tenantId, EntityId entityId, String attributeType) { public List<AttributeKvEntry> findAll(TenantId tenantId, EntityId entityId, String attributeType) {
return DaoUtil.convertDataList(Lists.newArrayList( return DaoUtil.convertDataList(Lists.newArrayList(
attributeKvRepository.findAllByEntityTypeAndEntityIdAndAttributeType( attributeKvRepository.findAllByEntityTypeAndEntityIdAndAttributeType(
entityId.getEntityType(), entityId.getEntityType(),
entityId.getId(), entityId.getId(),
attributeType))); attributeType)));
} }
@Override @Override
@ -188,6 +191,16 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
return futuresList; return futuresList;
} }
@Transactional
@Override
public List<Pair<String, String>> 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) { private AttributeKvCompositeKey getAttributeKvCompositeKey(EntityId entityId, String attributeType, String attributeKey) {
return new AttributeKvCompositeKey( return new AttributeKvCompositeKey(
entityId.getEntityType(), entityId.getEntityType(),

Loading…
Cancel
Save