Browse Source

Merge pull request #9418 from volodymyr-babak/hotfix/edge-event-processing-executor-service

Edge events are processed in a separate threads to avoid blocking cor…
pull/9431/head
Andrew Shvayka 3 years ago
committed by GitHub
parent
commit
5a1031a8db
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 161
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  2. 1
      application/src/main/resources/thingsboard.yml

161
application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java

@ -16,16 +16,13 @@
package org.thingsboard.server.service.edge; package org.thingsboard.server.service.edge;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
@ -59,7 +56,7 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit;
@Service @Service
@TbCoreComponent @TbCoreComponent
@ -128,17 +125,20 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Autowired @Autowired
protected ApplicationEventPublisher eventPublisher; protected ApplicationEventPublisher eventPublisher;
private ExecutorService dbCallBackExecutor; @Value("${actors.system.edge_dispatcher_pool_size:4}")
private int edgeDispatcherSize;
private ExecutorService executor;
@PostConstruct @PostConstruct
public void initExecutor() { public void initExecutor() {
dbCallBackExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-notifications")); executor = ThingsBoardExecutors.newWorkStealingPool(edgeDispatcherSize, "edge-notifications");
} }
@PreDestroy @PreDestroy
public void shutdownExecutor() { public void shutdownExecutor() {
if (dbCallBackExecutor != null) { if (executor != null) {
dbCallBackExecutor.shutdownNow(); executor.shutdownNow();
} }
} }
@ -157,79 +157,78 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) { public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) {
TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB())); TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB()));
log.debug("[{}] Pushing notification to edge {}", tenantId, edgeNotificationMsg); log.debug("[{}] Pushing notification to edge {}", tenantId, edgeNotificationMsg);
final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(60);
try { try {
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); executor.submit(() -> {
ListenableFuture<Void> future; try {
switch (type) { if (deadline < System.nanoTime()) {
case EDGE: log.warn("[{}] Skipping notification message because deadline reached {}", tenantId, edgeNotificationMsg);
future = edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg); return;
break; }
case ASSET: EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());
future = assetProcessor.processEntityNotification(tenantId, edgeNotificationMsg); switch (type) {
break; case EDGE:
case DEVICE: edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg);
future = deviceProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case ASSET:
case ENTITY_VIEW: assetProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = entityViewProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case DEVICE:
case DASHBOARD: deviceProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = dashboardProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case ENTITY_VIEW:
case RULE_CHAIN: entityViewProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = ruleChainProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case DASHBOARD:
case USER: dashboardProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = userProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case RULE_CHAIN:
case CUSTOMER: ruleChainProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg); break;
break; case USER:
case DEVICE_PROFILE: userProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = deviceProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case CUSTOMER:
case ASSET_PROFILE: customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg);
future = assetProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case DEVICE_PROFILE:
case OTA_PACKAGE: deviceProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = otaPackageProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case ASSET_PROFILE:
case WIDGETS_BUNDLE: assetProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = widgetBundleProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case OTA_PACKAGE:
case WIDGET_TYPE: otaPackageProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = widgetTypeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case WIDGETS_BUNDLE:
case QUEUE: widgetBundleProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = queueProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case WIDGET_TYPE:
case ALARM: widgetTypeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg); break;
break; case QUEUE:
case RELATION: queueProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
future = relationProcessor.processRelationNotification(tenantId, edgeNotificationMsg); break;
break; case ALARM:
case TENANT: alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg);
future = tenantEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case RELATION:
case TENANT_PROFILE: relationProcessor.processRelationNotification(tenantId, edgeNotificationMsg);
future = tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); break;
break; case TENANT:
default: tenantEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); break;
future = Futures.immediateFuture(null); case TENANT_PROFILE:
} tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
Futures.addCallback(future, new FutureCallback<>() { break;
@Override default:
public void onSuccess(@Nullable Void unused) { log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type);
callback.onSuccess(); }
} } catch (Exception e) {
callBackFailure(tenantId, edgeNotificationMsg, callback, e);
@Override
public void onFailure(Throwable throwable) {
callBackFailure(tenantId, edgeNotificationMsg, callback, throwable);
} }
}, dbCallBackExecutor); });
callback.onSuccess();
} catch (Exception e) { } catch (Exception e) {
callBackFailure(tenantId, edgeNotificationMsg, callback, e); callBackFailure(tenantId, edgeNotificationMsg, callback, e);
} }

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

@ -364,6 +364,7 @@ actors:
tenant_dispatcher_pool_size: "${ACTORS_SYSTEM_TENANT_DISPATCHER_POOL_SIZE:2}" tenant_dispatcher_pool_size: "${ACTORS_SYSTEM_TENANT_DISPATCHER_POOL_SIZE:2}"
device_dispatcher_pool_size: "${ACTORS_SYSTEM_DEVICE_DISPATCHER_POOL_SIZE:4}" device_dispatcher_pool_size: "${ACTORS_SYSTEM_DEVICE_DISPATCHER_POOL_SIZE:4}"
rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:8}" rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:8}"
edge_dispatcher_pool_size: "${ACTORS_SYSTEM_EDGE_DISPATCHER_POOL_SIZE:4}"
tenant: tenant:
create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}" create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}"
session: session:

Loading…
Cancel
Save