From 165148a1264e826821df3299f527abb00f17d69b Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 3 Jun 2022 12:18:47 +0300 Subject: [PATCH] Fixed and improved startup sequence --- .../server/actors/service/DefaultActorService.java | 4 ++-- .../service/queue/DefaultTbCoreConsumerService.java | 4 ++-- .../queue/processing/AbstractConsumerService.java | 4 ++-- .../transport/TbCoreTransportApiService.java | 4 ++-- .../queue/common/DefaultTbQueueRequestTemplate.java | 6 +++--- .../queue/discovery/DummyDiscoveryService.java | 4 ++-- .../queue/discovery/HashPartitionService.java | 8 +++++--- .../server/queue/discovery/ZkDiscoveryService.java | 4 ++-- .../queue/usagestats/DefaultTbApiUsageClient.java | 4 +--- .../thingsboard/server/queue/util/AfterStartUp.java | 13 ++++++++++++- .../lwm2m/server/DefaultLwM2mTransportService.java | 7 +------ .../lwm2m/server/client/LwM2mClientContextImpl.java | 2 +- .../server/model/LwM2MModelConfigServiceImpl.java | 2 +- .../server/transport/snmp/SnmpTransportContext.java | 2 +- .../transport/service/DefaultTransportService.java | 11 +++++------ .../service/TransportQueueRoutingInfoService.java | 9 ++++++--- 16 files changed, 48 insertions(+), 40 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 0809afb4b5..3d3bf2bbbb 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -35,6 +35,7 @@ import org.thingsboard.server.actors.stats.StatsActor; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; +import org.thingsboard.server.queue.util.AfterStartUp; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -113,8 +114,7 @@ public class DefaultActorService extends TbApplicationEventListener staleRequest, long currentNs) { if (currentNs >= staleRequest.getSubmitTime() + staleRequest.getTimeout()) { - log.warn("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); + log.debug("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } else { - log.error("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); + log.info("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } staleRequest.future.setException(new TimeoutException()); } @@ -173,7 +173,7 @@ public class DefaultTbQueueRequestTemplate expectedResponse = pendingRequests.remove(requestId); if (expectedResponse == null) { - log.warn("[{}] Invalid or stale request, response: {}", requestId, String.valueOf(response).replace("\n", " ")); + log.debug("[{}] Invalid or stale request, response: {}", requestId, String.valueOf(response).replace("\n", " ")); } else { expectedResponse.future.set(response); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java index 6718471abf..d124419135 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java @@ -22,6 +22,7 @@ import org.springframework.context.annotation.DependsOn; import org.springframework.context.event.EventListener; import org.springframework.core.annotation.Order; import org.springframework.stereotype.Service; +import org.thingsboard.server.queue.util.AfterStartUp; import java.util.Collections; @@ -40,8 +41,7 @@ public class DummyDiscoveryService implements DiscoveryService { this.partitionService = partitionService; } - @EventListener(ApplicationReadyEvent.class) - @Order(value = 1) + @AfterStartUp(order = AfterStartUp.DISCOVERY_SERVICE) public void onApplicationEvent(ApplicationReadyEvent event) { partitionService.recalculatePartitions(serviceInfoProvider.getServiceInfo(), Collections.emptyList()); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index e7db76bead..8a56a86bcc 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -31,6 +31,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo; import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent; +import org.thingsboard.server.queue.util.AfterStartUp; import javax.annotation.PostConstruct; import java.util.ArrayList; @@ -89,10 +90,10 @@ public class HashPartitionService implements PartitionService { @PostConstruct public void init() { this.hashFunction = forName(hashFunctionName); - partitionsInit(); } - private void partitionsInit() { + @AfterStartUp(order = AfterStartUp.QUEUE_INFO_INITIALIZATION) + public void partitionsInit() { QueueKey coreKey = new QueueKey(ServiceType.TB_CORE); partitionSizesMap.put(coreKey, corePartitions); partitionTopicsMap.put(coreKey, coreTopic); @@ -101,6 +102,7 @@ public class HashPartitionService implements PartitionService { String serviceType = serviceInfoProvider.getServiceType(); + if ("tb-transport".equals(serviceType)) { //If transport started earlier than tb-core int getQueuesRetries = 10; @@ -111,7 +113,7 @@ public class HashPartitionService implements PartitionService { queueRoutingInfoList = queueRoutingInfoService.getAllQueuesRoutingInfo(); break; } catch (Exception e) { - log.info("Failed to get queues routing info!"); + log.info("Failed to get queues routing info: {}!", e.getMessage()); getQueuesRetries--; } try { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java index 8afd4dc768..89eda16cad 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java @@ -41,6 +41,7 @@ import org.springframework.util.Assert; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent; +import org.thingsboard.server.queue.util.AfterStartUp; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -115,8 +116,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi .collect(Collectors.toList()); } - @EventListener(ApplicationReadyEvent.class) - @Order(value = 1) + @AfterStartUp(order = AfterStartUp.DISCOVERY_SERVICE) public void onApplicationEvent(ApplicationReadyEvent event) { if (stopped) { log.debug("Ignoring application ready event. Service is stopped."); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java index 60867f314f..11f01cba2b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java @@ -59,9 +59,7 @@ public class DefaultTbApiUsageClient implements TbApiUsageClient { private final EnumMap> stats = new EnumMap<>(ApiUsageRecordKey.class); - @Lazy - @Autowired - private PartitionService partitionService; + private final PartitionService partitionService; private final SchedulerComponent scheduler; private final TbQueueProducerProvider producerProvider; private TbQueueProducer> msgProducer; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java index 19e56e31cb..487e734c33 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java @@ -30,6 +30,17 @@ import java.lang.annotation.Target; @EventListener(ApplicationReadyEvent.class) @Order public @interface AfterStartUp { + + int QUEUE_INFO_INITIALIZATION = 1; + int DISCOVERY_SERVICE = 2; + + int ACTOR_SYSTEM = 9; + int REGULAR_SERVICE = 10; + + int BEFORE_TRANSPORT_SERVICE = Integer.MAX_VALUE - 1001; + int TRANSPORT_SERVICE = Integer.MAX_VALUE - 1000; + int AFTER_TRANSPORT_SERVICE = Integer.MAX_VALUE - 999; + @AliasFor(annotation = Order.class, attribute = "value") - int order() default Integer.MAX_VALUE; + int order(); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java index 9b49030f7e..d28ee5ee4c 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java @@ -20,15 +20,11 @@ import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.scandium.config.DtlsConfig; import org.eclipse.californium.scandium.config.DtlsConnectorConfig; import org.eclipse.californium.scandium.dtls.cipher.CipherSuite; -import org.eclipse.leshan.core.node.LwM2mNode; import org.eclipse.leshan.core.node.codec.DefaultLwM2mDecoder; import org.eclipse.leshan.core.node.codec.DefaultLwM2mEncoder; -import org.eclipse.leshan.core.request.SendRequest; import org.eclipse.leshan.server.californium.LeshanServer; import org.eclipse.leshan.server.californium.LeshanServerBuilder; import org.eclipse.leshan.server.californium.registration.CaliforniumRegistrationStore; -import org.eclipse.leshan.server.registration.Registration; -import org.eclipse.leshan.server.send.SendListener; import org.springframework.stereotype.Component; import org.thingsboard.server.cache.ota.OtaPackageDataCache; import org.thingsboard.server.common.data.DataConstants; @@ -44,7 +40,6 @@ import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PreDestroy; import java.security.cert.X509Certificate; -import java.util.Map; import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RECOMMENDED_CIPHER_SUITES_ONLY; import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RECOMMENDED_CURVES_ONLY; @@ -76,7 +71,7 @@ public class DefaultLwM2mTransportService implements LwM2MTransportService { private LeshanServer server; - @AfterStartUp(order = Integer.MAX_VALUE - 1) + @AfterStartUp(order = AfterStartUp.BEFORE_TRANSPORT_SERVICE) public void init() { this.server = getLhServer(); /* diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java index 260dcd7ab1..c25c2d3c9a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java @@ -84,7 +84,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { private final Map lwM2mClientsByRegistrationId = new ConcurrentHashMap<>(); private final Map profiles = new ConcurrentHashMap<>(); - @AfterStartUp(order = Integer.MAX_VALUE - 1) + @AfterStartUp(order = AfterStartUp.BEFORE_TRANSPORT_SERVICE) public void init() { String nodeId = context.getNodeId(); Set fetchedClients = clientStore.getAll(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java index 4432a0c83c..2634672a3b 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java @@ -69,7 +69,7 @@ public class LwM2MModelConfigServiceImpl implements LwM2MModelConfigService { private ConcurrentMap currentModelConfigs; - @AfterStartUp(order = Integer.MAX_VALUE - 1) + @AfterStartUp(order = AfterStartUp.BEFORE_TRANSPORT_SERVICE) private void init() { List models = modelStore.getAll(); log.debug("Fetched model configs: {}", models); 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 42190974b9..1dd46dc7cd 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 @@ -73,7 +73,7 @@ public class SnmpTransportContext extends TransportContext { private final Map sessions = new ConcurrentHashMap<>(); private final Collection allSnmpDevicesIds = new ConcurrentLinkedDeque<>(); - @AfterStartUp(order = Integer.MAX_VALUE) + @AfterStartUp(order = AfterStartUp.AFTER_TRANSPORT_SERVICE) public void fetchDevicesAndEstablishSessions() { log.info("Initializing SNMP devices sessions"); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 5b29d3c3ff..9f6ddc87a4 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -157,13 +157,10 @@ public class DefaultTransportService implements TransportService { @Autowired @Lazy private TbApiUsageClient apiUsageClient; - @Autowired - @Lazy - private PartitionService partitionService; - private final Map statsMap = new LinkedHashMap<>(); private final Gson gson = new Gson(); + private final PartitionService partitionService; private final TbTransportQueueFactory queueProvider; private final TbQueueProducerProvider producerProvider; @@ -197,7 +194,8 @@ public class DefaultTransportService implements TransportService { private volatile boolean stopped = false; - public DefaultTransportService(TbServiceInfoProvider serviceInfoProvider, + public DefaultTransportService(PartitionService partitionService, + TbServiceInfoProvider serviceInfoProvider, TbTransportQueueFactory queueProvider, TbQueueProducerProvider producerProvider, NotificationsTopicService notificationsTopicService, @@ -207,6 +205,7 @@ public class DefaultTransportService implements TransportService { TransportRateLimitService rateLimitService, DataDecodingEncodingService dataDecodingEncodingService, SchedulerComponent scheduler, TransportResourceCache transportResourceCache, ApplicationEventPublisher eventPublisher) { + this.partitionService = partitionService; this.serviceInfoProvider = serviceInfoProvider; this.queueProvider = queueProvider; this.producerProvider = producerProvider; @@ -240,7 +239,7 @@ public class DefaultTransportService implements TransportService { mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); } - @AfterStartUp + @AfterStartUp(order = AfterStartUp.TRANSPORT_SERVICE) private void start() { mainConsumerExecutor.execute(() -> { while (!stopped) { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportQueueRoutingInfoService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportQueueRoutingInfoService.java index ae1b6934c3..95fcef99f3 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportQueueRoutingInfoService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportQueueRoutingInfoService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.transport.service; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; @@ -33,9 +34,11 @@ import java.util.stream.Collectors; @ConditionalOnExpression("'${service.type:null}'=='tb-transport'") public class TransportQueueRoutingInfoService implements QueueRoutingInfoService { - @Lazy - @Autowired - private TransportService transportService; + private final TransportService transportService; + + public TransportQueueRoutingInfoService(@Lazy TransportService transportService) { + this.transportService = transportService; + } @Override public List getAllQueuesRoutingInfo() {