Armin Felder 2 weeks ago
committed by GitHub
parent
commit
f8fe0c40ae
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 9
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java
  2. 27
      common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java
  3. 100
      common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java

9
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java

@ -99,6 +99,15 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
public boolean destroy() {
cancelHighPriorityProcessing();
EdgeSessionState state = getState();
// Unblock the consumer loop first: processMsgs() may be parked on the send-downlink
// future of a batch that the (now gone) edge will never acknowledge. Completing it as
// interrupted lets the loop observe the stop flag instead of stalling consumer.stop().
if (state.getSendDownlinkMsgsFuture() != null && !state.getSendDownlinkMsgsFuture().isDone()) {
state.getSendDownlinkMsgsFuture().set(true);
}
if (state.getScheduledSendDownlinkTask() != null) {
state.getScheduledSendDownlinkTask().cancel(true);
}
try {
if (consumer != null) {
log.info("[{}][{}] Stopping edge event consumer...", state.getTenantId(), state.getEdgeId());

27
common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java

@ -25,6 +25,7 @@ import org.thingsboard.server.queue.TbQueueMsg;
import java.util.List;
import java.util.Set;
import java.util.concurrent.CancellationException;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
@ -36,6 +37,8 @@ import java.util.function.Supplier;
@Slf4j
public class QueueConsumerManager<M extends TbQueueMsg> {
public static final long DEFAULT_STOP_TIMEOUT_MS = TimeUnit.SECONDS.toMillis(10);
private final String name;
private final MsgPackProcessor<M> msgPackProcessor;
private final long pollInterval;
@ -43,6 +46,7 @@ public class QueueConsumerManager<M extends TbQueueMsg> {
private final String threadPrefix;
/** Optional poll gate: while {@code false} the loop skips polling so the position doesn't advance; {@code null} = always ready (default). */
private final BooleanSupplier readinessCheck;
private final long stopTimeoutMs;
@Getter
private final TbQueueConsumer<M> consumer;
@ -52,13 +56,15 @@ public class QueueConsumerManager<M extends TbQueueMsg> {
@Builder
public QueueConsumerManager(String name, MsgPackProcessor<M> msgPackProcessor,
long pollInterval, Supplier<TbQueueConsumer<M>> consumerCreator,
ExecutorService consumerExecutor, String threadPrefix, BooleanSupplier readinessCheck) {
ExecutorService consumerExecutor, String threadPrefix,
BooleanSupplier readinessCheck, Long stopTimeoutMs) {
this.name = name;
this.pollInterval = pollInterval;
this.msgPackProcessor = msgPackProcessor;
this.consumerExecutor = consumerExecutor;
this.threadPrefix = threadPrefix;
this.readinessCheck = readinessCheck;
this.stopTimeoutMs = stopTimeoutMs != null ? stopTimeoutMs : DEFAULT_STOP_TIMEOUT_MS;
this.consumer = consumerCreator.get();
}
@ -135,11 +141,22 @@ public class QueueConsumerManager<M extends TbQueueMsg> {
log.debug("[{}] Stopping consumer", name);
stopped = true;
consumer.unsubscribe();
if (consumerTask == null) {
return;
}
try {
if (consumerTask != null) {
consumerTask.get(10, TimeUnit.SECONDS);
}
} catch (InterruptedException | ExecutionException | TimeoutException e) {
consumerTask.get(stopTimeoutMs, TimeUnit.MILLISECONDS);
} catch (TimeoutException e) {
// The loop did not finish in time: it is most likely blocked inside the message pack
// processor (e.g. waiting for a downlink delivery that will never be acknowledged).
// A leftover consumer keeps its group membership alive via heartbeats and thereby
// blocks a replacement consumer of the same group from getting partitions assigned,
// so we must interrupt it instead of abandoning it.
log.warn("[{}] Consumer loop did not stop in {} ms, interrupting it", name, stopTimeoutMs);
consumerTask.cancel(true);
} catch (CancellationException e) {
log.debug("[{}] Consumer loop was already cancelled", name);
} catch (InterruptedException | ExecutionException e) {
log.error("[{}] Failed to await consumer loop stop", name, e);
}
}

100
common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java

@ -31,6 +31,7 @@ import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.BooleanSupplier;
@ -145,6 +146,98 @@ class QueueConsumerManagerTest {
.isTrue();
}
// --- stop() behaviour ---------------------------------------------------
// Stopping must work both while the loop is idle and while it is blocked inside the
// msg pack processor. A stop that cannot terminate the loop leaves a zombie consumer
// behind, which (for kafka) keeps its group membership alive via heartbeats and blocks
// the next consumer of the same group from getting partitions assigned (see edge
// downlink starvation after session handover).
private static final long SHORT_STOP_TIMEOUT_MS = 200;
@Test
void stopShouldCompleteGracefullyWhenConsumerLoopIsIdle() {
manager = launchManager(consumer, null, (msgs, c) -> c.commit());
await().atMost(5, TimeUnit.SECONDS)
.untilAsserted(() -> assertThat(consumer.getPollCount()).isPositive());
manager.stop();
assertThat(consumerLoopHasExited()).as("consumer loop thread released").isTrue();
}
@Test
void stopShouldInterruptConsumerLoopThatIsBlockedInProcessing() throws Exception {
CountDownLatch processingStarted = new CountDownLatch(1);
CountDownLatch neverReleased = new CountDownLatch(1);
consumer.enqueue(List.of(mock(TbQueueMsg.class)));
manager = launchManager(consumer, null, (msgs, c) -> {
processingStarted.countDown();
// simulates waiting for a downlink delivery that will never be acknowledged
neverReleased.await();
}, SHORT_STOP_TIMEOUT_MS);
assertThat(processingStarted.await(5, TimeUnit.SECONDS))
.as("processor entered its blocking section").isTrue();
long stopStarted = System.currentTimeMillis();
manager.stop();
assertThat(System.currentTimeMillis() - stopStarted)
.as("stop() must not hang much longer than the configured stop timeout")
.isLessThan(SHORT_STOP_TIMEOUT_MS + 3000);
assertThat(consumerLoopHasExited())
.as("blocked consumer loop must be interrupted, not abandoned").isTrue();
}
@Test
void stopShouldReturnImmediatelyWhenConsumerWasNeverLaunched() {
consumerExecutor = Executors.newSingleThreadExecutor();
manager = QueueConsumerManager.<TbQueueMsg>builder()
.name("test-consumer")
.pollInterval(POLL_INTERVAL_MS)
.consumerCreator(() -> consumer)
.consumerExecutor(consumerExecutor)
.msgPackProcessor((msgs, c) -> c.commit())
.build();
manager.stop();
assertThat(consumer.isStopped()).isTrue();
}
@Test
void consumerLoopShouldSurviveProcessorFailuresAndKeepPolling() {
AtomicInteger processedBatches = new AtomicInteger();
consumer.enqueue(List.of(mock(TbQueueMsg.class)));
consumer.enqueue(List.of(mock(TbQueueMsg.class)));
manager = launchManager(consumer, null, (msgs, c) -> {
if (processedBatches.incrementAndGet() == 1) {
throw new RuntimeException("transient failure on the first batch");
}
c.commit();
});
await().atMost(5, TimeUnit.SECONDS).untilAsserted(() ->
assertThat(processedBatches.get())
.as("loop keeps processing after a processor failure")
.isGreaterThanOrEqualTo(2));
}
/**
* The single-threaded consumer executor can only run this probe task once the
* consumer loop has actually finished - a hung loop keeps the thread forever.
*/
private boolean consumerLoopHasExited() {
try {
consumerExecutor.submit(() -> {
}).get(5, TimeUnit.SECONDS);
return true;
} catch (Exception e) {
return false;
}
}
private static void awaitReadinessGateEvaluated(AtomicInteger readinessChecks) {
await().atMost(5, TimeUnit.SECONDS)
.untilAsserted(() -> assertThat(readinessChecks.get())
@ -161,6 +254,12 @@ class QueueConsumerManagerTest {
private QueueConsumerManager<TbQueueMsg> launchManager(TestQueueConsumer consumer, BooleanSupplier readinessCheck,
QueueConsumerManager.MsgPackProcessor<TbQueueMsg> processor) {
return launchManager(consumer, readinessCheck, processor, null);
}
private QueueConsumerManager<TbQueueMsg> launchManager(TestQueueConsumer consumer, BooleanSupplier readinessCheck,
QueueConsumerManager.MsgPackProcessor<TbQueueMsg> processor,
Long stopTimeoutMs) {
consumerExecutor = Executors.newSingleThreadExecutor();
QueueConsumerManager<TbQueueMsg> queueConsumerManager = QueueConsumerManager.<TbQueueMsg>builder()
.name("test-consumer")
@ -168,6 +267,7 @@ class QueueConsumerManagerTest {
.consumerCreator(() -> consumer)
.consumerExecutor(consumerExecutor)
.readinessCheck(readinessCheck)
.stopTimeoutMs(stopTimeoutMs)
.msgPackProcessor(processor)
.build();
queueConsumerManager.subscribe();

Loading…
Cancel
Save