Browse Source

Fixed and improved startup sequence

pull/6633/head
Andrii Shvaika 4 years ago
parent
commit
165148a126
  1. 4
      application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  4. 4
      application/src/main/java/org/thingsboard/server/service/transport/TbCoreTransportApiService.java
  5. 6
      common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
  6. 4
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java
  7. 8
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  8. 4
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
  9. 4
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java
  10. 13
      common/queue/src/main/java/org/thingsboard/server/queue/util/AfterStartUp.java
  11. 7
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java
  12. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
  13. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java
  14. 2
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
  15. 11
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  16. 9
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportQueueRoutingInfoService.java

4
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<PartitionCha
}
}
@EventListener(ApplicationReadyEvent.class)
@Order(value = 2)
@AfterStartUp(order = AfterStartUp.ACTOR_SYSTEM)
public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) {
log.info("Received application ready event. Sending application init message to actor system");
appActor.tellWithHighPriority(new AppInitMsg());

4
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -61,6 +61,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.edge.EdgeNotificationService;
@ -170,8 +171,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
}
}
@EventListener(ApplicationReadyEvent.class)
@Order(value = 2)
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
public void onApplicationEvent(ApplicationReadyEvent event) {
super.onApplicationEvent(event);
launchUsageStatsConsumer();

4
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -39,6 +39,7 @@ import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
@ -92,8 +93,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
this.notificationsConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName(nfConsumerThreadName));
}
@EventListener(ApplicationReadyEvent.class)
@Order(value = 2)
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
public void onApplicationEvent(ApplicationReadyEvent event) {
log.info("Subscribing to notifications: {}", nfConsumer.getTopic());
this.nfConsumer.subscribe();

4
application/src/main/java/org/thingsboard/server/service/transport/TbCoreTransportApiService.java

@ -33,6 +33,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.queue.util.TbCoreComponent;
import javax.annotation.PostConstruct;
@ -91,8 +92,7 @@ public class TbCoreTransportApiService {
transportApiTemplate = builder.build();
}
@EventListener(ApplicationReadyEvent.class)
@Order(value = 2)
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE)
public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) {
log.info("Received application ready event. Starting polling for events.");
transportApiTemplate.init(transportApiService);

6
common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java

@ -156,9 +156,9 @@ public class DefaultTbQueueRequestTemplate<Request extends TbQueueMsg, Response
void setTimeoutException(UUID key, ResponseMetaData<Response> 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<Request extends TbQueueMsg, Response
log.trace("[{}] Response received: {}", requestId, String.valueOf(response).replace("\n", " ")); //TODO remove overhead
ResponseMetaData<Response> 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);
}

4
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());
}

8
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 {

4
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.");

4
common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java

@ -59,9 +59,7 @@ public class DefaultTbApiUsageClient implements TbApiUsageClient {
private final EnumMap<ApiUsageRecordKey, ConcurrentMap<OwnerId, AtomicLong>> stats = new EnumMap<>(ApiUsageRecordKey.class);
@Lazy
@Autowired
private PartitionService partitionService;
private final PartitionService partitionService;
private final SchedulerComponent scheduler;
private final TbQueueProducerProvider producerProvider;
private TbQueueProducer<TbProtoQueueMsg<ToUsageStatsServiceMsg>> msgProducer;

13
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();
}

7
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();
/*

2
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<String, LwM2mClient> lwM2mClientsByRegistrationId = new ConcurrentHashMap<>();
private final Map<UUID, Lwm2mDeviceProfileTransportConfiguration> profiles = new ConcurrentHashMap<>();
@AfterStartUp(order = Integer.MAX_VALUE - 1)
@AfterStartUp(order = AfterStartUp.BEFORE_TRANSPORT_SERVICE)
public void init() {
String nodeId = context.getNodeId();
Set<LwM2mClient> fetchedClients = clientStore.getAll();

2
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<String, LwM2MModelConfig> currentModelConfigs;
@AfterStartUp(order = Integer.MAX_VALUE - 1)
@AfterStartUp(order = AfterStartUp.BEFORE_TRANSPORT_SERVICE)
private void init() {
List<LwM2MModelConfig> models = modelStore.getAll();
log.debug("Fetched model configs: {}", models);

2
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<DeviceId, DeviceSessionContext> sessions = new ConcurrentHashMap<>();
private final Collection<DeviceId> allSnmpDevicesIds = new ConcurrentLinkedDeque<>();
@AfterStartUp(order = Integer.MAX_VALUE)
@AfterStartUp(order = AfterStartUp.AFTER_TRANSPORT_SERVICE)
public void fetchDevicesAndEstablishSessions() {
log.info("Initializing SNMP devices sessions");

11
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<String, Number> 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) {

9
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<QueueRoutingInfo> getAllQueuesRoutingInfo() {

Loading…
Cancel
Save