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]