From 90d8bef576708bad98717f35e08913c3d1547bd9 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Wed, 27 Nov 2019 18:38:55 +0200 Subject: [PATCH] Alarm performance improvements --- .../device/DeviceActorMessageProcessor.java | 10 +++-- .../src/main/resources/thingsboard.yml | 2 +- .../server/dao/alarm/BaseAlarmService.java | 37 +++++++------------ .../server/dao/sql/alarm/AlarmRepository.java | 7 +--- .../server/dao/sql/alarm/JpaAlarmDao.java | 4 +- .../resources/sql/schema-entities-idx.sql | 2 +- 6 files changed, 26 insertions(+), 36 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index aa1f8bc909..6784a2a0f3 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -223,6 +223,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } void process(ActorContext context, TransportToDeviceActorMsgWrapper wrapper) { + boolean reportDeviceActivity = false; TransportToDeviceActorMsg msg = wrapper.getMsg(); if (msg.hasSessionEvent()) { processSessionStateMsgs(msg.getSessionInfo(), msg.getSessionEvent()); @@ -235,11 +236,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } if (msg.hasPostAttributes()) { handlePostAttributesRequest(context, msg.getSessionInfo(), msg.getPostAttributes()); - reportLogicalDeviceActivity(); + reportDeviceActivity = true; } if (msg.hasPostTelemetry()) { handlePostTelemetryRequest(context, msg.getSessionInfo(), msg.getPostTelemetry()); - reportLogicalDeviceActivity(); + reportDeviceActivity = true; } if (msg.hasGetAttributes()) { handleGetAttributesRequest(context, msg.getSessionInfo(), msg.getGetAttributes()); @@ -249,11 +250,14 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } if (msg.hasToServerRPCCallRequest()) { handleClientSideRPCRequest(context, msg.getSessionInfo(), msg.getToServerRPCCallRequest()); - reportLogicalDeviceActivity(); + reportDeviceActivity = true; } if (msg.hasSubscriptionInfo()) { handleSessionActivity(context, msg.getSessionInfo(), msg.getSubscriptionInfo()); } + if (reportDeviceActivity) { + reportLogicalDeviceActivity(); + } } private void reportLogicalDeviceActivity() { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 1e28eb5ef9..b181572247 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -215,7 +215,7 @@ actors: timeout: "${ACTORS_SESSION_SYNC_TIMEOUT:10000}" rule: # Specify thread pool size for database request callbacks executor service - db_callback_thread_pool_size: "${ACTORS_RULE_DB_CALLBACK_THREAD_POOL_SIZE:1}" + db_callback_thread_pool_size: "${ACTORS_RULE_DB_CALLBACK_THREAD_POOL_SIZE:50}" # Specify thread pool size for javascript executor service js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:50}" # Specify thread pool size for mail sender executor service diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java index d83b068cbe..8a4279fe7f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java @@ -52,6 +52,7 @@ import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Set; import java.util.concurrent.ExecutionException; @@ -325,21 +326,21 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ private AlarmSeverity detectHighestSeverity(List alarms) { if (!alarms.isEmpty()) { List sorted = new ArrayList(alarms); - sorted.sort((p1, p2) -> p1.getSeverity().compareTo(p2.getSeverity())); + sorted.sort(Comparator.comparing(Alarm::getSeverity)); return sorted.get(0).getSeverity(); } else { return null; } } - private void deleteRelation(TenantId tenantId, EntityRelation alarmRelation) throws ExecutionException, InterruptedException { + private void deleteRelation(TenantId tenantId, EntityRelation alarmRelation) { log.debug("Deleting Alarm relation: {}", alarmRelation); - relationService.deleteRelationAsync(tenantId, alarmRelation).get(); + relationService.deleteRelation(tenantId, alarmRelation); } - private void createRelation(TenantId tenantId, EntityRelation alarmRelation) throws ExecutionException, InterruptedException { + private void createRelation(TenantId tenantId, EntityRelation alarmRelation) { log.debug("Creating Alarm relation: {}", alarmRelation); - relationService.saveRelationAsync(tenantId, alarmRelation).get(); + relationService.saveRelation(tenantId, alarmRelation); } private Alarm merge(Alarm existing, Alarm alarm) { @@ -376,28 +377,18 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ } private void createAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus status, boolean createAnyRelation) { - try { - if (createAnyRelation) { - createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + AlarmSearchStatus.ANY.name(), RelationTypeGroup.ALARM)); - } - createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM)); - createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM)); - createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM)); - } catch (ExecutionException | InterruptedException e) { - log.warn("[{}] Failed to create relation. Status: [{}]", alarmId, status); - throw new RuntimeException(e); + if (createAnyRelation) { + createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + AlarmSearchStatus.ANY.name(), RelationTypeGroup.ALARM)); } + createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM)); + createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM)); + createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM)); } private void deleteAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus status) { - try { - deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM)); - deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM)); - deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM)); - } catch (ExecutionException | InterruptedException e) { - log.warn("[{}] Failed to delete relation. Status: [{}]", alarmId, status); - throw new RuntimeException(e); - } + deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM)); + deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM)); + deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM)); } private void updateAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus oldStatus, AlarmStatus newStatus) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java index a3d5b45a39..756ebbf070 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java @@ -31,11 +31,8 @@ import java.util.List; @SqlDao public interface AlarmRepository extends CrudRepository { - @Query("SELECT a FROM AlarmEntity a WHERE a.tenantId = :tenantId AND a.originatorId = :originatorId " + - "AND a.originatorType = :entityType AND a.type = :alarmType ORDER BY a.type ASC, a.id DESC") - List findLatestByOriginatorAndType(@Param("tenantId") String tenantId, - @Param("originatorId") String originatorId, - @Param("entityType") EntityType entityType, + @Query("SELECT a FROM AlarmEntity a WHERE a.originatorId = :originatorId AND a.type = :alarmType ORDER BY startTs DESC") + List findLatestByOriginatorAndType(@Param("originatorId") String originatorId, @Param("alarmType") String alarmType, Pageable pageable); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index 4762c5c4ec..8aa9a81094 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -77,11 +77,9 @@ public class JpaAlarmDao extends JpaAbstractDao implements A public ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type) { return service.submit(() -> { List latest = alarmRepository.findLatestByOriginatorAndType( - UUIDConverter.fromTimeUUID(tenantId.getId()), UUIDConverter.fromTimeUUID(originator.getId()), - originator.getEntityType(), type, - new PageRequest(0, 1)); + PageRequest.of(0, 1)); return latest.isEmpty() ? null : DaoUtil.getData(latest.get(0)); }); } diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 9809219ff9..a74d45836d 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/dao/src/main/resources/sql/schema-entities-idx.sql @@ -14,7 +14,7 @@ -- limitations under the License. -- -CREATE INDEX IF NOT EXISTS idx_alarm_originator_alarm_type ON alarm(tenant_id, type, originator_type, originator_id); +CREATE INDEX IF NOT EXISTS idx_alarm_originator_alarm_type ON alarm(originator_id, type, startTs DESC); CREATE INDEX IF NOT EXISTS idx_event_type_entity_id ON event(tenant_id, event_type, entity_type, entity_id);