From 644ccbe82ef85f928f8d276f502ecf1177435ab9 Mon Sep 17 00:00:00 2001 From: Alexandru Stefanica Date: Wed, 16 Sep 2026 19:12:43 +0100 Subject: [PATCH] [FLINK-40657][connector-base] Release per-fetcher state in FutureCompletingBlockingQueue Fetcher indexes are never recycled, so the array of ConditionAndFlag grew without bound. Key the state by producer index instead and release it from the SplitFetcher shutdown hook. Generated-by: Claude Code (Claude Opus 5) --- .../reader/fetcher/SplitFetcherManager.java | 1 + .../FutureCompletingBlockingQueue.java | 142 +++++++++++-- .../fetcher/SplitFetcherManagerTest.java | 184 ++++++++++++++++ .../FutureCompletingBlockingQueueTest.java | 197 ++++++++++++++++++ .../reader/synchronization/QueueProbe.java | 29 +++ 5 files changed, 530 insertions(+), 23 deletions(-) create mode 100644 flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/QueueProbe.java diff --git a/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java b/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java index 6ea31d2d53baed..2f5fbe411553bc 100644 --- a/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java +++ b/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java @@ -260,6 +260,7 @@ protected synchronized SplitFetcher createSplitFetcher() { errorHandler, () -> { fetchers.remove(fetcherId); + elementsQueue.releaseProducer(fetcherId); fetchersToShutDown.decrementAndGet(); // We need this to synchronize status of fetchers to concurrent partners // as diff --git a/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java b/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java index 799fd0c4489b95..d023b4985ef6fd 100644 --- a/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java +++ b/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java @@ -27,7 +27,8 @@ import java.lang.reflect.Field; import java.util.ArrayDeque; -import java.util.Arrays; +import java.util.HashMap; +import java.util.Map; import java.util.Queue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; @@ -37,6 +38,7 @@ import java.util.concurrent.locks.ReentrantLock; import static org.apache.flink.util.Preconditions.checkArgument; +import static org.apache.flink.util.Preconditions.checkState; /** * A custom implementation of blocking queue in combination with a {@link CompletableFuture} that is @@ -69,6 +71,11 @@ * capacity limits, without interrupting the thread. This is done via the {@link * #wakeUpPuttingThread(int)} method. * + *

The queue keeps a small amount of per-producer state to support this. A producer that is + * permanently done with the queue must be handed to {@link #releaseProducer(int)}, the lifecycle + * counterpart of {@link #wakeUpPuttingThread(int)}, otherwise that state is retained for the + * lifetime of the queue. + * * @param the type of the elements in the queue. */ @Internal @@ -102,9 +109,15 @@ public class FutureCompletingBlockingQueue { @GuardedBy("lock") private final Queue notFull; - /** The per-thread conditions and wakeUp flags. */ + /** + * The per-producer conditions and wakeUp flags, keyed by the producer's thread index. + * + *

Entries are created on demand by {@link #conditionAndFlagFor(int)} and removed by {@link + * #releaseProducer(int)}, so this map is bounded by the peak number of live or in-flight + * producers rather than by the largest producer index ever allocated. + */ @GuardedBy("lock") - private ConditionAndFlag[] putConditionAndFlags; + private final Map putConditionAndFlags; public FutureCompletingBlockingQueue() { this(SourceReaderOptions.ELEMENT_QUEUE_CAPACITY.defaultValue()); @@ -115,7 +128,7 @@ public FutureCompletingBlockingQueue(int capacity) { this.capacity = capacity; this.queue = new ArrayDeque<>(capacity); this.lock = new ReentrantLock(); - this.putConditionAndFlags = new ConditionAndFlag[1]; + this.putConditionAndFlags = new HashMap<>(); this.notFull = new ArrayDeque<>(); // initially the queue is empty and thus unavailable @@ -317,6 +330,28 @@ int getNumberOfQueuedPutters() { } } + int numberOfProducerStates() { + lock.lock(); + try { + return putConditionAndFlags.size(); + } finally { + lock.unlock(); + } + } + + boolean hasProducerState(int producerIndex) { + lock.lock(); + try { + return putConditionAndFlags.containsKey(producerIndex); + } finally { + lock.unlock(); + } + } + + Lock lock() { + return lock; + } + /** * Gracefully wakes up the thread with the given {@code threadIndex} if it is blocked in adding * an element. to the queue. If the thread is blocked in {@link #put(int, Object)} it will @@ -330,12 +365,55 @@ int getNumberOfQueuedPutters() { public void wakeUpPuttingThread(int threadIndex) { lock.lock(); try { - maybeCreateCondition(threadIndex); - ConditionAndFlag caf = putConditionAndFlags[threadIndex]; - if (caf != null) { - caf.setWakeUp(true); - caf.condition().signal(); + // Creates the entry when absent, deliberately: the flag has to be sticky, so that a + // producer woken before it ever calls put() still observes the request and returns + // immediately instead of parking on a full queue. + final ConditionAndFlag caf = conditionAndFlagFor(threadIndex); + caf.setWakeUp(true); + caf.condition().signal(); + } finally { + lock.unlock(); + } + } + + /** + * Releases the per-producer wakeup state held for {@code threadIndex}. + * + *

Call this once the producer with that index is permanently finished with this queue. + * Without it the queue retains each condition created by a wakeup request or by a put attempt + * made while the queue was full. {@code SplitFetcherManager} hands out a fresh, never-recycled + * index per {@code SplitFetcher}, so for a source whose fetchers are short-lived that set + * otherwise grows without bound for the lifetime of the JVM. + * + *

The call is idempotent and safe for an index that was never used. + * + *

The caller must not be inside {@link #put(int, Object)} for this index. Releasing a + * producer that is still running would drop a pending wakeUp flag it has yet to observe, and it + * could then park on a full queue after having been told to stop. {@code SplitFetcher} + * satisfies this by running its shutdown hook only after its run loop has exited. + * + *

For the same reason this must be the last interaction with the queue for that index: a + * later {@link #wakeUpPuttingThread(int)} would recreate the entry, and nothing would remove it + * again. {@code SplitFetcher} satisfies this in its normal lifecycle because it clears the + * current task before its run loop exits and invokes the shutdown hook afterward. + * + * @param threadIndex The number identifying the producer thread, as passed to {@link #put(int, + * Object)}. + */ + public void releaseProducer(int threadIndex) { + lock.lock(); + try { + final ConditionAndFlag caf = putConditionAndFlags.get(threadIndex); + if (caf == null) { + return; + } + if (caf.hasWaitingPutter()) { + // The condition may already have been removed from notFull after being signalled, + // while its putter is still waiting to reacquire the lock. Dropping the entry in + // that interval would discard a wakeUp flag the putter has yet to observe. + return; } + putConditionAndFlags.remove(threadIndex); } finally { lock.unlock(); } @@ -370,15 +448,17 @@ private T dequeue() { @GuardedBy("lock") private void waitOnPut(int fetcherIndex) throws InterruptedException { - maybeCreateCondition(fetcherIndex); - Condition cond = putConditionAndFlags[fetcherIndex].condition(); - notFull.add(cond); + final ConditionAndFlag caf = conditionAndFlagFor(fetcherIndex); + final Condition cond = caf.condition(); + caf.startWaiting(); try { + notFull.add(cond); cond.await(); } finally { // drop the condition once the thread stops waiting, so a later signalNextPutter() // does not signal a putter that is no longer waiting notFull.remove(cond); + caf.stopWaiting(); } } @@ -389,22 +469,24 @@ private void signalNextPutter() { } } + /** + * Returns the state for {@code threadIndex}, creating it when absent. Only {@link + * #wakeUpPuttingThread(int)} and {@link #waitOnPut(int)} call this, so a producer that is never + * woken and never blocks costs nothing. + */ @GuardedBy("lock") - private void maybeCreateCondition(int threadIndex) { - if (putConditionAndFlags.length < threadIndex + 1) { - putConditionAndFlags = Arrays.copyOf(putConditionAndFlags, threadIndex + 1); - } - - if (putConditionAndFlags[threadIndex] == null) { - putConditionAndFlags[threadIndex] = new ConditionAndFlag(lock.newCondition()); - } + private ConditionAndFlag conditionAndFlagFor(int threadIndex) { + return putConditionAndFlags.computeIfAbsent( + threadIndex, ignored -> new ConditionAndFlag(lock.newCondition())); } @GuardedBy("lock") private boolean getAndResetWakeUpFlag(int threadIndex) { - maybeCreateCondition(threadIndex); - if (putConditionAndFlags[threadIndex].getWakeUp()) { - putConditionAndFlags[threadIndex].setWakeUp(false); + // Deliberately does not create the state: an absent entry cannot carry a wakeUp flag, and + // waitOnPut() creates it a moment later if this producer goes on to block. + final ConditionAndFlag caf = putConditionAndFlags.get(threadIndex); + if (caf != null && caf.getWakeUp()) { + caf.setWakeUp(false); return true; } return false; @@ -415,6 +497,7 @@ private boolean getAndResetWakeUpFlag(int threadIndex) { private static class ConditionAndFlag { private final Condition cond; private boolean wakeUp; + private int waitingPutters; private ConditionAndFlag(Condition cond) { this.cond = cond; @@ -429,6 +512,19 @@ private boolean getWakeUp() { return wakeUp; } + private void startWaiting() { + waitingPutters++; + } + + private void stopWaiting() { + checkState(waitingPutters > 0, "stopWaiting() without a matching startWaiting()"); + waitingPutters--; + } + + private boolean hasWaitingPutter() { + return waitingPutters > 0; + } + private void setWakeUp(boolean value) { wakeUp = value; } diff --git a/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManagerTest.java b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManagerTest.java index cf0fce5d42cf67..bcf8bb934de6e3 100644 --- a/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManagerTest.java +++ b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManagerTest.java @@ -31,7 +31,9 @@ import org.apache.flink.connector.base.source.reader.mocks.TestingSplitReader; import org.apache.flink.connector.base.source.reader.splitreader.SplitReader; import org.apache.flink.connector.base.source.reader.splitreader.SplitsChange; +import org.apache.flink.connector.base.source.reader.splitreader.SplitsRemoval; import org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue; +import org.apache.flink.connector.base.source.reader.synchronization.QueueProbe; import org.apache.flink.core.testutils.OneShotLatch; import org.apache.flink.util.MdcUtils; @@ -140,6 +142,80 @@ void testCloseCleansUpPreviouslyClosedFetcher() throws Exception { fetcherManager.close(Long.MAX_VALUE); } + /** + * Each completed fetcher lifecycle must release its queue state, so monotonically increasing + * fetcher ids do not accumulate historical state. + */ + @Test + void testFetcherShutdownReleasesWakeupStateAcrossLifecycles() throws Exception { + final String splitId = "testSplit"; + final Configuration config = new Configuration(); + config.set(SourceReaderOptions.ELEMENT_QUEUE_CAPACITY, 1); + + final SplitFetcherManager fetcherManager = + new SingleThreadFetcherManager<>( + () -> + new AwaitingReader<>( + new IOException("Should not happen"), + new RecordsBySplits<>( + Collections.emptyMap(), + Collections.singleton(splitId))), + config); + final FutureCompletingBlockingQueue> queue = + fetcherManager.getQueue(); + + try { + for (int expectedFetcherId = 0; expectedFetcherId < 3; expectedFetcherId++) { + fetcherManager.addSplits( + Collections.singletonList(new TestingSourceSplit(splitId))); + assertThat(fetcherManager.fetchers).hasSize(1); + assertThat(fetcherManager.fetchers.keySet().iterator().next()) + .isEqualTo(expectedFetcherId); + + waitUntil( + () -> queue.size() == 1, + Duration.ofSeconds(10), + "The data batch should have filled the element queue."); + waitUntil( + () -> { + fetcherManager.maybeShutdownFinishedFetchers(); + return fetcherManager.fetchers.isEmpty(); + }, + Duration.ofSeconds(10), + "The idle fetcher should have been removed."); + waitUntil( + () -> QueueProbe.liveProducerStates(queue) == 1, + Duration.ofSeconds(10), + "The final synchronization batch should have registered producer state while waiting for queue capacity."); + + assertThat(QueueProbe.liveProducerStates(queue)).isOne(); + + final RecordsWithSplitIds dataBatch = queue.poll(); + assertThat(dataBatch).isNotNull(); + dataBatch.recycle(); + + waitUntil( + () -> queue.size() == 1, + Duration.ofSeconds(10), + "The final synchronization batch should have been enqueued."); + final RecordsWithSplitIds synchronizationBatch = queue.poll(); + assertThat(synchronizationBatch).isNotNull(); + synchronizationBatch.recycle(); + + waitUntil( + () -> QueueProbe.liveProducerStates(queue) == 0, + Duration.ofSeconds(10), + "The shutdown hook should release the fetcher's queue state."); + } + } finally { + RecordsWithSplitIds batch; + while ((batch = queue.poll()) != null) { + batch.recycle(); + } + fetcherManager.close(10_000L); + } + } + /** * This test is somewhat testing the implementation instead of contract. This is because the * test is trying to make sure the element queue draining thread is not tight looping. @@ -266,6 +342,71 @@ private final void testExceptionPropagation( } } + /** + * A fetcher whose current task failed never enqueues the shutdown synchronization batch. The + * failed task remains published as its running task, keeping the fetcher non-idle and reachable + * by {@link SplitFetcherManager#close(long)}, which releases its {@code recordsProcessedLatch}. + * Otherwise the thread would remain stranded with its {@link SplitReader} unclosed and its + * queue state unreleased. + */ + @Test + void testFailedFetcherIsNotReapedAndIsCleanedUpOnClose() throws Exception { + final String splitId = "testSplit"; + final FailOnRemoveReader reader = new FailOnRemoveReader<>(); + final SingleThreadFetcherManager fetcherManager = + new SingleThreadFetcherManager<>(() -> reader, new Configuration()); + + boolean closed = false; + try { + fetcherManager.addSplits(Collections.singletonList(new TestingSourceSplit(splitId))); + + // Drive the fetcher into the failure path with nothing left assigned: RemoveSplitsTask + // drops the split from assignedSplits before handing the change to the reader, which + // throws. assignedSplits and taskQueue are then empty, but the failed task remains + // published as runningTask and prevents the fetcher from being reported idle. + fetcherManager.removeSplits(Collections.singletonList(new TestingSourceSplit(splitId))); + + // Wait for the failure to reach the error handler. There is no window to hit here: the + // state under test is the fetcher's terminal state, not a transient one. + waitUntil( + () -> { + try { + fetcherManager.checkErrors(); + return false; + } catch (RuntimeException expected) { + return true; + } + }, + Duration.ofSeconds(30), + "The fetcher failure should have been reported to the error handler."); + + fetcherManager.maybeShutdownFinishedFetchers(); + assertThat(fetcherManager.getNumAliveFetchers()) + .as("A failed fetcher must not be reaped as idle.") + .isOne(); + + final FutureCompletingBlockingQueue> queue = + fetcherManager.getQueue(); + final int fetcherId = fetcherManager.fetchers.keySet().iterator().next(); + queue.wakeUpPuttingThread(fetcherId); + assertThat(QueueProbe.liveProducerStates(queue)) + .as( + "The failed fetcher should have queue state for its shutdown hook to release.") + .isOne(); + + fetcherManager.close(30_000L); + closed = true; + assertThat(reader.isClosed()).as("The split reader should have been closed.").isTrue(); + assertThat(QueueProbe.liveProducerStates(queue)) + .as("The shutdown hook should have released the fetcher's queue state.") + .isZero(); + } finally { + if (!closed) { + fetcherManager.close(30_000L); + } + } + } + // ------------------------------------------------------------------------ // test helpers // ------------------------------------------------------------------------ @@ -370,6 +511,49 @@ public void wakeUp() { public void close() {} } + /** + * Accepts a split assignment but throws when asked to remove one, mirroring the many {@link + * SplitReader} implementations that do not support removal. {@code RemoveSplitsTask} empties + * {@code assignedSplits} before delegating here, so the fetcher fails with nothing assigned. + */ + private static final class FailOnRemoveReader implements SplitReader { + + private final OneShotLatch fetchBlocker = new OneShotLatch(); + private volatile boolean closed; + + @Override + public RecordsWithSplitIds fetch() { + // Stay inside fetch() until woken up, so the fetcher does not spin on empty fetches. + try { + fetchBlocker.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + return new RecordsBySplits<>(Collections.emptyMap(), Collections.emptySet()); + } + + @Override + public void handleSplitsChanges(SplitsChange splitsChanges) { + if (splitsChanges instanceof SplitsRemoval) { + throw new UnsupportedOperationException("Split removal is not supported."); + } + } + + @Override + public void wakeUp() { + fetchBlocker.trigger(); + } + + @Override + public void close() { + closed = true; + } + + boolean isClosed() { + return closed; + } + } + private static final class AwaitingReader implements SplitReader { diff --git a/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java index 65e7ff38b92395..9804f3a6667e97 100644 --- a/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java +++ b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java @@ -33,6 +33,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.Lock; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; @@ -147,6 +148,202 @@ void testWakeUpDoesNotStrandAnotherPutter() throws Exception { .isTrue(); } + /** + * Without {@link FutureCompletingBlockingQueue#releaseProducer(int)} the queue keeps one + * condition per producer index it has ever seen. Because {@code SplitFetcherManager} allocates + * a fresh, never-recycled index per {@code SplitFetcher}, a source with short-lived fetchers + * accumulates them for the lifetime of the JVM. + */ + @Test + void testReleaseProducerBoundsTheWakeupState() throws InterruptedException { + final FutureCompletingBlockingQueue queue = new FutureCompletingBlockingQueue<>(1); + + for (int producer = 0; producer < 3; producer++) { + // Fill the single slot, so the next put takes the full-queue path that registers + // wakeup state for this producer. + assertThat(queue.put(producer, producer)).isTrue(); + queue.wakeUpPuttingThread(producer); + assertThat(queue.put(producer, producer)).isFalse(); + queue.poll(); + + assertThat(queue.numberOfProducerStates()).isOne(); + + queue.releaseProducer(producer); + + assertThat(queue.numberOfProducerStates()).isZero(); + } + + assertThat(queue.getNumberOfQueuedPutters()).isZero(); + } + + /** Release is idempotent, tolerates unknown ids, and only removes the target state. */ + @Test + void testReleaseProducerOnlyRemovesTheTargetState() { + final FutureCompletingBlockingQueue queue = new FutureCompletingBlockingQueue<>(1); + + queue.wakeUpPuttingThread(3); + queue.wakeUpPuttingThread(4); + + queue.releaseProducer(7); + queue.releaseProducer(3); + queue.releaseProducer(3); + + assertThat(queue.hasProducerState(3)).isFalse(); + assertThat(queue.hasProducerState(4)).isTrue(); + assertThat(queue.numberOfProducerStates()).isOne(); + + queue.releaseProducer(4); + assertThat(queue.numberOfProducerStates()).isZero(); + } + + /** + * Releasing a producer that is currently parked in {@code waitOnPut} would discard the wakeUp + * flag it is about to read and could leave it parked for good, so the release must be refused + * while it is still waiting. + */ + @Test + void testReleaseProducerIsRefusedWhileTheProducerIsParked() throws Exception { + final FutureCompletingBlockingQueue queue = new FutureCompletingBlockingQueue<>(1); + + queue.put(0, 0); + + final AtomicBoolean parkedPutterResult = new AtomicBoolean(true); + final Thread parkedPutter = + new Thread(() -> parkedPutterResult.set(putUnchecked(queue, 1, 1)), "parkedPutter"); + parkedPutter.start(); + try { + CommonTestUtils.waitUntilCondition(() -> queue.getNumberOfQueuedPutters() == 1); + + queue.releaseProducer(1); + assertThat(queue.numberOfProducerStates()) + .as("must not drop wakeup state for a producer currently inside put()") + .isOne(); + + // The graceful wakeup still reaches it, which is what the refusal protects. + queue.wakeUpPuttingThread(1); + joinWithinTimeout(parkedPutter); + assertThat(parkedPutterResult).isFalse(); + + // Once it has left put(), the release goes through. + queue.releaseProducer(1); + assertThat(queue.numberOfProducerStates()).isZero(); + } finally { + if (parkedPutter.isAlive()) { + queue.wakeUpPuttingThread(1); + parkedPutter.interrupt(); + queue.poll(); + parkedPutter.join(TimeUnit.SECONDS.toMillis(10)); + } + queue.releaseProducer(1); + } + } + + /** An interrupted wait must not leave the producer permanently marked as waiting. */ + @Test + void testInterruptedPutterDoesNotPreventProducerRelease() throws Exception { + final FutureCompletingBlockingQueue queue = new FutureCompletingBlockingQueue<>(1); + assertThat(queue.put(0, 0)).isTrue(); + + final CompletableFuture putInterrupted = new CompletableFuture<>(); + final Thread putter = + new Thread( + () -> { + try { + queue.put(1, 1); + putInterrupted.complete(false); + } catch (InterruptedException expected) { + putInterrupted.complete(true); + } catch (Throwable failure) { + putInterrupted.completeExceptionally(failure); + } + }, + "interruptedPutter"); + putter.start(); + try { + CommonTestUtils.waitUntilCondition(() -> queue.getNumberOfQueuedPutters() == 1); + + putter.interrupt(); + assertThat(putInterrupted.get(10, TimeUnit.SECONDS)).isTrue(); + joinWithinTimeout(putter); + assertThat(queue.getNumberOfQueuedPutters()).isZero(); + assertThat(queue.hasProducerState(1)).isTrue(); + + queue.releaseProducer(1); + assertThat(queue.hasProducerState(1)).isFalse(); + } finally { + if (putter.isAlive()) { + putter.interrupt(); + queue.poll(); + putter.join(TimeUnit.SECONDS.toMillis(10)); + } + queue.releaseProducer(1); + } + } + + /** + * The companion to the test above, for the window that makes membership of {@code notFull} an + * unreliable answer to "is this producer still inside {@code put()}?". + * + *

{@code signalNextPutter()} removes a condition from {@code notFull} at signal time, not + * when its producer resumes, so between the signal and that producer reacquiring the lock it is + * still inside {@code put()} while absent from {@code notFull}. A release in that window would + * drop the state together with a wakeUp flag the producer has not read yet, and the producer + * would then build fresh state with no flag and park again. The {@code waitingPutters} counter + * spans the whole of {@code cond.await()} and therefore covers it. + * + *

The window is driven deterministically rather than raced for: the test thread takes the + * queue's own lock, so the signalled producer cannot reacquire it and cannot leave {@code + * put()} until the test releases it. + */ + @Test + void testReleaseProducerIsRefusedAfterSignalBeforePutterReacquiresLock() throws Exception { + final FutureCompletingBlockingQueue queue = new FutureCompletingBlockingQueue<>(1); + assertThat(queue.put(0, 0)).isTrue(); + + final AtomicBoolean putterResult = new AtomicBoolean(true); + final Thread putter = + new Thread(() -> putterResult.set(putUnchecked(queue, 1, 1)), "signalledPutter"); + putter.start(); + try { + CommonTestUtils.waitUntilCondition(() -> queue.getNumberOfQueuedPutters() == 1); + + final Lock queueLock = queue.lock(); + queueLock.lock(); + try { + assertThat(queue.poll()).isZero(); + assertThat(queue.getNumberOfQueuedPutters()).isZero(); + + queue.wakeUpPuttingThread(1); + queue.releaseProducer(1); + assertThat(queue.hasProducerState(1)).isTrue(); + + assertThat(queue.put(2, 2)).isTrue(); + } finally { + queueLock.unlock(); + } + + joinWithinTimeout(putter); + assertThat(putterResult).isFalse(); + + queue.releaseProducer(1); + assertThat(queue.numberOfProducerStates()).isZero(); + } finally { + if (putter.isAlive()) { + queue.wakeUpPuttingThread(1); + putter.interrupt(); + queue.poll(); + putter.join(TimeUnit.SECONDS.toMillis(10)); + } + queue.releaseProducer(1); + queue.poll(); + } + } + + private static void joinWithinTimeout(Thread thread) throws InterruptedException { + thread.join(TimeUnit.SECONDS.toMillis(10)); + assertThat(thread.isAlive()).as("The putting thread should have terminated").isFalse(); + } + private static boolean putUnchecked( FutureCompletingBlockingQueue queue, int threadIndex, int value) { try { diff --git a/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/QueueProbe.java b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/QueueProbe.java new file mode 100644 index 00000000000000..e6b8461ec6585b --- /dev/null +++ b/flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/QueueProbe.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.connector.base.source.reader.synchronization; + +/** Test-only bridge exposing queue state to tests in other packages. */ +public final class QueueProbe { + private QueueProbe() {} + + /** Returns the number of producers for which the queue currently holds wakeup state. */ + public static int liveProducerStates(FutureCompletingBlockingQueue queue) { + return queue.numberOfProducerStates(); + } +}