diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java index 44b025b5d7..95b9b821ec 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java @@ -15,8 +15,6 @@ */ package org.thingsboard.server.service.edge.rpc.session.manager; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -46,6 +44,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; @@ -68,6 +67,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan private final EdgeSessionsHolder sessions; private final AtomicReference> initMigrationAndProcessingFutureRef = new AtomicReference<>(); + private final AtomicReference> highPriorityProcessingFutureRef = new AtomicReference<>(); private final Lock initLock = new ReentrantLock(); private volatile boolean migrationProcessed; private volatile boolean isHighPriorityProcessing; @@ -99,6 +99,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan @Override public boolean destroy() { + cancelHighPriorityProcessing(); EdgeSessionState state = getState(); try { if (consumer != null) { @@ -123,7 +124,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan } private void scheduleMigrationAndProcessing() { - cancelScheduleEdgeEventsCheck(); + cancelMigrationAndProcessingInit(); EdgeSessionState state = getSession().getState(); TenantId tenantId = state.getTenantId(); EdgeId edgeId = state.getEdgeId(); @@ -139,7 +140,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan try { processRemainingPostgresEventsIfRequired(edgeId); if (migrationProcessed) { - initAndLaunchConsumerWithCallback(tenantId, edgeId, state); + initAndLaunchConsumerOrRetry(tenantId, edgeId, state); } } catch (Exception e) { log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, edgeId, e); @@ -152,19 +153,17 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan initMigrationAndProcessingFutureRef.set(initMigrationAndProcessingTask); } - private void initAndLaunchConsumerWithCallback(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) { - Futures.transform(initAndLaunchConsumer(tenantId, edgeId, state), success -> { - if (Boolean.FALSE.equals(success)) { - scheduleMigrationAndProcessing(); - } else { - launchProcessingOfHighPriorityEvents(tenantId, edgeId); - } - return null; - }, ctx.getGrpcCallbackExecutorService()); + private void initAndLaunchConsumerOrRetry(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) { + if (initAndLaunchConsumer(tenantId, edgeId, state)) { + launchProcessingOfHighPriorityEvents(tenantId, edgeId); + } else { + scheduleMigrationAndProcessing(); + } } private void launchProcessingOfHighPriorityEvents(TenantId tenantId, EdgeId edgeId) { - ctx.getEdgeEventProcessingExecutorService().scheduleAtFixedRate(() -> { + cancelHighPriorityProcessing(); + ScheduledFuture highPriorityProcessingTask = ctx.getEdgeEventProcessingExecutorService().scheduleAtFixedRate(() -> { try { isHighPriorityProcessing = true; session.processHighPriorityEvents(); @@ -173,41 +172,50 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, edgeId, e); } }, NO_INITIAL_DELAY_VALUE, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS); + highPriorityProcessingFutureRef.set(highPriorityProcessingTask); } - private ListenableFuture initAndLaunchConsumer(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) { + private boolean initAndLaunchConsumer(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) { if (!state.isConnected() || state.isSyncInProgress() || isHighPriorityProcessing) { log.warn("[{}][{}] Session is not ready (connected={}, syncInProgress={}, highPriority={}), skip starting edge event consumer", tenantId, edgeId, state.isConnected(), state.isSyncInProgress(), isHighPriorityProcessing); - return Futures.immediateFuture(Boolean.FALSE); + return false; } - if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) { - try { - if (consumerExecutor != null && !consumerExecutor.isShutdown()) { - try { - consumerExecutor.shutdown(); - awaitConsumerTermination(); - } catch (Exception e) { - log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edgeId, e); - } + if (isConsumerInitializationRequired()) { + initConsumerAndExecutor(tenantId, edgeId, state); + } + return true; + } + + private void initConsumerAndExecutor(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) { + try { + if (consumerExecutor != null && !consumerExecutor.isShutdown()) { + try { + consumerExecutor.shutdown(); + awaitConsumerTermination(); + } catch (Exception e) { + log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edgeId, e); } - this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); - this.consumer = QueueConsumerManager.>builder() - .name("TB Edge events [" + edgeId + "]") - .msgPackProcessor(this::processMsgs) - .pollInterval(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()) - .consumerCreator(() -> tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId)) - .consumerExecutor(consumerExecutor) - .threadPrefix("edge-events-" + edgeId) - .build(); - consumer.subscribe(); - consumer.launch(); - } catch (Exception e) { - destroy(); - log.warn("[{}][{}] Failed to start edge event consumer", state.getSessionId(), edgeId, e); } + this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); + this.consumer = QueueConsumerManager.>builder() + .name("TB Edge events [" + edgeId + "]") + .msgPackProcessor(this::processMsgs) + .pollInterval(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()) + .consumerCreator(() -> tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId)) + .consumerExecutor(consumerExecutor) + .threadPrefix("edge-events-" + edgeId) + .build(); + consumer.subscribe(); + consumer.launch(); + } catch (Exception e) { + destroy(); + log.warn("[{}][{}] Failed to start edge event consumer", state.getSessionId(), edgeId, e); } - return Futures.immediateFuture(Boolean.TRUE); + } + + private boolean isConsumerInitializationRequired() { + return consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped()); } private void processMsgs(List> msgs, TbQueueConsumer> consumer) { @@ -240,13 +248,20 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan } } - private void cancelScheduleEdgeEventsCheck() { - EdgeId edgeId = getState().getEdgeId(); - log.trace("[{}] cancelling edge migration & processing init for edge", edgeId); + private void cancelMigrationAndProcessingInit() { + log.trace("[{}] cancelling edge migration & processing init for edge", getState().getEdgeId()); + cancelFuture(initMigrationAndProcessingFutureRef); + } + + private void cancelHighPriorityProcessing() { + log.trace("[{}] cancelling high priority processing task for edge", getState().getEdgeId()); + cancelFuture(highPriorityProcessingFutureRef); + } - ScheduledFuture sf = initMigrationAndProcessingFutureRef.getAndSet(null); - if (sf != null && !sf.isCancelled() && !sf.isDone()) { - sf.cancel(true); + private void cancelFuture(AtomicReference> futureReference) { + Future f = futureReference.getAndSet(null); + if (f != null && !f.isCancelled() && !f.isDone()) { + f.cancel(true); } }