From ab0710c13e9a6a346a0278f754644a3f8a0449ad Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 26 Jul 2022 19:15:25 +0300 Subject: [PATCH] Cleanup event tables by partitions --- .../service/ttl/EventsCleanUpService.java | 23 +--- .../src/main/resources/thingsboard.yml | 5 +- .../server/dao/event/EventService.java | 2 +- .../server/common/data/event/EventType.java | 13 ++- .../server/dao/event/BaseEventService.java | 4 +- .../server/dao/event/EventDao.java | 8 +- .../dao/sql/event/EventCleanupRepository.java | 2 +- .../server/dao/sql/event/JpaBaseEventDao.java | 23 ++-- .../sql/event/SqlEventCleanupRepository.java | 104 +++++++++++++++--- .../service/event/BaseEventServiceTest.java | 4 +- .../test/resources/cassandra-test.properties | 2 - 11 files changed, 132 insertions(+), 58 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java index 1fd3892f39..08ee57119c 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java @@ -56,26 +56,9 @@ public class EventsCleanUpService extends AbstractCleanUpService { public void cleanUp() { if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { long ts = System.currentTimeMillis(); - long regularEventStartTs; - long regularEventEndTs; - long debugEventStartTs; - long debugEventEndTs; - - if (ttlInSec > 0) { - regularEventEndTs = ts - TimeUnit.SECONDS.toMillis(ttlInSec); - regularEventStartTs = regularEventEndTs - 2 * executionIntervalInMs; - } else { - regularEventStartTs = regularEventEndTs = 0; - } - - if (debugTtlInSec > 0) { - debugEventEndTs = ts - TimeUnit.SECONDS.toMillis(debugTtlInSec); - debugEventStartTs = debugEventEndTs - 2 * executionIntervalInMs; - } else { - debugEventStartTs = debugEventEndTs = 0; - } - - eventService.cleanupEvents(regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs); + long regularEventExpTs = ttlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(ttlInSec) : 0; + long debugEventExpTs = debugTtlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(debugTtlInSec) : 0; + eventService.cleanupEvents(regularEventExpTs, debugEventExpTs); } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 17b996934e..5015cc2175 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -221,9 +221,6 @@ cassandra: ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" ts_key_value_partitions_max_cache_size: "${TS_KV_PARTITIONS_MAX_CACHE_SIZE:100000}" ts_key_value_ttl: "${TS_KV_TTL:0}" - events_ttl: "${TS_EVENTS_TTL:0}" - # Specify TTL of debug log in seconds. The current value corresponds to one week - debug_events_ttl: "${DEBUG_EVENTS_TTL:604800}" buffer_size: "${CASSANDRA_QUERY_BUFFER_SIZE:200000}" concurrent_limit: "${CASSANDRA_QUERY_CONCURRENT_LIMIT:1000}" permit_max_wait_time: "${PERMIT_MAX_WAIT_TIME:120000}" @@ -287,7 +284,7 @@ sql: ts_key_value_ttl: "${SQL_TTL_TS_TS_KEY_VALUE_TTL:0}" # Number of seconds events: enabled: "${SQL_TTL_EVENTS_ENABLED:true}" - execution_interval_ms: "${SQL_TTL_EVENTS_EXECUTION_INTERVAL:2220000}" # Number of milliseconds (max random initial delay and fixed period). # 37minutes to avoid common interval spikes + execution_interval_ms: "${SQL_TTL_EVENTS_EXECUTION_INTERVAL:3600000}" # Number of milliseconds (max random initial delay and fixed period). events_ttl: "${SQL_TTL_EVENTS_EVENTS_TTL:0}" # Number of seconds debug_events_ttl: "${SQL_TTL_EVENTS_DEBUG_EVENTS_TTL:604800}" # Number of seconds. The current value corresponds to one week edge_events: diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java index b5f4c98e25..2b89959b56 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java @@ -42,6 +42,6 @@ public interface EventService { void removeEvents(TenantId tenantId, EntityId entityId, EventFilter eventFilter, Long startTime, Long endTime); - void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs); + void cleanupEvents(long regularEventExpTs, long debugEventExpTs); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java index 5fa6fb7370..a7909d28b6 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java @@ -18,16 +18,27 @@ package org.thingsboard.server.common.data.event; import lombok.Getter; public enum EventType { - ERROR("error_event", "ERROR"), LC_EVENT("lc_event", "LC_EVENT"), STATS("stats_event", "STATS"), DEBUG_RULE_NODE("rule_node_debug_event", "DEBUG_RULE_NODE"), DEBUG_RULE_CHAIN("rule_chain_debug_event", "DEBUG_RULE_CHAIN"); + ERROR("error_event", "ERROR"), + LC_EVENT("lc_event", "LC_EVENT"), + STATS("stats_event", "STATS"), + DEBUG_RULE_NODE("rule_node_debug_event", "DEBUG_RULE_NODE", true), + DEBUG_RULE_CHAIN("rule_chain_debug_event", "DEBUG_RULE_CHAIN", true); @Getter private final String table; @Getter private final String oldName; + @Getter + private final boolean debug; EventType(String table, String oldName) { + this(table, oldName, false); + } + + EventType(String table, String oldName, boolean debug) { this.table = table; this.oldName = oldName; + this.debug = debug; } } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java index 85088667f1..302e3f0fe8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java @@ -133,8 +133,8 @@ public class BaseEventService implements EventService { } @Override - public void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs) { - eventDao.cleanupEvents(regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs); + public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) { + eventDao.cleanupEvents(regularEventExpTs, debugEventExpTs); } private PageData convert(EntityType entityType, PageData pd) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java index 718cb41d7c..b7ca1282c7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java @@ -67,10 +67,8 @@ public interface EventDao { /** * Executes stored procedure to cleanup old events. Uses separate ttl for debug and other events. - * @param regularEventStartTs the start time of the interval to use to delete non debug events - * @param regularEventEndTs the end time of the interval to use to delete non debug events - * @param debugEventStartTs the start time of the interval to use to delete debug events - * @param debugEventEndTs the end time of the interval to use to delete debug events + * @param regularEventExpTs the expiration time of the regular events + * @param debugEventExpTs the expiration time of the debug events */ - void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs); + void cleanupEvents(long regularEventExpTs, long debugEventExpTs); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventCleanupRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventCleanupRepository.java index 757d69f794..8631db46ff 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventCleanupRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventCleanupRepository.java @@ -17,6 +17,6 @@ package org.thingsboard.server.dao.sql.event; public interface EventCleanupRepository { - void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs); + void cleanupEvents(long eventExpTime, boolean debug); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java index 26117637d5..0a39ab1be0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java @@ -65,7 +65,9 @@ import java.util.function.Function; @Component public class JpaBaseEventDao implements EventDao { - private static final long PARTITION_DURATION = TimeUnit.HOURS.toMillis(1); + public static final long REGULAR_PARTITION_DURATION = TimeUnit.DAYS.toMillis(1); + public static final long DEBUG_PARTITION_DURATION = TimeUnit.HOURS.toMillis(1); + private final Map> partitionsByEventType = new ConcurrentHashMap<>(); private static final ReentrantLock partitionCreationLock = new ReentrantLock(); @@ -169,10 +171,11 @@ public class JpaBaseEventDao implements EventDao { private void savePartitionIfNotExist(Event event) { var partitionsMap = partitionsByEventType.get(event.getType()); - long partitionStartTs = event.getCreatedTime() - (event.getCreatedTime() % PARTITION_DURATION); + var partitionDuration = event.getType().isDebug() ? DEBUG_PARTITION_DURATION : REGULAR_PARTITION_DURATION; + long partitionStartTs = event.getCreatedTime() - (event.getCreatedTime() % partitionDuration); if (partitionsMap.get(partitionStartTs) == null) { - long partitionEndTs = partitionStartTs + PARTITION_DURATION; - savePartition(partitionsMap, new SqlPartition(event.getType().getTable(), partitionStartTs, partitionEndTs, Long.toString(partitionStartTs))); + savePartition(partitionsMap, new SqlPartition(event.getType().getTable(), partitionStartTs, + partitionStartTs + partitionDuration, Long.toString(partitionStartTs))); } } @@ -314,9 +317,15 @@ public class JpaBaseEventDao implements EventDao { } @Override - public void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs) { - log.info("Going to cleanup old events. Interval for regular events: [{}:{}], for debug events: [{}:{}]", regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs); - eventCleanupRepository.cleanupEvents(regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs); + public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) { + if (regularEventExpTs > 0) { + log.info("Going to cleanup regular events with exp time: {}", regularEventExpTs); + eventCleanupRepository.cleanupEvents(regularEventExpTs, false); + } + if (debugEventExpTs > 0) { + log.info("Going to cleanup debug events with exp time: {}", debugEventExpTs); + eventCleanupRepository.cleanupEvents(debugEventExpTs, true); + } } private void parseUUID(String src, String paramName) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java index 6a17f7c740..54b4565097 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java @@ -16,38 +16,116 @@ package org.thingsboard.server.dao.sql.event; import lombok.extern.slf4j.Slf4j; +import org.postgresql.util.PSQLException; import org.springframework.stereotype.Repository; +import org.thingsboard.server.common.data.event.EventType; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; -import java.util.concurrent.TimeUnit; +import java.util.ArrayList; +import java.util.List; + +import static org.thingsboard.server.dao.sql.event.JpaBaseEventDao.DEBUG_PARTITION_DURATION; +import static org.thingsboard.server.dao.sql.event.JpaBaseEventDao.REGULAR_PARTITION_DURATION; @Slf4j @Repository public class SqlEventCleanupRepository extends JpaAbstractDaoListeningExecutorService implements EventCleanupRepository { + private static final String SELECT_PARTITIONS_STMT = "SELECT tablename from pg_tables WHERE schemaname = 'public' and tablename like concat(?, '_%')"; + private static final int PSQL_VERSION_14 = 140000; + + private volatile Integer currentServerVersion; + @Override - public void cleanupEvents(long regularEventStartTs, long regularEventEndTs, long debugEventStartTs, long debugEventEndTs) { + public void cleanupEvents(long eventExpTime, boolean debug) { + for (EventType eventType : EventType.values()) { + if (eventType.isDebug() == debug) { + cleanupEvents(eventType, eventExpTime); + } + } + } + + private void cleanupEvents(EventType eventType, long eventExpTime) { + var partitionDuration = eventType.isDebug() ? DEBUG_PARTITION_DURATION : REGULAR_PARTITION_DURATION; + List partitions = fetchPartitions(eventType); + for (var partitionTs : partitions) { + var partitionEndTs = partitionTs + partitionDuration; + if (partitionEndTs < eventExpTime) { + log.info("[{}] Detaching expired partition: [{}-{}]", eventType, partitionTs, partitionEndTs); + if (detachAndDropPartition(eventType, partitionTs)) { + log.info("[{}] Detached expired partition: {}", eventType, partitionTs); + } + } else { + log.debug("[{}] Skip valid partition: {}", eventType, partitionTs); + } + } + } + + private List fetchPartitions(EventType eventType) { + List partitions = new ArrayList<>(); try (Connection connection = dataSource.getConnection(); - PreparedStatement stmt = connection.prepareStatement("call cleanup_events_by_ttl(?,?,?,?,?)")) { - stmt.setLong(1, regularEventStartTs); - stmt.setLong(2, regularEventEndTs); - stmt.setLong(3, debugEventStartTs); - stmt.setLong(4, debugEventEndTs); - stmt.setLong(5, 0); - stmt.setQueryTimeout((int) TimeUnit.HOURS.toSeconds(1)); + PreparedStatement stmt = connection.prepareStatement(SELECT_PARTITIONS_STMT)) { + stmt.setString(1, eventType.getTable()); stmt.execute(); - printWarnings(stmt); - try (ResultSet resultSet = stmt.getResultSet()){ - resultSet.next(); - log.info("Total events removed by TTL: [{}]", resultSet.getLong(1)); + try (ResultSet resultSet = stmt.getResultSet()) { + while (resultSet.next()) { + String partitionTableName = resultSet.getString(1); + String partitionTsStr = partitionTableName.substring(eventType.getTable().length() + 1); + try { + partitions.add(Long.parseLong(partitionTsStr)); + } catch (NumberFormatException nfe) { + log.warn("Failed to parse table name: {}", partitionTableName); + } + } } } catch (SQLException e) { log.error("SQLException occurred during events TTL task execution ", e); } + return partitions; + } + + private boolean detachAndDropPartition(EventType eventType, long partitionTs) { + String tablePartition = eventType.getTable() + "_" + partitionTs; + String detachPsqlStmtStr = "ALTER TABLE " + eventType.getTable() + " DETACH PARTITION " + tablePartition; + if (getCurrentServerVersion() >= PSQL_VERSION_14) { + detachPsqlStmtStr += " CONCURRENTLY"; + } + + String dropStmtStr = "DROP TABLE " + tablePartition; + try (Connection connection = dataSource.getConnection(); + PreparedStatement detachStmt = connection.prepareStatement(detachPsqlStmtStr); + PreparedStatement dropStmt = connection.prepareStatement(dropStmtStr)) { + detachStmt.execute(); + dropStmt.execute(); + return true; + } catch (SQLException e) { + log.error("[{}] SQLException occurred during detach and drop of the partition: {}", eventType, partitionTs, e); + } + return false; + } + + private synchronized int getCurrentServerVersion() { + if (currentServerVersion == null) { + try (Connection connection = dataSource.getConnection(); + PreparedStatement versionStmt = connection.prepareStatement("SELECT current_setting('server_version_num')")) { + versionStmt.execute(); + try (ResultSet resultSet = versionStmt.getResultSet()) { + while (resultSet.next()) { + currentServerVersion = resultSet.getInt(1); + } + } + } catch (SQLException e) { + log.warn("SQLException occurred during fetch of the server version", e); + } + if (currentServerVersion == null) { + currentServerVersion = 0; + } + } + return currentServerVersion; } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java index eee70f6624..4bb315e5dc 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java @@ -96,7 +96,7 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent3.getUuidId())); Assert.assertFalse(events.hasNext()); - eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, timeBeforeStartTime - 1, timeAfterEndTime + 1); + eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1); } @Test @@ -126,7 +126,7 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent.getUuidId())); Assert.assertFalse(events.hasNext()); - eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, timeBeforeStartTime - 1, timeAfterEndTime + 1); + eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1); } private EventInfo saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception { diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 4b3ea0a74d..5765153a6f 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -58,8 +58,6 @@ cassandra.query.ts_key_value_partitions_max_cache_size=100000 cassandra.query.ts_key_value_ttl=0 -cassandra.query.debug_events_ttl=604800 - cassandra.query.max_limit_per_request=1000 cassandra.query.buffer_size=100000 cassandra.query.concurrent_limit=1000