MartijnVisser commented on code in PR #29218: URL: https://github.com/apache/flink/pull/29218#discussion_r4167443385
########## flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/QueueProbe.java: ########## @@ -0,0 +1,148 @@ +/* + * 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; + +import java.lang.reflect.Array; +import java.lang.reflect.Field; +import java.util.Map; +import java.util.Queue; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Supplier; + +/** + * Test-only access to the internal state of {@link FutureCompletingBlockingQueue}. + * + * <p>Reflection is used deliberately rather than adding accessors to the queue: a {@code + * VisibleForTesting} method in connector production code registers as a dependency on non-public + * Flink API and would need a new entry in the frozen architecture-violation store. + * + * <p>The producer-state probes handle both a map and an array representation so that they stay + * meaningful against an implementation that releases state by nulling an array slot. Such an + * implementation bounds {@link #liveProducerStates} but not {@link #producerStateStorageSize}, + * which is why both are exposed. + */ +public final class QueueProbe { + + private static final Field PRODUCER_STATES_FIELD = field("putConditionAndFlags"); Review Comment: This could call package-private accessors on the queue instead of reflecting into its fields; without `@VisibleForTesting` they don't touch the ArchUnit store. `queuedPutters` duplicates `getNumberOfQueuedPutters()` and the array branches can go. ########## flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java: ########## @@ -147,6 +135,233 @@ 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<Integer> 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(QueueProbe.liveProducerStates(queue)).isOne(); + assertThat(QueueProbe.producerStateStorageSize(queue)).isOne(); + + queue.releaseProducer(producer); + + assertThat(QueueProbe.liveProducerStates(queue)).isZero(); + assertThat(QueueProbe.producerStateStorageSize(queue)).isZero(); + } + + assertThat(queue.getNumberOfQueuedPutters()).isZero(); + } + + /** + * The storage assertions are what separate releasing the state from merely clearing it. An + * implementation that kept the {@code ConditionAndFlag[]} and only nulled the released slot + * would satisfy every live-count assertion above while still growing its backing array to the + * largest index ever seen, so this pins the storage down as well. + */ + @Test + void testSparseProducerIndexDoesNotExpandStorage() throws InterruptedException { + final FutureCompletingBlockingQueue<Integer> queue = new FutureCompletingBlockingQueue<>(1); + + assertThat(queue.put(100_000, 1)).isTrue(); + assertThat(QueueProbe.liveProducerStates(queue)).isZero(); + assertThat(QueueProbe.producerStateStorageSize(queue)).isZero(); + assertThat(queue.poll()).isOne(); + + queue.wakeUpPuttingThread(100_000); + + assertThat(QueueProbe.liveProducerStates(queue)).isOne(); + assertThat(QueueProbe.producerStateStorageSize(queue)).isOne(); + + queue.releaseProducer(100_000); + + assertThat(QueueProbe.liveProducerStates(queue)).isZero(); + assertThat(QueueProbe.producerStateStorageSize(queue)).isZero(); + } + + /** Release is idempotent, tolerates unknown ids, and only removes the target state. */ + @Test + void testReleaseProducerOnlyRemovesTheTargetState() { + final FutureCompletingBlockingQueue<Integer> queue = new FutureCompletingBlockingQueue<>(1); + + queue.wakeUpPuttingThread(3); + queue.wakeUpPuttingThread(4); + + queue.releaseProducer(7); + queue.releaseProducer(3); + queue.releaseProducer(3); + + assertThat(QueueProbe.containsProducerState(queue, 3)).isFalse(); + assertThat(QueueProbe.containsProducerState(queue, 4)).isTrue(); + assertThat(QueueProbe.liveProducerStates(queue)).isOne(); + + queue.releaseProducer(4); + assertThat(QueueProbe.liveProducerStates(queue)).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 + @Timeout(value = 60, unit = TimeUnit.SECONDS) Review Comment: Please remove the new `@Timeout`s here and in `SplitFetcherManagerTest`, so a hang still gets the CI thread dump: https://flink.apache.org/how-to-contribute/code-style-and-quality-common/#avoid-timeouts-in-junit-tests ########## flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueueTest.java: ########## @@ -86,24 +87,11 @@ void testPollEmptyQueue() throws InterruptedException { @Test void testWakeUpPut() throws InterruptedException { Review Comment: Why was this test rewritten? The original passes unchanged on this branch. ########## flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcher.java: ########## @@ -188,24 +198,34 @@ boolean runOnce() { } // execute the task outside of lock, so that it can be woken up - boolean taskFinished; + boolean taskFinished = false; + boolean taskRunCompleted = false; Review Comment: I would prefer to revert this change. A fetcher failure is terminal (`checkErrors()` keeps throwing), so the reader and its queue are torn down. With `SplitFetcher` reverted, only the two tests asserting it fail. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
