diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 7da26fdc81..dfdd3a407c 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1295,6 +1295,8 @@ transport: ignore_type_cast_errors: "${SNMP_RESPONSE_IGNORE_TYPE_CAST_ERRORS:false}" # Thread pool size for scheduler that executes device querying tasks scheduler_thread_pool_size: "${SNMP_SCHEDULER_THREAD_POOL_SIZE:4}" + # Maximum number of retry attempts for a single SNMP devices batch during bootstrap. + batch_retries: "${SNMP_BOOTSTRAP_RETRIES:8}" stats: # Enable/Disable the collection of transport statistics enabled: "${TB_TRANSPORT_STATS_ENABLED:true}" diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java index 56b4e00162..93ce3d1278 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java @@ -15,11 +15,14 @@ */ package org.thingsboard.server.transport.snmp; +import jakarta.annotation.PreDestroy; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.event.EventListener; import org.springframework.stereotype.Component; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; @@ -53,9 +56,12 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentLinkedDeque; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; @TbSnmpTransportComponent @Component @@ -72,30 +78,52 @@ public class SnmpTransportContext extends TransportContext { private final SnmpAuthService snmpAuthService; private final Map sessions = new ConcurrentHashMap<>(); - private final Collection allSnmpDevicesIds = new ConcurrentLinkedDeque<>(); + private final Set allSnmpDevicesIds = ConcurrentHashMap.newKeySet(); + private final ExecutorService snmpExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("snmp-bootstrap")); + + @Value("${transport.snmp.batch_retries}") + private int snmpBootstrapBatchRetries; @AfterStartUp(order = AfterStartUp.AFTER_TRANSPORT_SERVICE) public void fetchDevicesAndEstablishSessions() { - log.info("Initializing SNMP devices sessions"); + snmpExecutor.execute(this::bootstrapWithRetries); + } + private void bootstrapWithRetries() { + log.info("Initializing SNMP devices sessions"); int batchIndex = 0; int batchSize = 512; boolean nextBatchExists = true; while (nextBatchExists) { - TransportProtos.GetSnmpDevicesResponseMsg snmpDevicesResponse = protoEntityService.getSnmpDevicesIds(batchIndex, batchSize); - snmpDevicesResponse.getIdsList().stream() - .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++; + for (int attempt = 1; attempt <= snmpBootstrapBatchRetries; attempt++) { + try { + TransportProtos.GetSnmpDevicesResponseMsg snmpDevicesResponse = protoEntityService.getSnmpDevicesIds(batchIndex, batchSize); + snmpDevicesResponse.getIdsList().stream() + .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++; + break; + } catch (Exception e) { + if (attempt >= snmpBootstrapBatchRetries) { + log.error("SNMP bootstrap: batch {} failed after {} attempts.", batchIndex, attempt, e); + return; + } + log.warn("SNMP bootstrap: batch {} attempt {}/{} failed.", batchIndex, attempt, snmpBootstrapBatchRetries, e); + try { + TimeUnit.SECONDS.sleep(10); + } catch (InterruptedException ex) { + log.warn("SNMP bootstrap interrupted. Stopping bootstrap task."); + return; + } + } + } } - - log.debug("Found all SNMP devices ids: {}", allSnmpDevicesIds); + log.debug("Found SNMP devices ids: {}", allSnmpDevicesIds); } private void establishDeviceSession(Device device) { @@ -300,4 +328,9 @@ public class SnmpTransportContext extends TransportContext { return sessions.values(); } + @PreDestroy + public void destroy() { + snmpExecutor.shutdown(); + } + } diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 79aee31921..567654cce4 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -151,6 +151,8 @@ transport: ignore_type_cast_errors: "${SNMP_RESPONSE_IGNORE_TYPE_CAST_ERRORS:false}" # Thread pool size for scheduler that executes device querying tasks scheduler_thread_pool_size: "${SNMP_SCHEDULER_THREAD_POOL_SIZE:4}" + # Maximum number of retry attempts for a single SNMP devices batch during bootstrap. + batch_retries: "${SNMP_BOOTSTRAP_RETRIES:8}" sessions: # Session inactivity timeout is a global configuration parameter that defines how long the device transport session will be opened after the last message arrives from the device. # The parameter value is in milliseconds.