Browse Source

Alarm performance improvements

pull/2216/head
Andrew Shvayka 7 years ago
parent
commit
90d8bef576
  1. 10
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 2
      application/src/main/resources/thingsboard.yml
  3. 37
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  4. 7
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java
  5. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java
  6. 2
      dao/src/main/resources/sql/schema-entities-idx.sql

10
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) { void process(ActorContext context, TransportToDeviceActorMsgWrapper wrapper) {
boolean reportDeviceActivity = false;
TransportToDeviceActorMsg msg = wrapper.getMsg(); TransportToDeviceActorMsg msg = wrapper.getMsg();
if (msg.hasSessionEvent()) { if (msg.hasSessionEvent()) {
processSessionStateMsgs(msg.getSessionInfo(), msg.getSessionEvent()); processSessionStateMsgs(msg.getSessionInfo(), msg.getSessionEvent());
@ -235,11 +236,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
if (msg.hasPostAttributes()) { if (msg.hasPostAttributes()) {
handlePostAttributesRequest(context, msg.getSessionInfo(), msg.getPostAttributes()); handlePostAttributesRequest(context, msg.getSessionInfo(), msg.getPostAttributes());
reportLogicalDeviceActivity(); reportDeviceActivity = true;
} }
if (msg.hasPostTelemetry()) { if (msg.hasPostTelemetry()) {
handlePostTelemetryRequest(context, msg.getSessionInfo(), msg.getPostTelemetry()); handlePostTelemetryRequest(context, msg.getSessionInfo(), msg.getPostTelemetry());
reportLogicalDeviceActivity(); reportDeviceActivity = true;
} }
if (msg.hasGetAttributes()) { if (msg.hasGetAttributes()) {
handleGetAttributesRequest(context, msg.getSessionInfo(), msg.getGetAttributes()); handleGetAttributesRequest(context, msg.getSessionInfo(), msg.getGetAttributes());
@ -249,11 +250,14 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
if (msg.hasToServerRPCCallRequest()) { if (msg.hasToServerRPCCallRequest()) {
handleClientSideRPCRequest(context, msg.getSessionInfo(), msg.getToServerRPCCallRequest()); handleClientSideRPCRequest(context, msg.getSessionInfo(), msg.getToServerRPCCallRequest());
reportLogicalDeviceActivity(); reportDeviceActivity = true;
} }
if (msg.hasSubscriptionInfo()) { if (msg.hasSubscriptionInfo()) {
handleSessionActivity(context, msg.getSessionInfo(), msg.getSubscriptionInfo()); handleSessionActivity(context, msg.getSessionInfo(), msg.getSubscriptionInfo());
} }
if (reportDeviceActivity) {
reportLogicalDeviceActivity();
}
} }
private void reportLogicalDeviceActivity() { private void reportLogicalDeviceActivity() {

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

@ -215,7 +215,7 @@ actors:
timeout: "${ACTORS_SESSION_SYNC_TIMEOUT:10000}" timeout: "${ACTORS_SESSION_SYNC_TIMEOUT:10000}"
rule: rule:
# Specify thread pool size for database request callbacks executor service # 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 # Specify thread pool size for javascript executor service
js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:50}" js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:50}"
# Specify thread pool size for mail sender executor service # Specify thread pool size for mail sender executor service

37
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.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Comparator;
import java.util.List; import java.util.List;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
@ -325,21 +326,21 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
private AlarmSeverity detectHighestSeverity(List<AlarmInfo> alarms) { private AlarmSeverity detectHighestSeverity(List<AlarmInfo> alarms) {
if (!alarms.isEmpty()) { if (!alarms.isEmpty()) {
List<AlarmInfo> sorted = new ArrayList(alarms); List<AlarmInfo> sorted = new ArrayList(alarms);
sorted.sort((p1, p2) -> p1.getSeverity().compareTo(p2.getSeverity())); sorted.sort(Comparator.comparing(Alarm::getSeverity));
return sorted.get(0).getSeverity(); return sorted.get(0).getSeverity();
} else { } else {
return null; 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); 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); log.debug("Creating Alarm relation: {}", alarmRelation);
relationService.saveRelationAsync(tenantId, alarmRelation).get(); relationService.saveRelation(tenantId, alarmRelation);
} }
private Alarm merge(Alarm existing, Alarm alarm) { 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) { private void createAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus status, boolean createAnyRelation) {
try { if (createAnyRelation) {
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 + 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);
} }
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) { 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.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.getClearSearchStatus().name(), RelationTypeGroup.ALARM)); deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().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);
}
} }
private void updateAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus oldStatus, AlarmStatus newStatus) { private void updateAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus oldStatus, AlarmStatus newStatus) {

7
dao/src/main/java/org/thingsboard/server/dao/sql/alarm/AlarmRepository.java

@ -31,11 +31,8 @@ import java.util.List;
@SqlDao @SqlDao
public interface AlarmRepository extends CrudRepository<AlarmEntity, String> { public interface AlarmRepository extends CrudRepository<AlarmEntity, String> {
@Query("SELECT a FROM AlarmEntity a WHERE a.tenantId = :tenantId AND a.originatorId = :originatorId " + @Query("SELECT a FROM AlarmEntity a WHERE a.originatorId = :originatorId AND a.type = :alarmType ORDER BY startTs DESC")
"AND a.originatorType = :entityType AND a.type = :alarmType ORDER BY a.type ASC, a.id DESC") List<AlarmEntity> findLatestByOriginatorAndType(@Param("originatorId") String originatorId,
List<AlarmEntity> findLatestByOriginatorAndType(@Param("tenantId") String tenantId,
@Param("originatorId") String originatorId,
@Param("entityType") EntityType entityType,
@Param("alarmType") String alarmType, @Param("alarmType") String alarmType,
Pageable pageable); Pageable pageable);
} }

4
dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java

@ -77,11 +77,9 @@ public class JpaAlarmDao extends JpaAbstractDao<AlarmEntity, Alarm> implements A
public ListenableFuture<Alarm> findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type) { public ListenableFuture<Alarm> findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type) {
return service.submit(() -> { return service.submit(() -> {
List<AlarmEntity> latest = alarmRepository.findLatestByOriginatorAndType( List<AlarmEntity> latest = alarmRepository.findLatestByOriginatorAndType(
UUIDConverter.fromTimeUUID(tenantId.getId()),
UUIDConverter.fromTimeUUID(originator.getId()), UUIDConverter.fromTimeUUID(originator.getId()),
originator.getEntityType(),
type, type,
new PageRequest(0, 1)); PageRequest.of(0, 1));
return latest.isEmpty() ? null : DaoUtil.getData(latest.get(0)); return latest.isEmpty() ? null : DaoUtil.getData(latest.get(0));
}); });
} }

2
dao/src/main/resources/sql/schema-entities-idx.sql

@ -14,7 +14,7 @@
-- limitations under the License. -- 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); CREATE INDEX IF NOT EXISTS idx_event_type_entity_id ON event(tenant_id, event_type, entity_type, entity_id);

Loading…
Cancel
Save