|
|
|
@ -18,26 +18,29 @@ package org.thingsboard.server.service.queue; |
|
|
|
import akka.actor.ActorRef; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
|
|
|
import org.springframework.boot.context.event.ApplicationReadyEvent; |
|
|
|
import org.springframework.context.event.EventListener; |
|
|
|
import org.springframework.scheduling.annotation.Scheduled; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
|
|
import org.thingsboard.rule.engine.api.RpcError; |
|
|
|
import org.thingsboard.server.actors.ActorSystemContext; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.msg.TbActorMsg; |
|
|
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
|
|
|
import org.thingsboard.server.common.msg.queue.TbMsgCallback; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.*; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCResponseProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.LocalSubscriptionServiceMsgProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionMgrMsgProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeUpdateProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionCloseProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TbTimeSeriesUpdateProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg; |
|
|
|
import org.thingsboard.server.queue.TbQueueConsumer; |
|
|
|
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|
|
|
import org.thingsboard.server.queue.discovery.PartitionChangeEvent; |
|
|
|
import org.thingsboard.server.queue.provider.TbCoreQueueProvider; |
|
|
|
import org.thingsboard.server.queue.util.TbMonolithOrCoreComponent; |
|
|
|
import org.thingsboard.server.service.encoding.DataDecodingEncodingService; |
|
|
|
import org.thingsboard.server.service.queue.processing.AbstractConsumerService; |
|
|
|
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; |
|
|
|
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; |
|
|
|
import org.thingsboard.server.service.state.DeviceStateService; |
|
|
|
@ -47,7 +50,6 @@ import org.thingsboard.server.service.subscription.TbSubscriptionUtils; |
|
|
|
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; |
|
|
|
|
|
|
|
import javax.annotation.PostConstruct; |
|
|
|
import javax.annotation.PreDestroy; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.UUID; |
|
|
|
@ -55,15 +57,14 @@ import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
|
import java.util.concurrent.CountDownLatch; |
|
|
|
import java.util.concurrent.ExecutorService; |
|
|
|
import java.util.concurrent.Executors; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.function.Function; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
@Service |
|
|
|
@ConditionalOnExpression("'${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core'") |
|
|
|
@TbMonolithOrCoreComponent |
|
|
|
@Slf4j |
|
|
|
public class DefaultTbCoreConsumerService implements TbCoreConsumerService { |
|
|
|
public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCoreMsg, ToCoreNotificationMsg> implements TbCoreConsumerService { |
|
|
|
|
|
|
|
@Value("${queue.core.poll_interval}") |
|
|
|
private long pollDuration; |
|
|
|
@ -72,56 +73,33 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { |
|
|
|
@Value("${queue.core.stats.enabled:false}") |
|
|
|
private boolean statsEnabled; |
|
|
|
|
|
|
|
private final ActorSystemContext actorContext; |
|
|
|
private final DeviceStateService stateService; |
|
|
|
private final TbLocalSubscriptionService localSubscriptionService; |
|
|
|
private final SubscriptionManagerService subscriptionManagerService; |
|
|
|
private final DataDecodingEncodingService encodingService; |
|
|
|
private final TbQueueConsumer<TbProtoQueueMsg<ToCoreMsg>> mainConsumer; |
|
|
|
private final TbQueueConsumer<TbProtoQueueMsg<ToCoreNotificationMsg>> nfConsumer; |
|
|
|
private final TbCoreDeviceRpcService tbCoreDeviceRpcService; |
|
|
|
private final TbCoreConsumerStats stats = new TbCoreConsumerStats(); |
|
|
|
private volatile ExecutorService mainConsumerExecutor; |
|
|
|
private volatile ExecutorService notificationsConsumerExecutor; |
|
|
|
private volatile boolean stopped = false; |
|
|
|
|
|
|
|
public DefaultTbCoreConsumerService(TbCoreQueueProvider tbCoreQueueProvider, ActorSystemContext actorContext, |
|
|
|
DeviceStateService stateService, TbLocalSubscriptionService localSubscriptionService, |
|
|
|
SubscriptionManagerService subscriptionManagerService, DataDecodingEncodingService encodingService, |
|
|
|
TbCoreDeviceRpcService tbCoreDeviceRpcService) { |
|
|
|
this.mainConsumer = tbCoreQueueProvider.getToCoreMsgConsumer(); |
|
|
|
this.nfConsumer = tbCoreQueueProvider.getToCoreNotificationsMsgConsumer(); |
|
|
|
this.actorContext = actorContext; |
|
|
|
super(actorContext, encodingService, |
|
|
|
tbCoreQueueProvider.getToCoreMsgConsumer(), tbCoreQueueProvider.getToCoreNotificationsMsgConsumer()); |
|
|
|
this.stateService = stateService; |
|
|
|
this.localSubscriptionService = localSubscriptionService; |
|
|
|
this.subscriptionManagerService = subscriptionManagerService; |
|
|
|
this.encodingService = encodingService; |
|
|
|
this.tbCoreDeviceRpcService = tbCoreDeviceRpcService; |
|
|
|
} |
|
|
|
|
|
|
|
@PostConstruct |
|
|
|
public void init() { |
|
|
|
this.mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-core-consumer")); |
|
|
|
this.notificationsConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-core-notifications-consumer")); |
|
|
|
super.init("tb-core-consumer", "tb-core-notifications-consumer"); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { |
|
|
|
if (partitionChangeEvent.getServiceKey().getServiceType() == ServiceType.TB_CORE) { |
|
|
|
log.info("Subscribing to partitions: {}", partitionChangeEvent.getPartitions()); |
|
|
|
this.mainConsumer.subscribe(partitionChangeEvent.getPartitions()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@EventListener(ApplicationReadyEvent.class) |
|
|
|
public void onApplicationEvent(ApplicationReadyEvent event) { |
|
|
|
log.info("Subscribing to notifications: {}", mainConsumer.getTopic()); |
|
|
|
this.nfConsumer.subscribe(); |
|
|
|
launchNotificationsConsumer(); |
|
|
|
launchMainConsumer(); |
|
|
|
} |
|
|
|
|
|
|
|
private void launchMainConsumer() { |
|
|
|
protected void launchMainConsumer() { |
|
|
|
log.info("Launching main consumer"); |
|
|
|
mainConsumerExecutor.execute(() -> { |
|
|
|
while (!stopped) { |
|
|
|
try { |
|
|
|
@ -134,6 +112,7 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { |
|
|
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToCoreMsg>> failedMap = new ConcurrentHashMap<>(); |
|
|
|
CountDownLatch processingTimeoutLatch = new CountDownLatch(1); |
|
|
|
pendingMap.forEach((id, msg) -> { |
|
|
|
log.info("[{}] Creating main callback for message: {}", id, msg.getValue()); |
|
|
|
TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); |
|
|
|
try { |
|
|
|
ToCoreMsg toCoreMsg = msg.getValue(); |
|
|
|
@ -177,57 +156,38 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private void launchNotificationsConsumer() { |
|
|
|
notificationsConsumerExecutor.execute(() -> { |
|
|
|
while (!stopped) { |
|
|
|
try { |
|
|
|
List<TbProtoQueueMsg<ToCoreNotificationMsg>> msgs = nfConsumer.poll(pollDuration); |
|
|
|
if (msgs.isEmpty()) { |
|
|
|
continue; |
|
|
|
} |
|
|
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToCoreNotificationMsg>> pendingMap = msgs.stream().collect( |
|
|
|
Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity())); |
|
|
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToCoreNotificationMsg>> failedMap = new ConcurrentHashMap<>(); |
|
|
|
CountDownLatch processingTimeoutLatch = new CountDownLatch(1); |
|
|
|
pendingMap.forEach((id, msg) -> { |
|
|
|
TbMsgCallback callback = new MsgPackCallback<>(id, processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>(), failedMap); |
|
|
|
try { |
|
|
|
ToCoreNotificationMsg toCoreMsg = msg.getValue(); |
|
|
|
if (toCoreMsg.hasToLocalSubscriptionServiceMsg()) { |
|
|
|
log.trace("[{}] Forwarding message to local subscription service {}", id, toCoreMsg.getToLocalSubscriptionServiceMsg()); |
|
|
|
forwardToLocalSubMgrService(toCoreMsg.getToLocalSubscriptionServiceMsg(), callback); |
|
|
|
} else if (toCoreMsg.hasFromDeviceRpcResponse()) { |
|
|
|
log.trace("[{}] Forwarding message to RPC service {}", id, toCoreMsg.getFromDeviceRpcResponse()); |
|
|
|
forwardToCoreRpcService(toCoreMsg.getFromDeviceRpcResponse(), callback); |
|
|
|
} else if (toCoreMsg.getComponentLifecycleMsg() != null && !toCoreMsg.getComponentLifecycleMsg().isEmpty()) { |
|
|
|
Optional<TbActorMsg> actorMsg = encodingService.decode(toCoreMsg.getComponentLifecycleMsg().toByteArray()); |
|
|
|
if (actorMsg.isPresent()) { |
|
|
|
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); |
|
|
|
actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); |
|
|
|
} |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
} catch (Throwable e) { |
|
|
|
log.warn("[{}] Failed to process notification: {}", id, msg, e); |
|
|
|
callback.onFailure(e); |
|
|
|
} |
|
|
|
}); |
|
|
|
if (!processingTimeoutLatch.await(packProcessingTimeout, TimeUnit.MILLISECONDS)) { |
|
|
|
pendingMap.forEach((id, msg) -> log.warn("[{}] Timeout to process notification: {}", id, msg.getValue())); |
|
|
|
failedMap.forEach((id, msg) -> log.warn("[{}] Failed to process notification: {}", id, msg.getValue())); |
|
|
|
} |
|
|
|
nfConsumer.commit(); |
|
|
|
} catch (Exception e) { |
|
|
|
log.warn("Failed to obtain notifications from queue.", e); |
|
|
|
try { |
|
|
|
Thread.sleep(pollDuration); |
|
|
|
} catch (InterruptedException e2) { |
|
|
|
log.trace("Failed to wait until the server has capacity to handle new notifications", e2); |
|
|
|
} |
|
|
|
} |
|
|
|
@Override |
|
|
|
protected ServiceType getServiceType() { |
|
|
|
return ServiceType.TB_CORE; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
protected long getNotificationPollDuration() { |
|
|
|
return pollDuration; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
protected long getNotificationPackProcessingTimeout() { |
|
|
|
return packProcessingTimeout; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
protected void handleNotification(UUID id, TbProtoQueueMsg<ToCoreNotificationMsg> msg, TbMsgCallback callback) throws Exception { |
|
|
|
ToCoreNotificationMsg toCoreMsg = msg.getValue(); |
|
|
|
if (toCoreMsg.hasToLocalSubscriptionServiceMsg()) { |
|
|
|
log.trace("[{}] Forwarding message to local subscription service {}", id, toCoreMsg.getToLocalSubscriptionServiceMsg()); |
|
|
|
forwardToLocalSubMgrService(toCoreMsg.getToLocalSubscriptionServiceMsg(), callback); |
|
|
|
} else if (toCoreMsg.hasFromDeviceRpcResponse()) { |
|
|
|
log.trace("[{}] Forwarding message to RPC service {}", id, toCoreMsg.getFromDeviceRpcResponse()); |
|
|
|
forwardToCoreRpcService(toCoreMsg.getFromDeviceRpcResponse(), callback); |
|
|
|
} else if (toCoreMsg.getComponentLifecycleMsg() != null && !toCoreMsg.getComponentLifecycleMsg().isEmpty()) { |
|
|
|
Optional<TbActorMsg> actorMsg = encodingService.decode(toCoreMsg.getComponentLifecycleMsg().toByteArray()); |
|
|
|
if (actorMsg.isPresent()) { |
|
|
|
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); |
|
|
|
actorContext.getAppActor().tell(actorMsg.get(), ActorRef.noSender()); |
|
|
|
} |
|
|
|
log.info("Tb Core Notifications Consumer stopped."); |
|
|
|
}); |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbMsgCallback callback) { |
|
|
|
@ -245,23 +205,6 @@ public class DefaultTbCoreConsumerService implements TbCoreConsumerService { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@PreDestroy |
|
|
|
public void destroy() { |
|
|
|
stopped = true; |
|
|
|
if (mainConsumer != null) { |
|
|
|
mainConsumer.unsubscribe(); |
|
|
|
} |
|
|
|
if (nfConsumer != null) { |
|
|
|
nfConsumer.unsubscribe(); |
|
|
|
} |
|
|
|
if (mainConsumerExecutor != null) { |
|
|
|
mainConsumerExecutor.shutdownNow(); |
|
|
|
} |
|
|
|
if (notificationsConsumerExecutor != null) { |
|
|
|
notificationsConsumerExecutor.shutdownNow(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void forwardToLocalSubMgrService(LocalSubscriptionServiceMsgProto msg, TbMsgCallback callback) { |
|
|
|
if (msg.hasSubUpdate()) { |
|
|
|
localSubscriptionService.onSubscriptionUpdate(msg.getSubUpdate().getSessionId(), TbSubscriptionUtils.fromProto(msg.getSubUpdate()), callback); |
|
|
|
|