Browse Source

Refactor SNMP devices' sessions establishing

pull/4520/head
Viacheslav Klimov 5 years ago
committed by Andrew Shvayka
parent
commit
1de97ad0e1
  1. 34
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
  2. 21
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/ProtoTransportEntityService.java

34
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java

@ -55,7 +55,6 @@ import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedDeque; import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.stream.Collectors;
@TbSnmpTransportComponent @TbSnmpTransportComponent
@Component @Component
@ -72,25 +71,30 @@ public class SnmpTransportContext extends TransportContext {
private final SnmpAuthService snmpAuthService; private final SnmpAuthService snmpAuthService;
private final Map<DeviceId, DeviceSessionContext> sessions = new ConcurrentHashMap<>(); private final Map<DeviceId, DeviceSessionContext> sessions = new ConcurrentHashMap<>();
private Collection<DeviceId> allSnmpDevicesIds = new ConcurrentLinkedDeque<>(); private final Collection<DeviceId> allSnmpDevicesIds = new ConcurrentLinkedDeque<>();
@AfterStartUp(order = 2) @AfterStartUp(order = 2)
public void initDevicesSessions() { public void fetchDevicesAndEstablishSessions() {
log.info("Initializing SNMP devices sessions"); log.info("Initializing SNMP devices sessions");
allSnmpDevicesIds = protoEntityService.getAllSnmpDevicesIds().stream()
.map(DeviceId::new)
.collect(Collectors.toList());
log.trace("Found all SNMP devices ids: {}", allSnmpDevicesIds);
List<DeviceId> managedDevicesIds = allSnmpDevicesIds.stream() int batchIndex = 0;
.filter(deviceId -> balancingService.isManagedByCurrentTransport(deviceId.getId())) int batchSize = 512;
.collect(Collectors.toList()); boolean nextBatchExists = true;
log.info("SNMP devices managed by current SNMP transport: {}", managedDevicesIds);
managedDevicesIds.stream() while (nextBatchExists) {
.map(protoEntityService::getDeviceById) TransportProtos.GetSnmpDevicesResponseMsg snmpDevicesResponse = protoEntityService.getSnmpDevicesIds(batchIndex, batchSize);
.collect(Collectors.toList()) snmpDevicesResponse.getIdsList().stream()
.forEach(this::establishDeviceSession); .map(id -> new DeviceId(UUID.fromString(id)))
.peek(allSnmpDevicesIds::add)
.filter(deviceId -> balancingService.isManagedByCurrentTransport(deviceId.getId()))
.map(protoEntityService::getDeviceById)
.forEach(device -> getExecutor().execute(() -> establishDeviceSession(device)));
nextBatchExists = snmpDevicesResponse.getHasNextPage();
batchIndex++;
}
log.debug("Found all SNMP devices ids: {}", allSnmpDevicesIds);
} }
private void establishDeviceSession(Device device) { private void establishDeviceSession(Device device) {

21
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/ProtoTransportEntityService.java

@ -81,26 +81,7 @@ public class ProtoTransportEntityService {
.orElseThrow(() -> new IllegalArgumentException("Device credentials not found")); .orElseThrow(() -> new IllegalArgumentException("Device credentials not found"));
} }
public List<UUID> getAllSnmpDevicesIds() { public TransportProtos.GetSnmpDevicesResponseMsg getSnmpDevicesIds(int page, int pageSize) {
List<UUID> result = new ArrayList<>();
int page = 0;
int pageSize = 512;
boolean hasNextPage = true;
while (hasNextPage) {
TransportProtos.GetSnmpDevicesResponseMsg responseMsg = requestSnmpDevicesIds(page, pageSize);
result.addAll(responseMsg.getIdsList().stream()
.map(UUID::fromString)
.collect(Collectors.toList()));
hasNextPage = responseMsg.getHasNextPage();
page++;
}
return result;
}
private TransportProtos.GetSnmpDevicesResponseMsg requestSnmpDevicesIds(int page, int pageSize) {
TransportProtos.GetSnmpDevicesRequestMsg requestMsg = TransportProtos.GetSnmpDevicesRequestMsg.newBuilder() TransportProtos.GetSnmpDevicesRequestMsg requestMsg = TransportProtos.GetSnmpDevicesRequestMsg.newBuilder()
.setPage(page) .setPage(page)
.setPageSize(pageSize) .setPageSize(pageSize)

Loading…
Cancel
Save