Browse Source

Upgrade script for the events migration

pull/7001/head
Andrii Shvaika 4 years ago
parent
commit
a38edcd62b
  1. 215
      application/src/main/data/upgrade/3.4.0/schema_update.sql
  2. 3
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  3. 12
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  4. 8
      application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java
  5. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java
  6. 15
      common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java
  7. 3
      common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java
  8. 6
      common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java
  9. 4
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  10. 7
      dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
  11. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java
  12. 21
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  13. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java
  14. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java
  15. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java
  16. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java
  17. 2
      dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
  18. 77
      dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java
  19. 214
      dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java

215
application/src/main/data/upgrade/3.4.0/schema_update.sql

@ -0,0 +1,215 @@
--
-- Copyright © 2016-2022 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.
--
CREATE TABLE IF NOT EXISTS rule_node_debug_event (
id uuid NOT NULL,
tenant_id uuid NOT NULL ,
ts bigint NOT NULL,
entity_id uuid NOT NULL,
service_id varchar,
e_type varchar,
e_entity_id uuid,
e_entity_type varchar,
e_msg_id uuid,
e_msg_type varchar,
e_data_type varchar,
e_relation_type varchar,
e_data varchar,
e_metadata varchar,
e_error varchar
) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS rule_chain_debug_event (
id uuid NOT NULL,
tenant_id uuid NOT NULL,
ts bigint NOT NULL,
entity_id uuid NOT NULL,
service_id varchar NOT NULL,
e_message varchar,
e_error varchar
) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS stats_event (
id uuid NOT NULL,
tenant_id uuid NOT NULL,
ts bigint NOT NULL,
entity_id uuid NOT NULL,
service_id varchar NOT NULL,
e_messages_processed bigint NOT NULL,
e_errors_occurred bigint NOT NULL
) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS lc_event (
id uuid NOT NULL,
tenant_id uuid NOT NULL,
ts bigint NOT NULL,
entity_id uuid NOT NULL,
service_id varchar NOT NULL,
e_type varchar NOT NULL,
e_success boolean NOT NULL,
e_error varchar
) PARTITION BY RANGE (ts);
CREATE TABLE IF NOT EXISTS error_event (
id uuid NOT NULL,
tenant_id uuid NOT NULL,
ts bigint NOT NULL,
entity_id uuid NOT NULL,
service_id varchar NOT NULL,
e_method varchar NOT NULL,
e_error varchar
) PARTITION BY RANGE (ts);
CREATE INDEX IF NOT EXISTS idx_rule_node_debug_event_main
ON rule_node_debug_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95);
CREATE INDEX IF NOT EXISTS idx_rule_chain_debug_event_main
ON rule_chain_debug_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95);
CREATE INDEX IF NOT EXISTS idx_stats_event_main
ON stats_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95);
CREATE INDEX IF NOT EXISTS idx_lc_event_main
ON lc_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95);
CREATE INDEX IF NOT EXISTS idx_error_event_main
ON error_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95);
-- Useful to migrate old events to the new table structure;
CREATE OR REPLACE PROCEDURE migrate_regular_events(IN start_ts_in_ms bigint, IN partition_size_in_hours int)
LANGUAGE plpgsql AS
$$
DECLARE
partition_size_in_ms bigint;
p record;
table_name varchar;
BEGIN
partition_size_in_ms = partition_size_in_hours * 3600 * 1000;
FOR p IN SELECT DISTINCT event_type as event_type, (created_time - created_time % partition_size_in_ms) as partition_ts FROM event e WHERE e.event_type in ('STATS', 'LC_EVENT', 'ERROR') and ts > start_ts_in_ms
LOOP
IF p.event_type = 'STATS' THEN
table_name := 'stats_event';
ELSEIF p.event_type = 'LC_EVENT' THEN
table_name := 'lc_event';
ELSEIF p.event_type = 'ERROR' THEN
table_name := 'error_event';
END IF;
RAISE NOTICE '[%] Partition to create : [%-%]', table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms);
EXECUTE format('CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ( %s ) TO ( %s )', table_name, p.partition_ts, table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms));
END LOOP;
INSERT INTO stats_event
SELECT id,
tenant_id,
ts,
entity_id,
body::json ->> 'server',
(body::json ->> 'messagesProcessed')::bigint,
(body::json ->> 'errorsOccurred')::bigint
FROM event
WHERE ts > start_ts_in_ms
AND event_type = 'STATS'
ON CONFLICT DO NOTHING;
INSERT INTO lc_event
SELECT id,
tenant_id,
ts,
entity_id,
body::json ->> 'server',
body::json ->> 'event',
(body::json ->> 'success')::boolean,
body::json ->> 'error'
FROM event
WHERE ts > start_ts_in_ms
AND event_type = 'LC_EVENT'
ON CONFLICT DO NOTHING;
INSERT INTO error_event
SELECT id,
tenant_id,
ts,
entity_id,
body::json ->> 'server',
body::json ->> 'method',
body::json ->> 'error'
FROM event
WHERE ts > start_ts_in_ms
AND event_type = 'ERROR'
ON CONFLICT DO NOTHING;
END
$$;
-- Useful to migrate old debug events to the new table structure;
CREATE OR REPLACE PROCEDURE migrate_debug_events(IN start_ts_in_ms bigint, IN partition_size_in_hours int)
LANGUAGE plpgsql AS
$$
DECLARE
partition_size_in_ms bigint;
p record;
table_name varchar;
BEGIN
partition_size_in_ms = partition_size_in_hours * 3600 * 1000;
FOR p IN SELECT DISTINCT event_type as event_type, (created_time - created_time % partition_size_in_ms) as partition_ts FROM event e WHERE e.event_type in ('DEBUG_RULE_NODE', 'DEBUG_RULE_CHAIN') and ts > start_ts_in_ms
LOOP
IF p.event_type = 'DEBUG_RULE_NODE' THEN
table_name := 'rule_node_debug_event';
ELSEIF p.event_type = 'DEBUG_RULE_CHAIN' THEN
table_name := 'rule_chain_debug_event';
END IF;
RAISE NOTICE '[%] Partition to create : [%-%]', table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms);
EXECUTE format('CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ( %s ) TO ( %s )', table_name, p.partition_ts, table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms));
END LOOP;
INSERT INTO rule_node_debug_event
SELECT id,
tenant_id,
ts,
entity_id,
body::json ->> 'server',
body::json ->> 'type',
(body::json ->> 'entityId')::uuid,
body::json ->> 'entityName',
(body::json ->> 'msgId')::uuid,
body::json ->> 'msgType',
body::json ->> 'dataType',
body::json ->> 'relationType',
body::json ->> 'data',
body::json ->> 'metadata',
body::json ->> 'error'
FROM event
WHERE ts > start_ts_in_ms
AND event_type = 'DEBUG_RULE_NODE'
ON CONFLICT DO NOTHING;
INSERT INTO rule_chain_debug_event
SELECT id,
tenant_id,
ts,
entity_id,
body::json ->> 'server',
body::json ->> 'message',
body::json ->> 'error'
FROM event
WHERE ts > start_ts_in_ms
AND event_type = 'DEBUG_RULE_CHAIN'
ON CONFLICT DO NOTHING;
END
$$;

