diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java index 5e327a72db..83855bc8f8 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java @@ -62,17 +62,15 @@ public class NotificationsCleanUpService extends AbstractCleanUpService { public void cleanUp() { long expTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec); long partitionDurationMs = TimeUnit.HOURS.toMillis(partitionSizeInHours); - if (!isSystemTenantPartitionMine()) { + if (isSystemTenantPartitionMine()) { + partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); + } else { partitioningRepository.cleanupPartitionsCache(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); - return; } - long lastRemovedNotificationTs = partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); - if (lastRemovedNotificationTs > 0) { - long gap = TimeUnit.MINUTES.toMillis(10); - long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; - cleanUpNotificationRequests(requestExpTime); - } + long gap = TimeUnit.MINUTES.toMillis(10); + long requestExpTime = expTime - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; + cleanUpNotificationRequests(requestExpTime); } private void cleanUpNotificationRequests(long expirationTime) { @@ -80,13 +78,15 @@ public class NotificationsCleanUpService extends AbstractCleanUpService { int totalRemoved = 0; int tenantsProcessed = 0; - // Clean up SYSADMIN's notification requests: - try { - totalRemoved += cleanUpByTenant(TenantId.SYS_TENANT_ID, expirationTime); - } catch (Exception e) { - log.warn("Failed to clean up notification requests for sysadmin {}", TenantId.SYS_TENANT_ID, e); + // Clean up SYSADMIN's notification requests on the system node only + if (isSystemTenantPartitionMine()) { + try { + totalRemoved += cleanUpByTenant(TenantId.SYS_TENANT_ID, expirationTime); + } catch (Exception e) { + log.warn("Failed to clean up notification requests for sysadmin {}", TenantId.SYS_TENANT_ID, e); + } } - // Clean up notification requests for tenants + // Each node cleans up notification requests for its own tenants PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); for (TenantId tenantId : tenants) { try { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java index 7be106219b..3f76170c84 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java @@ -21,6 +21,7 @@ import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.dao.model.sql.RpcEntity; @@ -34,6 +35,7 @@ public interface RpcRepository extends JpaRepository { Page findAllByTenantId(UUID tenantId, Pageable pageable); + @Transactional @Modifying @Query(value = "DELETE FROM rpc WHERE id IN " + "(SELECT id FROM rpc WHERE tenant_id = :tenantId AND created_time < :expirationTime LIMIT :batchSize)",