Browse Source

Merge pull request #10380 from thingsboard/fix/device-state-executor

Add callback executor for device state service
pull/10385/head
Andrew Shvayka 3 years ago
committed by GitHub
parent
commit
749b60943f
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 16
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

16
application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

@ -191,6 +191,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
private int telemetryTtl; private int telemetryTtl;
private ListeningExecutorService deviceStateExecutor; private ListeningExecutorService deviceStateExecutor;
private ListeningExecutorService deviceStateCallbackExecutor;
final ConcurrentMap<DeviceId, DeviceStateData> deviceStates = new ConcurrentHashMap<>(); final ConcurrentMap<DeviceId, DeviceStateData> deviceStates = new ConcurrentHashMap<>();
@ -199,6 +200,8 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
super.init(); super.init();
deviceStateExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( deviceStateExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "device-state")); Math.max(4, Runtime.getRuntime().availableProcessors()), "device-state"));
deviceStateCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "device-state-callback"));
scheduledExecutor.scheduleWithFixedDelay(this::checkStates, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS); scheduledExecutor.scheduleWithFixedDelay(this::checkStates, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS);
scheduledExecutor.scheduleWithFixedDelay(this::reportActivityStats, defaultActivityStatsIntervalInSec, defaultActivityStatsIntervalInSec, TimeUnit.SECONDS); scheduledExecutor.scheduleWithFixedDelay(this::reportActivityStats, defaultActivityStatsIntervalInSec, defaultActivityStatsIntervalInSec, TimeUnit.SECONDS);
} }
@ -209,6 +212,9 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
if (deviceStateExecutor != null) { if (deviceStateExecutor != null) {
deviceStateExecutor.shutdownNow(); deviceStateExecutor.shutdownNow();
} }
if (deviceStateCallbackExecutor != null) {
deviceStateCallbackExecutor.shutdownNow();
}
} }
@Override @Override
@ -372,7 +378,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
log.warn("Failed to register device to the state service", t); log.warn("Failed to register device to the state service", t);
callback.onFailure(t); callback.onFailure(t);
} }
}, deviceStateExecutor); }, deviceStateCallbackExecutor);
} else if (proto.getUpdated()) { } else if (proto.getUpdated()) {
DeviceStateData stateData = getOrFetchDeviceStateData(device.getId()); DeviceStateData stateData = getOrFetchDeviceStateData(device.getId());
TbMsgMetaData md = new TbMsgMetaData(); TbMsgMetaData md = new TbMsgMetaData();
@ -635,10 +641,10 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
ListenableFuture<DeviceStateData> future; ListenableFuture<DeviceStateData> future;
if (persistToTelemetry) { if (persistToTelemetry) {
ListenableFuture<List<TsKvEntry>> tsData = tsService.findLatest(TenantId.SYS_TENANT_ID, device.getId(), PERSISTENT_ATTRIBUTES); ListenableFuture<List<TsKvEntry>> tsData = tsService.findLatest(TenantId.SYS_TENANT_ID, device.getId(), PERSISTENT_ATTRIBUTES);
future = Futures.transform(tsData, extractDeviceStateData(device), deviceStateExecutor); future = Futures.transform(tsData, extractDeviceStateData(device), MoreExecutors.directExecutor());
} else { } else {
ListenableFuture<List<AttributeKvEntry>> attrData = attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), SERVER_SCOPE, PERSISTENT_ATTRIBUTES); ListenableFuture<List<AttributeKvEntry>> attrData = attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), SERVER_SCOPE, PERSISTENT_ATTRIBUTES);
future = Futures.transform(attrData, extractDeviceStateData(device), deviceStateExecutor); future = Futures.transform(attrData, extractDeviceStateData(device), MoreExecutors.directExecutor());
} }
return transformInactivityTimeout(future); return transformInactivityTimeout(future);
} }
@ -656,8 +662,8 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
} }
}); });
return deviceStateData; return deviceStateData;
}, deviceStateExecutor); }, MoreExecutors.directExecutor());
}, deviceStateExecutor); }, deviceStateCallbackExecutor);
} }
private <T extends KvEntry> Function<List<T>, DeviceStateData> extractDeviceStateData(Device device) { private <T extends KvEntry> Function<List<T>, DeviceStateData> extractDeviceStateData(Device device) {

Loading…
Cancel
Save