3
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -221,6 +221,9 @@ public class ThingsboardInstallService {
log.info("Upgrading ThingsBoard from version 3.3.4 to 3.4.0 ...");
databaseEntitiesUpgradeService.upgradeDatabase("3.3.4");
dataUpdateService.updateData("3.3.4");
case "3.4.0":
log.info("Upgrading ThingsBoard from version 3.4.0 to 3.4.1 ...");
databaseEntitiesUpgradeService.upgradeDatabase("3.4.0");
log.info("Updating system data...");
systemDataLoaderService.updateSystemWidgets();
break;

12
application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java

@ -596,6 +596,18 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
log.error("Failed updating schema!!!", e);
}
break;
case "3.4.0":
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating schema ...");
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.4.0", SCHEMA_UPDATE_SQL);
loadSql(schemaUpdateFile, conn);
log.info("Updating schema settings...");
conn.createStatement().execute("UPDATE tb_schema_settings SET schema_version = 3004001;");
log.info("Schema updated.");
} catch (Exception e) {
log.error("Failed updating schema!!!", e);
}
break;
default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
}

8
application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java

@ -25,7 +25,6 @@ import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.concurrent.TimeUnit;
@TbCoreComponent
@Slf4j
@Service
public class EventsCleanUpService extends AbstractCleanUpService {
@ -39,9 +38,6 @@ public class EventsCleanUpService extends AbstractCleanUpService {
@Value("${sql.ttl.events.debug_events_ttl}")
private long debugTtlInSec;
@Value("${sql.ttl.events.execution_interval_ms}")
private long executionIntervalInMs;
@Value("${sql.ttl.events.enabled}")
private boolean ttlTaskExecutionEnabled;
@ -54,11 +50,11 @@ public class EventsCleanUpService extends AbstractCleanUpService {
@Scheduled(initialDelayString = RANDOM_DELAY_INTERVAL_MS_EXPRESSION, fixedDelayString = "${sql.ttl.events.execution_interval_ms}")
public void cleanUp() {
if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) {
if (ttlTaskExecutionEnabled) {
long ts = System.currentTimeMillis();
long regularEventExpTs = ttlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(ttlInSec) : 0;
long debugEventExpTs = debugTtlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(debugTtlInSec) : 0;
eventService.cleanupEvents(regularEventExpTs, debugEventExpTs);
eventService.cleanupEvents(regularEventExpTs, debugEventExpTs, isSystemTenantPartitionMine());
}
}

3
common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java

@ -26,7 +26,6 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import java.util.List;
import java.util.Optional;
public interface EventService {
@ -42,6 +41,6 @@ public interface EventService {
void removeEvents(TenantId tenantId, EntityId entityId, EventFilter eventFilter, Long startTime, Long endTime);
void cleanupEvents(long regularEventExpTs, long debugEventExpTs);
void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb);
}

15
common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.data.event;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.Builder;
import lombok.EqualsAndHashCode;
import lombok.Getter;
@ -40,9 +41,11 @@ public class ErrorEvent extends Event {
this.error = error;
}
@Getter @Setter
@Getter
@Setter
private String method;
@Getter @Setter
@Getter
@Setter
private String error;
@Override
@ -52,6 +55,12 @@ public class ErrorEvent extends Event {
@Override
public EventInfo toInfo(EntityType entityType) {
return null;
EventInfo eventInfo = super.toInfo(entityType);
var json = (ObjectNode) eventInfo.getBody();
json.put("method", method);
if (error != null) {
json.put("error", error);
}
return eventInfo;
}
}

3
common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java

@ -84,7 +84,8 @@ public class RuleNodeDebugEvent extends Event {
var json = (ObjectNode) eventInfo.getBody();
json.put("type", eventType);
if (eventEntity != null) {
json.put("entityId", eventEntity.getId().toString()).put("entityType", eventEntity.getEntityType().name());
json.put("entityId", eventEntity.getId().toString())
.put("entityType", eventEntity.getEntityType().name());
}
if (msgId != null) {
json.put("msgId", msgId.toString());

6
common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.data.event;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.Builder;
import lombok.EqualsAndHashCode;
import lombok.Getter;
@ -51,6 +52,9 @@ public class StatisticsEvent extends Event {
@Override
public EventInfo toInfo(EntityType entityType) {
return null;
EventInfo eventInfo = super.toInfo(entityType);
var json = (ObjectNode) eventInfo.getBody();
json.put("messagesProcessed", messagesProcessed).put("errorsOccurred", errorsOccurred);
return eventInfo;
}
}

4
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -125,8 +125,8 @@ public class BaseEventService implements EventService {
}
@Override
public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) {
eventDao.cleanupEvents(regularEventExpTs, debugEventExpTs);
public void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb) {
eventDao.cleanupEvents(regularEventExpTs, debugEventExpTs, cleanupDb);
}
private PageData<EventInfo> convert(EntityType entityType, PageData<? extends Event> pd) {

7
dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java

@ -16,15 +16,11 @@
package org.thingsboard.server.dao.event;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.EventInfo;
import org.thingsboard.server.common.data.event.Event;
import org.thingsboard.server.common.data.event.EventFilter;
import org.thingsboard.server.common.data.event.EventType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.Dao;
import java.util.List;
import java.util.UUID;
@ -70,8 +66,9 @@ public interface EventDao {
* Executes stored procedure to cleanup old events. Uses separate ttl for debug and other events.
* @param regularEventExpTs the expiration time of the regular events
* @param debugEventExpTs the expiration time of the debug events
* @param cleanupDb
*/
void cleanupEvents(long regularEventExpTs, long debugEventExpTs);
void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb);
/**
* Removes all events for the specified entity and time interval

4
dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java

@ -35,8 +35,8 @@ import java.util.UUID;
public interface ErrorEventRepository extends EventRepository<ErrorEventEntity, ErrorEvent>, JpaRepository<ErrorEventEntity, UUID> {
@Override
@Query("SELECT e FROM ErrorEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC")
List<ErrorEventEntity> findLatestEvents(UUID tenantId, UUID entityId, int limit);
@Query(nativeQuery = true, value = "SELECT * FROM error_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit")
List<ErrorEventEntity> findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit);
@Override
@Query("SELECT e FROM ErrorEventEntity e WHERE " +

21
dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java

@ -427,14 +427,29 @@ public class JpaBaseEventDao implements EventDao {
}
@Override
public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) {
public void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb) {
if (regularEventExpTs > 0) {
log.info("Going to cleanup regular events with exp time: {}", regularEventExpTs);
eventCleanupRepository.cleanupEvents(regularEventExpTs, false);
if (cleanupDb) {
eventCleanupRepository.cleanupEvents(regularEventExpTs, false);
}
cleanupPartitions(regularEventExpTs, false);
}
if (debugEventExpTs > 0) {
log.info("Going to cleanup debug events with exp time: {}", debugEventExpTs);
eventCleanupRepository.cleanupEvents(debugEventExpTs, true);
if (cleanupDb) {
eventCleanupRepository.cleanupEvents(debugEventExpTs, true);
}
cleanupPartitions(debugEventExpTs, true);
}
}
private void cleanupPartitions(long expTime, boolean isDebug) {
for (EventType eventType : EventType.values()) {
if (eventType.isDebug() == isDebug) {
Map<Long, SqlPartition> partitions = partitionsByEventType.get(eventType);
partitions.keySet().removeIf(startTs -> startTs + partitionConfiguration.getPartitionSizeInMs(eventType) < expTime);
}
}
}

5
dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java

@ -24,6 +24,7 @@ import org.springframework.data.repository.query.Param;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.event.LifecycleEvent;
import org.thingsboard.server.dao.model.sql.LifecycleEventEntity;
import org.thingsboard.server.dao.model.sql.RuleChainDebugEventEntity;
import java.util.List;
import java.util.UUID;
@ -31,8 +32,8 @@ import java.util.UUID;
public interface LifecycleEventRepository extends EventRepository<LifecycleEventEntity, LifecycleEvent>, JpaRepository<LifecycleEventEntity, UUID> {
@Override
@Query("SELECT e FROM LifecycleEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC")
List<LifecycleEventEntity> findLatestEvents(UUID tenantId, UUID entityId, int limit);
@Query(nativeQuery = true, value = "SELECT * FROM lc_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit")
List<LifecycleEventEntity> findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit);
@Query("SELECT e FROM LifecycleEventEntity e WHERE " +

4
dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java

@ -34,8 +34,8 @@ import java.util.UUID;
public interface RuleChainDebugEventRepository extends EventRepository<RuleChainDebugEventEntity, RuleChainDebugEvent>, JpaRepository<RuleChainDebugEventEntity, UUID> {
@Override
@Query("SELECT e FROM RuleChainDebugEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC")
List<RuleChainDebugEventEntity> findLatestEvents(UUID tenantId, UUID entityId, int limit);
@Query(nativeQuery = true, value = "SELECT * FROM rule_chain_debug_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit")
List<RuleChainDebugEventEntity> findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit);
@Override
@Query("SELECT e FROM RuleChainDebugEventEntity e WHERE " +

5
dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java

@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.event.ErrorEvent;
import org.thingsboard.server.common.data.event.RuleNodeDebugEvent;
import org.thingsboard.server.dao.model.sql.ErrorEventEntity;
import org.thingsboard.server.dao.model.sql.RuleNodeDebugEventEntity;
import org.thingsboard.server.dao.model.sql.StatisticsEventEntity;
import java.util.List;
import java.util.UUID;
@ -34,8 +35,8 @@ import java.util.UUID;
public interface RuleNodeDebugEventRepository extends EventRepository<RuleNodeDebugEventEntity, RuleNodeDebugEvent>, JpaRepository<RuleNodeDebugEventEntity, UUID> {
@Override
@Query("SELECT e FROM RuleNodeDebugEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC")
List<RuleNodeDebugEventEntity> findLatestEvents(UUID tenantId, UUID entityId, int limit);
@Query(nativeQuery = true, value = "SELECT * FROM rule_node_debug_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit")
List<RuleNodeDebugEventEntity> findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit);
@Override
@Query("SELECT e FROM RuleNodeDebugEventEntity e WHERE " +

4
dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java

@ -31,8 +31,8 @@ import java.util.UUID;
public interface StatisticsEventRepository extends EventRepository<StatisticsEventEntity, StatisticsEvent>, JpaRepository<StatisticsEventEntity, UUID> {
@Override
@Query("SELECT e FROM LifecycleEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC")
List<StatisticsEventEntity> findLatestEvents(UUID tenantId, UUID entityId, int limit);
@Query(nativeQuery = true, value = "SELECT * FROM stats_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit")
List<StatisticsEventEntity> findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit);
@Query("SELECT e FROM StatisticsEventEntity e WHERE " +
"e.tenantId = :tenantId " +

2
dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java

@ -188,7 +188,7 @@ public abstract class AbstractServiceTest {
}
protected RuleNodeDebugEvent generateEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) throws IOException {
protected RuleNodeDebugEvent generateEvent(TenantId tenantId, EntityId entityId) throws IOException {
if (tenantId == null) {
tenantId = TenantId.fromUUID(Uuids.timeBased());
}

77
dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java

@ -19,7 +19,6 @@ import com.datastax.oss.driver.api.core.uuid.Uuids;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EventInfo;
import org.thingsboard.server.common.data.event.Event;
import org.thingsboard.server.common.data.event.EventType;
@ -35,7 +34,7 @@ import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.service.AbstractServiceTest;
import java.text.ParseException;
import java.util.Optional;
import java.util.List;
import static org.apache.commons.lang3.time.DateFormatUtils.ISO_DATETIME_TIME_ZONE_FORMAT;
@ -57,16 +56,14 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
@Test
public void saveEvent() throws Exception {
TenantId tenantId = new TenantId(Uuids.timeBased());
DeviceId devId = new DeviceId(Uuids.timeBased());
RuleNodeDebugEvent event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString());
RuleNodeDebugEvent event = generateEvent(tenantId, devId);
eventService.saveAsync(event).get();
throw new RuntimeException("fix me!");
// Optional<EventInfo> loaded = eventService.findEvent(event.getTenantId(), event.getEntityId(), event.getType(), event.getUid());
// Assert.assertTrue(loaded.isPresent());
// Assert.assertNotNull(loaded.get());
// Assert.assertEquals(event.getEntityId(), loaded.get().getEntityId());
// Assert.assertEquals(event.getType(), loaded.get().getType());
// Assert.assertEquals(event.getBody(), loaded.get().getBody());
List<EventInfo> loaded = eventService.findLatestEvents(event.getTenantId(), devId, event.getType(), 1);
Assert.assertNotNull(loaded);
Assert.assertEquals(1, loaded.size());
Assert.assertEquals(event.getData(), loaded.get(0).getBody().get("data").asText());
}
@Test
@ -74,29 +71,29 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
CustomerId customerId = new CustomerId(Uuids.timeBased());
TenantId tenantId = TenantId.fromUUID(Uuids.timeBased());
saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId);
EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId);
EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId);
EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId);
Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId);
Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId);
Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId);
saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId);
TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime"), startTime, endTime);
TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("ts"), startTime, endTime);
PageData<EventInfo> events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink);
PageData<EventInfo> events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink);
Assert.assertNotNull(events.getData());
Assert.assertTrue(events.getData().size() == 2);
Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent.getUuidId()));
Assert.assertTrue(events.getData().get(1).getUuidId().equals(savedEvent2.getUuidId()));
Assert.assertEquals(2, events.getData().size());
Assert.assertEquals(savedEvent.getUuidId(), events.getData().get(0).getUuidId());
Assert.assertEquals(savedEvent2.getUuidId(), events.getData().get(1).getUuidId());
Assert.assertTrue(events.hasNext());
events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink.nextPageLink());
events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink.nextPageLink());
Assert.assertNotNull(events.getData());
Assert.assertTrue(events.getData().size() == 1);
Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent3.getUuidId()));
Assert.assertEquals(1, events.getData().size());
Assert.assertEquals(savedEvent3.getUuidId(), events.getData().get(0).getUuidId());
Assert.assertFalse(events.hasNext());
eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1);
eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, true);
}
@Test
@ -104,36 +101,36 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
CustomerId customerId = new CustomerId(Uuids.timeBased());
TenantId tenantId = TenantId.fromUUID(Uuids.timeBased());
saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId);
EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId);
EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId);
EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId);
Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId);
Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId);
Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId);
saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId);
TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime);
TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("ts", SortOrder.Direction.DESC), startTime, endTime);
PageData<EventInfo> events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink);
PageData<EventInfo> events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink);
Assert.assertNotNull(events.getData());
Assert.assertTrue(events.getData().size() == 2);
Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent3.getUuidId()));
Assert.assertTrue(events.getData().get(1).getUuidId().equals(savedEvent2.getUuidId()));
Assert.assertEquals(2, events.getData().size());
Assert.assertEquals(savedEvent3.getUuidId(), events.getData().get(0).getUuidId());
Assert.assertEquals(savedEvent2.getUuidId(), events.getData().get(1).getUuidId());
Assert.assertTrue(events.hasNext());
events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink.nextPageLink());
events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink.nextPageLink());
Assert.assertNotNull(events.getData());
Assert.assertTrue(events.getData().size() == 1);
Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent.getUuidId()));
Assert.assertEquals(1, events.getData().size());
Assert.assertEquals(savedEvent.getUuidId(), events.getData().get(0).getUuidId());
Assert.assertFalse(events.hasNext());
eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1);
eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, true);
}
private EventInfo saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception {
throw new RuntimeException("fix me!");
// EventInfo event = generateEvent(tenantId, entityId, DataConstants.STATS, null);
// event.setId(new EventId(Uuids.startOf(time)));
// eventService.saveAsync(event).get();
// return event;
private Event saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception {
RuleNodeDebugEvent event = generateEvent(tenantId, entityId);
event.setId(new EventId(Uuids.timeBased()));
event.setCreatedTime(time);
eventService.saveAsync(event).get();
return event;
}
}

214
dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java

@ -17,16 +17,26 @@ package org.thingsboard.server.dao.sql.event;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import lombok.extern.slf4j.Slf4j;
import org.junit.After;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.event.Event;
import org.thingsboard.server.common.data.event.EventType;
import org.thingsboard.server.common.data.event.StatisticsEvent;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.AbstractJpaDaoTest;
import org.thingsboard.server.dao.event.EventDao;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
/**
* Created by Valerii Sosliuk on 5/5/2017.
*/
@Slf4j
public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
@ -34,131 +44,77 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
private EventDao eventDao;
UUID tenantId = Uuids.timeBased();
@After
public void deleteEvents() {
throw new RuntimeException("fix me!");
// List<EventInfo> events = eventDao.find(TenantId.fromUUID(tenantId));
// for (EventInfo event : events) {
// eventDao.removeById(TenantId.fromUUID(tenantId), event.getUuidId());
// }
@Test
public void findEvent() throws InterruptedException, ExecutionException, TimeoutException {
UUID entityId = Uuids.timeBased();
Event event1 = getStatsEvent(Uuids.timeBased(), tenantId, entityId);
eventDao.saveAsync(event1).get(1, TimeUnit.MINUTES);
Thread.sleep(2);
Event event2 = getStatsEvent(Uuids.timeBased(), tenantId, entityId);
eventDao.saveAsync(event2).get(1, TimeUnit.MINUTES);
List<? extends Event> foundEvents = eventDao.findLatestEvents(tenantId, entityId, EventType.STATS, 1);
assertNotNull("Events expected to be not null", foundEvents);
assertEquals(1, foundEvents.size());
assertEquals(event2, foundEvents.get(0));
}
@Test
public void findEventsByEntityIdAndPageLink() throws Exception {
UUID entityId1 = Uuids.timeBased();
UUID entityId2 = Uuids.timeBased();
long startTime = System.currentTimeMillis();
Event event1 = getStatsEvent(Uuids.timeBased(), tenantId, entityId1);
eventDao.saveAsync(event1).get(1, TimeUnit.MINUTES);
Thread.sleep(2);
Event event2 = getStatsEvent(Uuids.timeBased(), tenantId, entityId2);
eventDao.saveAsync(event2).get(1, TimeUnit.MINUTES);
long endTime = System.currentTimeMillis();
PageData<? extends Event> events1 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, new TimePageLink(30));
assertEquals(1, events1.getData().size());
PageData<? extends Event> events2 = eventDao.findEvents(tenantId, entityId2, EventType.STATS, new TimePageLink(30));
assertEquals(1, events2.getData().size());
PageData<? extends Event> events3 = eventDao.findEvents(tenantId, Uuids.timeBased(), EventType.STATS, new TimePageLink(30));
assertEquals(0, events3.getData().size());
TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null);
PageData<? extends Event> events12 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink2);
assertEquals(1, events12.getData().size());
assertEquals(event1, events12.getData().get(0));
TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime);
PageData<? extends Event> events13 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink3);
assertEquals(1, events13.getData().size());
assertEquals(event1, events13.getData().get(0));
TimePageLink pageLink4 = new TimePageLink(5, 0, "", null, startTime, endTime);
PageData<? extends Event> events14 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink4);
assertEquals(1, events14.getData().size());
assertEquals(event1, events14.getData().get(0));
pageLink4 = pageLink4.nextPageLink();
PageData<? extends Event> events6 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink4);
assertEquals(0, events6.getData().size());
}
// @Test
// public void findEvent() {
// UUID entityId = Uuids.timeBased();
// EventInfo savedEvent = eventDao.save(TenantId.fromUUID(tenantId), getEvent(entityId, tenantId, entityId));
// EventInfo foundEvent = eventDao.findEvent(tenantId, new DeviceId(entityId), DataConstants.STATS, savedEvent.getUid());
// assertNotNull("Event expected to be not null", foundEvent);
// assertEquals(savedEvent.getId(), foundEvent.getId());
// }
//
// @Test
// public void findEventsByEntityIdAndPageLink() throws Exception {
// UUID entityId1 = Uuids.timeBased();
// UUID entityId2 = Uuids.timeBased();
// long startTime = System.currentTimeMillis();
// long endTime = createEventsTwoEntities(tenantId, entityId1, entityId2, 20);
//
// TimePageLink pageLink1 = new TimePageLink(30);
// PageData<EventInfo> events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink1);
// assertEquals(10, events1.getData().size());
//
// TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null);
// PageData<EventInfo> events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink2);
// assertEquals(10, events2.getData().size());
//
// TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime);
// PageData<EventInfo> events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink3);
// assertEquals(10, events3.getData().size());
//
// TimePageLink pageLink4 = new TimePageLink(5, 0, "", null, startTime, endTime);
// PageData<EventInfo> events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4);
// assertEquals(5, events4.getData().size());
//
// pageLink4 = pageLink4.nextPageLink();
// PageData<EventInfo> events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4);
// assertEquals(5, events5.getData().size());
//
// pageLink4 = pageLink4.nextPageLink();
// PageData<EventInfo> events6 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4);
// assertEquals(0, events6.getData().size());
//
// }
//
// @Test
// public void findEventsByEntityIdAndEventTypeAndPageLink() throws Exception {
// UUID entityId1 = Uuids.timeBased();
// UUID entityId2 = Uuids.timeBased();
// long startTime = System.currentTimeMillis();
// long endTime = createEventsTwoEntitiesTwoTypes(tenantId, entityId1, entityId2, 20);
//
// TimePageLink pageLink1 = new TimePageLink(30);
// PageData<EventInfo> events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink1);
// assertEquals(5, events1.getData().size());
//
// TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null);
// PageData<EventInfo> events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink2);
// assertEquals(5, events2.getData().size());
//
// TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime);
// PageData<EventInfo> events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink3);
// assertEquals(5, events3.getData().size());
//
// TimePageLink pageLink4 = new TimePageLink(4, 0, "", null, startTime, endTime);
// PageData<EventInfo> events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4);
// assertEquals(4, events4.getData().size());
//
// pageLink4 = pageLink4.nextPageLink();
// PageData<EventInfo> events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4);
// assertEquals(1, events5.getData().size());
// }
//
// private long createEventsTwoEntitiesTwoTypes(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception {
// for (int i = 0; i < count / 2; i++) {
// String type = i % 2 == 0 ? STATS : ALARM;
// UUID eventId1 = Uuids.timeBased();
// EventInfo event1 = getEvent(eventId1, tenantId, entityId1, type);
// eventDao.saveAsync(event1).get();
// UUID eventId2 = Uuids.timeBased();
// EventInfo event2 = getEvent(eventId2, tenantId, entityId2, type);
// eventDao.saveAsync(event2).get();
// }
// return System.currentTimeMillis();
// }
//
// private long createEventsTwoEntities(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception {
// for (int i = 0; i < count / 2; i++) {
// UUID eventId1 = Uuids.timeBased();
// EventInfo event1 = getEvent(eventId1, tenantId, entityId1);
// eventDao.saveAsync(event1).get();
// UUID eventId2 = Uuids.timeBased();
// EventInfo event2 = getEvent(eventId2, tenantId, entityId2);
// eventDao.saveAsync(event2).get();
// }
// return System.currentTimeMillis();
// }
//
// private EventInfo getEvent(UUID eventId, UUID tenantId, UUID entityId, String type) {
// EventInfo event = getEvent(eventId, tenantId, entityId);
// event.setType(type);
// return event;
// }
//
// private EventInfo getEvent(UUID eventId, UUID tenantId, UUID entityId) {
// EventInfo event = new EventInfo();
// event.setId(new EventId(eventId));
// event.setTenantId(TenantId.fromUUID(tenantId));
// EntityId deviceId = new DeviceId(entityId);
// event.setEntityId(deviceId);
// event.setUid(event.getId().getId().toString());
// event.setType(STATS);
// ObjectMapper mapper = new ObjectMapper();
// try {
// JsonNode jsonNode = mapper.readTree("{\"key\":\"value\"}");
// event.setBody(jsonNode);
// } catch (IOException e) {
// log.error(e.getMessage(), e);
// }
// return event;
// }
private Event getStatsEvent(UUID eventId, UUID tenantId, UUID entityId) {
StatisticsEvent.StatisticsEventBuilder event = StatisticsEvent.builder();
event.id(eventId);
event.ts(System.currentTimeMillis());
event.tenantId(new TenantId(tenantId));
event.entityId(entityId);
event.serviceId("server A");
event.messagesProcessed(1);
event.errorsOccurred(0);
return event.build();
}
}

Loading…
Cancel
Save