Savonitar commented on code in PR #318:
URL: 
https://github.com/apache/flink-connector-kafka/pull/318#discussion_r4157621571


##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()

Review Comment:
   Looks like this method and 
[testBoundedSourceWithoutPartitionsSignalsNoMoreSplits 
](https://github.com/apache/flink-connector-kafka/pull/318/changes#diff-3ba711e9cc1b504e425c47244c177882aff434cb8f3d7344bd7a72209b2b013bR1116)
 overlap a lot: "when" block (I mean try with resources) and assertion. 
   Is it better to merge them to one parameterized test? Or it is intentional 
separation to clearly distinguish restores state from a fresh empty source?
   



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no

Review Comment:
   > by enabling partition discovery
   > makes the partition change non-empty 
   
   if I understand it correctly, it says that if we enable partition discovery, 
it will not be blocked, but will make change non-empty? 
   Wouldn't the test just hang in that case? It is a bit misleading, if so. 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on
+     * an empty partition change: it has no readers to tell that the input 
ended, and treating the

Review Comment:
   > it has no readers to tell that the input ended
   
   Not fully following, here 
https://github.com/apache/flink-connector-kafka/pull/318/changes#diff-3ba711e9cc1b504e425c47244c177882aff434cb8f3d7344bd7a72209b2b013bR1087
 we have a reader 
   ```
   registerReader(context, enumerator, READER0);
   ```
   
   Is it a contradiction? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on
+     * an empty partition change: it has no readers to tell that the input 
ended, and treating the
+     * discovery as finished would make a later restore initialise the 
partitions that appeared
+     * meanwhile from the earliest offset instead of the configured one. It 
must keep returning
+     * early, and this test fails if the condition in checkPartitionChanges is 
ever widened past
+     * bounded sources.
+     */
+    @Test
+    public void testUnboundedSourceKeepsSkippingAnEmptyPartitionChange() 
throws Throwable {
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                false)) {
+            enumerator.start();
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        try {

Review Comment:
   do we need this try-catch block? why?



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on
+     * an empty partition change: it has no readers to tell that the input 
ended, and treating the
+     * discovery as finished would make a later restore initialise the 
partitions that appeared
+     * meanwhile from the earliest offset instead of the configured one. It 
must keep returning
+     * early, and this test fails if the condition in checkPartitionChanges is 
ever widened past
+     * bounded sources.
+     */
+    @Test
+    public void testUnboundedSourceKeepsSkippingAnEmptyPartitionChange() 
throws Throwable {
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                false)) {
+            enumerator.start();
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        try {
+                            softly.assertThat(
+                                            
enumerator.snapshotState(1L).initialDiscoveryFinished())
+                                    .as("An empty change must not mark the 
discovery as finished")
+                                    .isFalse();
+                        } catch (Exception e) {
+                            throw new IllegalStateException(e);
+                        }
+                        softly.assertThat(context.hasNoMoreSplits(READER0))

Review Comment:
   Unbounded never emits NoMoreSplitsEvent, right?
   In that case, could you please clarify the motivation of this check 🤔  ? It 
is some sort of regressions/precautious? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java:
##########
@@ -2195,6 +2195,56 @@ private DynamicKafkaSourceEnumerator createEnumerator(
                 applyPropertiesConsumer);
     }
 
+    /**
+     * A bounded DynamicKafkaSource restored with every partition already 
assigned must still tell
+     * its readers that no more splits are coming. Each sub-enumerator runs 
its one-time discovery
+     * and finds nothing new. One that returns early on that empty change 
never marks the discovery
+     * as finished, so the readers wait forever and the job never finishes.
+     */
+    @Test
+    public void testBoundedRestoreSignalsNoMoreSplitsWithoutPartitionChanges() 
throws Throwable {
+        final DynamicKafkaSourceEnumState restoredState = getCheckpointState();
+
+        Properties properties = new Properties();
+        // A bounded source never runs periodic discovery; 
DynamicKafkaSourceBuilder forces this.

Review Comment:
   i see here we partially rebuild existing helper 
https://github.com/apache/flink-connector-kafka/blob/ae656dd6fe23f19f2e7e79cacd2d5e0341f3706b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java#L2154
 
   does it make sense to just create an overload for existing helper method 
(e.g. with second argument) and use it? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))

Review Comment:
   Should we also verify it is not notified many times? because 
`hasNoMoreSplits` is a simple yes/no check. 
   e.g.
   ```
   assertThat(context.events)
                     .as("Nothing is assigned and each reader is notified 
exactly once")
                     .containsExactly(
                             "signalNoMoreSplits:" + READER0, 
"signalNoMoreSplits:" + READER1);
   ```
   



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBoundedRestoreITCase.java:
##########
@@ -0,0 +1,241 @@
+/*
+ * 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.kafka.source;
+
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.RestartStrategyOptions;
+import org.apache.flink.configuration.StateRecoveryOptions;
+import 
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import 
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
+import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
+import org.apache.flink.core.execution.SavepointFormatType;
+import org.apache.flink.core.testutils.CommonTestUtils;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.jobmaster.JobResult;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
+import org.apache.flink.test.junit5.InjectMiniCluster;
+import org.apache.flink.test.junit5.MiniClusterExtension;
+
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * A bounded {@link KafkaSource} must still finish when the job is restored 
from a savepoint, a
+ * retained checkpoint, or after a JobManager failover.

Review Comment:
   > retained checkpoint, or after a JobManager failover.
   
   Do we have these tests? Or we have only savepoint path tested? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBoundedRestoreITCase.java:
##########
@@ -0,0 +1,241 @@
+/*
+ * 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.kafka.source;
+
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.RestartStrategyOptions;
+import org.apache.flink.configuration.StateRecoveryOptions;
+import 
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import 
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
+import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
+import org.apache.flink.core.execution.SavepointFormatType;
+import org.apache.flink.core.testutils.CommonTestUtils;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.jobmaster.JobResult;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
+import org.apache.flink.test.junit5.InjectMiniCluster;
+import org.apache.flink.test.junit5.MiniClusterExtension;
+
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * A bounded {@link KafkaSource} must still finish when the job is restored 
from a savepoint, a
+ * retained checkpoint, or after a JobManager failover.
+ *
+ * <p>All three recreate the enumerator from state, and the recreated 
enumerator finds every
+ * subscribed partition already assigned. On an unfixed enumerator that empty 
partition change makes
+ * {@code checkPartitionChanges} return before it reaches the only place that 
marks the discovery as
+ * finished, so the restored readers consume up to their stopping offset and 
then wait forever for a
+ * {@code NoMoreSplitsEvent} that is never sent.
+ *
+ * <p>Do not weaken this test by enabling partition discovery, by putting the 
stopping offset within
+ * the records produced before the savepoint, or by running it in batch mode: 
the first two stop the
+ * restored partition change from being empty, and the third removes the 
checkpoint coordinator that
+ * the savepoint needs.
+ */
+@ResourceLock("KafkaTestBase")
+public class KafkaSourceBoundedRestoreITCase {
+
+    private static final String TOPIC = 
"KafkaSourceBoundedRestoreITCase-topic";
+
+    /** Records produced before the savepoint. */
+    private static final int RECORDS_BEFORE_SAVEPOINT = 10;
+
+    /** Records produced while the job is stopped, which the restored job must 
still consume. */
+    private static final int RECORDS_AFTER_SAVEPOINT = 5;
+
+    private static final int TOTAL_RECORDS = RECORDS_BEFORE_SAVEPOINT + 
RECORDS_AFTER_SAVEPOINT;
+    private static final Duration TIMEOUT = Duration.ofMinutes(2);
+
+    /** Values seen by the pipeline in the current run, in the order they were 
seen. */
+    private static final List<Integer> COLLECTED = new 
CopyOnWriteArrayList<>();
+
+    @TempDir private Path savepointBasePath;
+
+    @RegisterExtension
+    public static final MiniClusterExtension MINI_CLUSTER =
+            new MiniClusterExtension(
+                    new MiniClusterResourceConfiguration.Builder()
+                            .setNumberTaskManagers(1)
+                            .setNumberSlotsPerTaskManager(2)
+                            .build());
+
+    @BeforeEach
+    public void setup() throws Throwable {
+        COLLECTED.clear();
+        KafkaSourceTestEnv.setup();
+        KafkaSourceTestEnv.createTestTopic(TOPIC, 1, 1);
+        produce(0, RECORDS_BEFORE_SAVEPOINT);
+    }
+
+    @AfterEach
+    public void tearDown() throws Exception {
+        KafkaSourceTestEnv.tearDown();
+    }
+
+    private static void produce(int fromValue, int count) throws Throwable {
+        List<ProducerRecord<String, Integer>> records = new ArrayList<>();
+        for (int i = fromValue; i < fromValue + count; i++) {
+            records.add(new ProducerRecord<>(TOPIC, 0, "key-" + i, i));
+        }
+        KafkaSourceTestEnv.produceToKafka(records);
+    }
+
+    /**
+     * The job graph is built by one helper for both runs so that the restored 
job has the identical
+     * topology, and therefore the identical operator IDs, as the job the 
savepoint was taken from.
+     */
+    private JobGraph getJobGraph(Configuration extraConf) {
+        KafkaSource<Integer> source =
+                KafkaSource.<Integer>builder()
+                        
.setBootstrapServers(KafkaSourceTestEnv.brokerConnectionStrings)
+                        .setTopics(TOPIC)
+                        .setGroupId("KafkaSourceBoundedRestoreITCase")
+                        .setStartingOffsets(OffsetsInitializer.earliest())
+                        // Beyond what has been produced so far, so the first 
run cannot finish on
+                        // its own and the savepoint is always taken from a 
running source.
+                        .setBounded(
+                                OffsetsInitializer.offsets(
+                                        Collections.singletonMap(
+                                                new TopicPartition(TOPIC, 0),
+                                                (long) TOTAL_RECORDS)))
+                        .setDeserializer(
+                                KafkaRecordDeserializationSchema.valueOnly(
+                                        IntegerDeserializer.class))
+                        .build();
+
+        Configuration configuration = new Configuration();
+        configuration.addAll(extraConf);
+        // A hang must surface as a timeout on the job, not as restart churn.
+        configuration.set(RestartStrategyOptions.RESTART_STRATEGY, "disable");
+
+        StreamExecutionEnvironment env =
+                
StreamExecutionEnvironment.getExecutionEnvironment(configuration);
+        env.setParallelism(1);
+        DataStream<Integer> stream =
+                env.fromSource(source, WatermarkStrategy.noWatermarks(), 
"kafka-source")
+                        .uid("kafka-source")
+                        .map(
+                                (MapFunction<Integer, Integer>)
+                                        value -> {
+                                            COLLECTED.add(value);
+                                            return value;
+                                        })
+                        .uid("collector");
+        stream.sinkTo(new DiscardingSink<>()).uid("sink");
+        return env.getStreamGraph().getJobGraph();
+    }
+
+    /**
+     * Waits for the job to finish. A job that fails or is cancelled instead 
rethrows its own cause,
+     * so that an unrelated failure is reported as itself rather than as the 
hang under test.
+     */
+    private static void awaitJobFinished(MiniCluster miniCluster, JobID jobId) 
throws Exception {
+        final JobResult jobResult;
+        try {
+            jobResult =
+                    miniCluster
+                            .requestJobResult(jobId)
+                            .get(TIMEOUT.toMillis(), TimeUnit.MILLISECONDS);
+        } catch (TimeoutException e) {
+            throw new AssertionError(
+                    "The restored bounded job did not finish; the readers were 
never told that no "

Review Comment:
   Curious, why do we assume that any timeout exception is exactly:
   > The restored bounded job did not finish



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on

Review Comment:
   > The mirror of the two tests above
   
   1. Is it true? (or one of them is below?).
   2. Despite ^ , does it make sense to rephrase this comment? E.g. what if in 
the future people reorder tests? It is misleading.



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on
+     * an empty partition change: it has no readers to tell that the input 
ended, and treating the
+     * discovery as finished would make a later restore initialise the 
partitions that appeared
+     * meanwhile from the earliest offset instead of the configured one. It 
must keep returning
+     * early, and this test fails if the condition in checkPartitionChanges is 
ever widened past
+     * bounded sources.
+     */
+    @Test
+    public void testUnboundedSourceKeepsSkippingAnEmptyPartitionChange() 
throws Throwable {
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                false)) {
+            enumerator.start();
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        try {
+                            softly.assertThat(
+                                            
enumerator.snapshotState(1L).initialDiscoveryFinished())
+                                    .as("An empty change must not mark the 
discovery as finished")
+                                    .isFalse();
+                        } catch (Exception e) {
+                            throw new IllegalStateException(e);
+                        }
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("An unbounded source never runs out of 
splits")
+                                .isFalse();
+                        
softly.assertThat(context.getSplitsAssignmentSequence())
+                                .as("There is no partition to assign")
+                                .isEmpty();
+                    });
+        }
+    }
+
+    /**
+     * The same early return leaves a bounded source that subscribes to no 
existing partition
+     * hanging on a fresh start: the partition change is empty from the first 
discovery on, so the
+     * readers are never told that no more splits are coming.
+     */
+    @Test
+    public void testBoundedSourceWithoutPartitionsSignalsNoMoreSplits() throws 
Throwable {
+        // A single explicit name, because the helper joins the names into a 
pattern and an empty
+        // collection would compile to a pattern that matches every topic.
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                false)) {
+            enumerator.start();
+
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("There is no partition to assign")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * A restored bounded enumerator must not tell a reader that the input 
ended before the
+     * discovery that follows the restore has assigned that reader its splits. 
A partition created
+     * while the job was down is only found by that discovery, so signalling 
from addReader on the
+     * strength of the restored initialDiscoveryFinished flag would finish the 
reader first and
+     * assign to it afterwards. That is the failure that apache/flink#21909 
was rejected for, and
+     * this test pins the ordering that avoids it.
+     */
+    @Test
+    public void testRestoredBoundedEnumeratorAssignsBeforeItSignals() throws 
Throwable {
+        // Only TOPIC1 was assigned before the savepoint; TOPIC2 stands in for 
the partitions that
+        // appeared while the job was down and that only the post-restore 
discovery can find.
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(Collections.singleton(TOPIC1)).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (RecordingSplitEnumeratorContext context =
+                        new RecordingSplitEnumeratorContext(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+            // Registers before the discovery callback, which is the ordering 
a restore produces.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+
+            assertThat(context.events)

Review Comment:
   Here we have multiple assertions, does it better to have one:
   ```
   assertThat(context.events)
             .as("The reader must get its new splits before it is told the 
input ended")
             .containsExactly("assignSplits:" + READER0, "signalNoMoreSplits:" 
+ READER0);
   ```
   ? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBoundedRestoreITCase.java:
##########
@@ -0,0 +1,241 @@
+/*
+ * 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.kafka.source;
+
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.RestartStrategyOptions;
+import org.apache.flink.configuration.StateRecoveryOptions;
+import 
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import 
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
+import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
+import org.apache.flink.core.execution.SavepointFormatType;
+import org.apache.flink.core.testutils.CommonTestUtils;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.jobmaster.JobResult;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
+import org.apache.flink.test.junit5.InjectMiniCluster;
+import org.apache.flink.test.junit5.MiniClusterExtension;
+
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * A bounded {@link KafkaSource} must still finish when the job is restored 
from a savepoint, a
+ * retained checkpoint, or after a JobManager failover.
+ *
+ * <p>All three recreate the enumerator from state, and the recreated 
enumerator finds every
+ * subscribed partition already assigned. On an unfixed enumerator that empty 
partition change makes
+ * {@code checkPartitionChanges} return before it reaches the only place that 
marks the discovery as
+ * finished, so the restored readers consume up to their stopping offset and 
then wait forever for a
+ * {@code NoMoreSplitsEvent} that is never sent.
+ *
+ * <p>Do not weaken this test by enabling partition discovery, by putting the 
stopping offset within

Review Comment:
   > Do not weaken this test by enabling partition discovery
   
   If people enable partition discovery, what will happen? AFAIK it is forced 
as -1 for Bounded 
https://github.com/apache/flink-connector-kafka/blob/main/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java#L514
   
   Just curious, do we anyway want to have this comment? If yes, could you 
please share why? 



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java:
##########
@@ -976,6 +1002,241 @@ public void testCheckSourceIntegrityFromProperties() 
throws Exception {
         }
     }
 
+    /**
+     * A bounded enumerator that is restored with every subscribed partition 
already assigned must
+     * still tell its readers that no more splits are coming. The restored 
enumerator runs its
+     * one-time discovery, finds nothing new, and on an unfixed enumerator 
returns early before
+     * reaching the only place that sets {@code noMoreNewPartitionSplits}, so 
the readers wait for a
+     * {@code NoMoreSplitsEvent} that never arrives and the job never finishes.
+     *
+     * <p>Do not weaken this test by enabling partition discovery or by 
leaving a partition out of
+     * the restored state: either makes the partition change non-empty and the 
early return is no
+     * longer taken.
+     */
+    @Test
+    public void 
testRestoredBoundedEnumeratorSignalsNoMoreSplitsWithoutPartitionChanges()
+            throws Throwable {
+        Set<KafkaPartitionSplit> restoredAssignedSplits =
+                
KafkaSourceTestEnv.getPartitionsForTopics(PRE_EXISTING_TOPICS).stream()
+                        .map(tp -> new KafkaPartitionSplit(tp, 0L))
+                        .collect(Collectors.toSet());
+        assertThat(restoredAssignedSplits).isNotEmpty();
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                PRE_EXISTING_TOPICS,
+                                restoredAssignedSplits,
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                true)) {
+            enumerator.start();
+
+            // A reader that re-registers before the discovery callback 
returns.
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            // A reader that re-registers after it.
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("Every partition is already assigned, so nothing may 
be assigned again")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * The mirror of the two tests above. An unbounded source with one-time 
discovery cannot act on
+     * an empty partition change: it has no readers to tell that the input 
ended, and treating the
+     * discovery as finished would make a later restore initialise the 
partitions that appeared
+     * meanwhile from the earliest offset instead of the configured one. It 
must keep returning
+     * early, and this test fails if the condition in checkPartitionChanges is 
ever widened past
+     * bounded sources.
+     */
+    @Test
+    public void testUnboundedSourceKeepsSkippingAnEmptyPartitionChange() 
throws Throwable {
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                false)) {
+            enumerator.start();
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        try {
+                            softly.assertThat(
+                                            
enumerator.snapshotState(1L).initialDiscoveryFinished())
+                                    .as("An empty change must not mark the 
discovery as finished")
+                                    .isFalse();
+                        } catch (Exception e) {
+                            throw new IllegalStateException(e);
+                        }
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("An unbounded source never runs out of 
splits")
+                                .isFalse();
+                        
softly.assertThat(context.getSplitsAssignmentSequence())
+                                .as("There is no partition to assign")
+                                .isEmpty();
+                    });
+        }
+    }
+
+    /**
+     * The same early return leaves a bounded source that subscribes to no 
existing partition
+     * hanging on a fresh start: the partition change is empty from the first 
discovery on, so the
+     * readers are never told that no more splits are coming.
+     */
+    @Test
+    public void testBoundedSourceWithoutPartitionsSignalsNoMoreSplits() throws 
Throwable {
+        // A single explicit name, because the helper joins the names into a 
pattern and an empty
+        // collection would compile to a pattern that matches every topic.
+        Collection<String> noSuchTopic = 
Collections.singleton("topic-that-does-not-exist");
+
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                -1,
+                                noSuchTopic,
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                new Properties(),
+                                OffsetsInitializer.earliest(),
+                                Boundedness.BOUNDED,
+                                false)) {
+            enumerator.start();
+
+            registerReader(context, enumerator, READER0);
+            runOneTimePartitionDiscovery(context);
+            registerReader(context, enumerator, READER1);
+
+            assertThat(context.getSplitsAssignmentSequence())
+                    .as("There is no partition to assign")
+                    .isEmpty();
+            SoftAssertions.assertSoftly(
+                    softly -> {
+                        softly.assertThat(context.hasNoMoreSplits(READER0))
+                                .as("Reader registered before the discovery 
callback")
+                                .isTrue();
+                        softly.assertThat(context.hasNoMoreSplits(READER1))
+                                .as("Reader registered after the discovery 
callback")
+                                .isTrue();
+                    });
+        }
+    }
+
+    /**
+     * A restored bounded enumerator must not tell a reader that the input 
ended before the
+     * discovery that follows the restore has assigned that reader its splits. 
A partition created
+     * while the job was down is only found by that discovery, so signalling 
from addReader on the
+     * strength of the restored initialDiscoveryFinished flag would finish the 
reader first and
+     * assign to it afterwards. That is the failure that apache/flink#21909 
was rejected for, and

Review Comment:
   > apache/flink#21909
   
   is it same  as 
https://github.com/apache/flink-connector-kafka/pull/318/changes#paneldiscussion_r4153236816?
 
   Do we want to have the reference to PR here?



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBoundedRestoreITCase.java:
##########
@@ -0,0 +1,241 @@
+/*
+ * 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.kafka.source;
+
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.RestartStrategyOptions;
+import org.apache.flink.configuration.StateRecoveryOptions;
+import 
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import 
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
+import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
+import org.apache.flink.core.execution.SavepointFormatType;
+import org.apache.flink.core.testutils.CommonTestUtils;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.jobmaster.JobResult;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
+import org.apache.flink.test.junit5.InjectMiniCluster;
+import org.apache.flink.test.junit5.MiniClusterExtension;
+
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * A bounded {@link KafkaSource} must still finish when the job is restored 
from a savepoint, a
+ * retained checkpoint, or after a JobManager failover.
+ *
+ * <p>All three recreate the enumerator from state, and the recreated 
enumerator finds every
+ * subscribed partition already assigned. On an unfixed enumerator that empty 
partition change makes
+ * {@code checkPartitionChanges} return before it reaches the only place that 
marks the discovery as
+ * finished, so the restored readers consume up to their stopping offset and 
then wait forever for a
+ * {@code NoMoreSplitsEvent} that is never sent.
+ *
+ * <p>Do not weaken this test by enabling partition discovery, by putting the 
stopping offset within
+ * the records produced before the savepoint, or by running it in batch mode: 
the first two stop the
+ * restored partition change from being empty, and the third removes the 
checkpoint coordinator that
+ * the savepoint needs.
+ */
+@ResourceLock("KafkaTestBase")
+public class KafkaSourceBoundedRestoreITCase {
+
+    private static final String TOPIC = 
"KafkaSourceBoundedRestoreITCase-topic";
+
+    /** Records produced before the savepoint. */
+    private static final int RECORDS_BEFORE_SAVEPOINT = 10;
+
+    /** Records produced while the job is stopped, which the restored job must 
still consume. */
+    private static final int RECORDS_AFTER_SAVEPOINT = 5;
+
+    private static final int TOTAL_RECORDS = RECORDS_BEFORE_SAVEPOINT + 
RECORDS_AFTER_SAVEPOINT;
+    private static final Duration TIMEOUT = Duration.ofMinutes(2);
+
+    /** Values seen by the pipeline in the current run, in the order they were 
seen. */
+    private static final List<Integer> COLLECTED = new 
CopyOnWriteArrayList<>();
+
+    @TempDir private Path savepointBasePath;
+
+    @RegisterExtension
+    public static final MiniClusterExtension MINI_CLUSTER =
+            new MiniClusterExtension(
+                    new MiniClusterResourceConfiguration.Builder()
+                            .setNumberTaskManagers(1)
+                            .setNumberSlotsPerTaskManager(2)
+                            .build());
+
+    @BeforeEach
+    public void setup() throws Throwable {
+        COLLECTED.clear();
+        KafkaSourceTestEnv.setup();
+        KafkaSourceTestEnv.createTestTopic(TOPIC, 1, 1);
+        produce(0, RECORDS_BEFORE_SAVEPOINT);
+    }
+
+    @AfterEach
+    public void tearDown() throws Exception {
+        KafkaSourceTestEnv.tearDown();

Review Comment:
   Does it make sense to cleanup jobs here (which will improve failure path 
cleanup)?
   e.g. we could do:
   ```
   miniCluster.cancelJob(job.getJobId()).get(...);
   miniCluster.requestJobResult(job.getJobId()).get(...);
   ```



-- 
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]

Reply via email to