This is an automated email from the ASF dual-hosted git repository.

1996fanrui pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git

commit a5b7bc38f506fad6e0ef046d0b6d1451490445b1
Author: Efrat Levitan <[email protected]>
AuthorDate: Sun Jul 12 20:24:50 2026 +0300

    [FLINK-40128][connector] Implement topic integrity awareness across kafka 
subscribers
    
    For user convenience, and to comply with kafka topic-name based APIs, flink 
will figure out the topicId on its own (i.e user will not need to provide it).
    
    To ensure topic integrity over recoveries, topic names → ids mapping will 
be preserved in KafkaSourceEnumeratorState.topicIntegrityMapping checkpointed 
state, and propagated to kafka subscriber post recovery.
    
    For older state versions, or when the feature is disabled, this field holds 
an empty dataset so by default no additional state is acquired.
    
    The implementation avoids additional calls to kafka server, and normally 
bases of existing retrieved metadata (see disclaimer),
    
    On every metadata fetch (on startup / upon partition discovery interval), 
the topic Id from kafka server is compared against the stored topic id. A 
mismatch will trigger a global failure.
---
 .../c0d94764-76a0-4c50-b617-70b1754c4612           | 18 ++---
 .../source/enumerator/KafkaSourceEnumState.java    | 25 +++++++
 .../enumerator/KafkaSourceEnumStateSerializer.java | 75 ++++++++++++++++++++-
 .../source/enumerator/KafkaSourceEnumerator.java   | 38 ++++++++++-
 .../subscriber/PartitionSetSubscriber.java         | 23 +++++--
 .../enumerator/subscriber/TopicListSubscriber.java | 23 +++++--
 .../subscriber/TopicPatternSubscriber.java         | 23 +++++--
 .../kafka/source/KafkaSourceBuilderTest.java       | 51 +++++++++++++++
 .../KafkaSourceEnumStateSerializerTest.java        | 42 +++++++++++-
 .../enumerator/KafkaSourceEnumeratorTest.java      | 41 ++++++++++++
 .../enumerator/subscriber/KafkaSubscriberTest.java | 76 +++++++++++++++++++++-
 11 files changed, 407 insertions(+), 28 deletions(-)

diff --git 
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
 
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
index 478ac02a..0a64cb1b 100644
--- 
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
+++ 
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
@@ -13,14 +13,14 @@ Method 
<org.apache.flink.connector.kafka.dynamic.source.DynamicKafkaSource.getKa
 Method 
<org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroupManager.close()>
 calls method 
<org.apache.flink.runtime.metrics.groups.AbstractMetricGroup.close()> in 
(KafkaClusterMetricGroupManager.java:73)
 Method 
<org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroupManager.close(java.lang.String)>
 calls method 
<org.apache.flink.runtime.metrics.groups.AbstractMetricGroup.close()> in 
(KafkaClusterMetricGroupManager.java:62)
 Method 
<org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroupManager.register(java.lang.String,
 
org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroup)>
 checks instanceof 
<org.apache.flink.runtime.metrics.groups.AbstractMetricGroup> in 
(KafkaClusterMetricGroupManager.java:42)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls constructor 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.<init>(int)>
 in (DynamicKafkaSourceReader.java:479)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:476)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:489)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls constructor 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.<init>(int)>
 in (DynamicKafkaSourceReader.java:592)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:589)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.completeAndResetAvailabilityHelper()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:602)
 Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.getAvailabilityHelperSize()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(DynamicKafkaSourceReader.java:0)
 Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.isActivelyConsumingSplits()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(DynamicKafkaSourceReader.java:0)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.isAvailable()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:385)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.isAvailable()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.resetToUnAvailable()>
 in (DynamicKafkaSourceReader.java:383)
-Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.syncAvailabilityHelperWithReaders()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.anyOf(int,
 java.util.concurrent.CompletableFuture)> in (DynamicKafkaSourceReader.java:500)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.isAvailable()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.getAvailableFuture()>
 in (DynamicKafkaSourceReader.java:437)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.isAvailable()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.resetToUnAvailable()>
 in (DynamicKafkaSourceReader.java:435)
+Method 
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.syncAvailabilityHelperWithReaders()>
 calls method 
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.anyOf(int,
 java.util.concurrent.CompletableFuture)> in (DynamicKafkaSourceReader.java:613)
 Method 
<org.apache.flink.connector.kafka.sink.ExactlyOnceKafkaWriter.getProducerPool()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(ExactlyOnceKafkaWriter.java:0)
 Method 
<org.apache.flink.connector.kafka.sink.ExactlyOnceKafkaWriter.getTransactionalIdPrefix()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(ExactlyOnceKafkaWriter.java:0)
 Method 
<org.apache.flink.connector.kafka.sink.KafkaSink.addPostCommitTopology(org.apache.flink.streaming.api.datastream.DataStream)>
 calls method <org.apache.flink.api.dag.Transformation.getCoLocationGroupKey()> 
in (KafkaSink.java:183)
@@ -42,16 +42,18 @@ Method 
<org.apache.flink.connector.kafka.source.KafkaSource.getStoppingOffsetsIn
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumStateSerializer.serializeV1(java.util.Collection)>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumStateSerializer.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumStateSerializer.serializeV2(java.util.Collection,
 boolean)> is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumStateSerializer.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumStateSerializer.serializeV3(org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumState)>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumStateSerializer.java:0)
+Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumStateSerializer.serializeV4(org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumState)>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumStateSerializer.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.deepCopyProperties(java.util.Properties,
 java.util.Properties)> is annotated with 
<org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumerator.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getPartitionChange(java.util.Set,
 boolean)> is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumerator.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getPendingPartitionSplitAssignment()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumerator.java:0)
 Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getSplitOwner(org.apache.kafka.common.TopicPartition,
 int)> is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumerator.java:0)
+Method 
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.topicIntegrityCheckEnabled()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceEnumerator.java:0)
 Method 
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.consumer()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaPartitionSplitReader.java:0)
 Method 
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.getPollTimeout()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaPartitionSplitReader.java:0)
 Method 
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.setConsumerClientRack(java.util.Properties,
 java.lang.String)> is annotated with 
<org.apache.flink.annotation.VisibleForTesting> in 
(KafkaPartitionSplitReader.java:0)
 Method 
<org.apache.flink.connector.kafka.source.reader.KafkaSourceReader.getNumAliveFetchers()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceReader.java:0)
 Method 
<org.apache.flink.connector.kafka.source.reader.KafkaSourceReader.getOffsetsToCommit()>
 is annotated with <org.apache.flink.annotation.VisibleForTesting> in 
(KafkaSourceReader.java:0)
 Method 
<org.apache.flink.streaming.connectors.kafka.table.DynamicKafkaRecordSerializationSchema.createProjectedRow(org.apache.flink.table.data.RowData,
 org.apache.flink.types.RowKind, 
[Lorg.apache.flink.table.data.RowData$FieldGetter;)> has parameter of type 
<[Lorg.apache.flink.table.data.RowData$FieldGetter;> in 
(DynamicKafkaRecordSerializationSchema.java:0)
-Method 
<org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptionsUtil.createKeyFormatProjection(org.apache.flink.configuration.ReadableConfig,
 org.apache.flink.table.types.DataType)> calls method 
<org.apache.flink.table.types.logical.utils.LogicalTypeChecks.getFieldNames(org.apache.flink.table.types.logical.LogicalType)>
 in (KafkaConnectorOptionsUtil.java:520)
-Method 
<org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptionsUtil.createValueFormatProjection(org.apache.flink.configuration.ReadableConfig,
 org.apache.flink.table.types.DataType)> calls method 
<org.apache.flink.table.types.logical.utils.LogicalTypeChecks.getFieldCount(org.apache.flink.table.types.logical.LogicalType)>
 in (KafkaConnectorOptionsUtil.java:564)
+Method 
<org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptionsUtil.createKeyFormatProjection(org.apache.flink.configuration.ReadableConfig,
 org.apache.flink.table.types.DataType)> calls method 
<org.apache.flink.table.types.logical.utils.LogicalTypeChecks.getFieldNames(org.apache.flink.table.types.logical.LogicalType)>
 in (KafkaConnectorOptionsUtil.java:521)
+Method 
<org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptionsUtil.createValueFormatProjection(org.apache.flink.configuration.ReadableConfig,
 org.apache.flink.table.types.DataType)> calls method 
<org.apache.flink.table.types.logical.utils.LogicalTypeChecks.getFieldCount(org.apache.flink.table.types.logical.LogicalType)>
 in (KafkaConnectorOptionsUtil.java:565)
 Method 
<org.apache.flink.streaming.connectors.kafka.table.KafkaDynamicSink.getFieldGetters(java.util.List,
 [I)> has return type <[Lorg.apache.flink.table.data.RowData$FieldGetter;> in 
(KafkaDynamicSink.java:0)
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumState.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumState.java
index 20472ac3..5de1bffa 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumState.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumState.java
@@ -22,7 +22,9 @@ import org.apache.flink.annotation.Internal;
 import org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.HashSet;
+import java.util.Map;
 import java.util.Set;
 import java.util.stream.Collectors;
 
@@ -38,16 +40,34 @@ public class KafkaSourceEnumState {
      */
     private final boolean initialDiscoveryFinished;
 
+    private final Map<String, String> trackedTopicIdsByName;
+
     public KafkaSourceEnumState(
             Set<SplitAndAssignmentStatus> splits, boolean 
initialDiscoveryFinished) {
+        this(splits, initialDiscoveryFinished, Collections.emptyMap());
+    }
+
+    public KafkaSourceEnumState(
+            Set<SplitAndAssignmentStatus> splits,
+            boolean initialDiscoveryFinished,
+            Map<String, String> trackedTopicIdsByName) {
         this.splits = splits;
         this.initialDiscoveryFinished = initialDiscoveryFinished;
+        this.trackedTopicIdsByName = trackedTopicIdsByName;
     }
 
     public KafkaSourceEnumState(
             Collection<KafkaPartitionSplit> assignedSplits,
             Collection<KafkaPartitionSplit> unassignedSplits,
             boolean initialDiscoveryFinished) {
+        this(assignedSplits, unassignedSplits, initialDiscoveryFinished, 
Collections.emptyMap());
+    }
+
+    public KafkaSourceEnumState(
+            Collection<KafkaPartitionSplit> assignedSplits,
+            Collection<KafkaPartitionSplit> unassignedSplits,
+            boolean initialDiscoveryFinished,
+            Map<String, String> trackedTopicIdsByName) {
         this.splits = new HashSet<>();
         splits.addAll(
                 assignedSplits.stream()
@@ -64,6 +84,7 @@ public class KafkaSourceEnumState {
                                                 topicPartition, 
AssignmentStatus.UNASSIGNED))
                         .collect(Collectors.toSet()));
         this.initialDiscoveryFinished = initialDiscoveryFinished;
+        this.trackedTopicIdsByName = trackedTopicIdsByName;
     }
 
     public Set<SplitAndAssignmentStatus> splits() {
@@ -82,6 +103,10 @@ public class KafkaSourceEnumState {
         return initialDiscoveryFinished;
     }
 
+    public Map<String, String> trackedTopicIdsByName() {
+        return trackedTopicIdsByName;
+    }
+
     private Collection<KafkaPartitionSplit> filterByAssignmentStatus(
             AssignmentStatus assignmentStatus) {
         return splits.stream()
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializer.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializer.java
index c42194b7..e1a98820 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializer.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializer.java
@@ -33,6 +33,7 @@ import java.io.DataInputStream;
 import java.io.DataOutputStream;
 import java.io.IOException;
 import java.util.Collection;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Map;
 import java.util.Set;
@@ -62,9 +63,16 @@ public class KafkaSourceEnumStateSerializer
      */
     private static final int VERSION_2 = 2;
 
+    /**
+     * state of VERSION_3 contains splits with status: ASSIGNED or 
UNASSIGNED_INITIAL and
+     * initialDiscoveryFinished.
+     */
     private static final int VERSION_3 = 3;
 
-    private static final int CURRENT_VERSION = VERSION_3;
+    /** state of version 4 contains additional trackedTopicIdsByName field for 
topic id tracking. */
+    private static final int VERSION_4 = 4;
+
+    private static final int CURRENT_VERSION = VERSION_4;
 
     private static final KafkaPartitionSplitSerializer SPLIT_SERIALIZER =
             new KafkaPartitionSplitSerializer();
@@ -76,7 +84,7 @@ public class KafkaSourceEnumStateSerializer
 
     @Override
     public byte[] serialize(KafkaSourceEnumState enumState) throws IOException 
{
-        return serializeV3(enumState);
+        return serializeV4(enumState);
     }
 
     @VisibleForTesting
@@ -102,6 +110,8 @@ public class KafkaSourceEnumStateSerializer
     @Override
     public KafkaSourceEnumState deserialize(int version, byte[] serialized) 
throws IOException {
         switch (version) {
+            case VERSION_4:
+                return deserializeVersion4(serialized);
             case VERSION_3:
                 return deserializeVersion3(serialized);
             case VERSION_2:
@@ -246,4 +256,65 @@ public class KafkaSourceEnumStateSerializer
             return new KafkaSourceEnumState(partitions, 
initialDiscoveryFinished);
         }
     }
+
+    @VisibleForTesting
+    static byte[] serializeV4(KafkaSourceEnumState enumState) throws 
IOException {
+        Set<SplitAndAssignmentStatus> splits = enumState.splits();
+        boolean initialDiscoveryFinished = 
enumState.initialDiscoveryFinished();
+        Map<String, String> trackedTopicIdsByName = 
enumState.trackedTopicIdsByName();
+        try (ByteArrayOutputStream baos = new ByteArrayOutputStream();
+                DataOutputStream out = new DataOutputStream(baos)) {
+            out.writeInt(splits.size());
+            out.writeInt(SPLIT_SERIALIZER.getVersion());
+            for (SplitAndAssignmentStatus split : splits) {
+                final byte[] splitBytes = 
SPLIT_SERIALIZER.serialize(split.split());
+                out.writeInt(splitBytes.length);
+                out.write(splitBytes);
+                out.writeInt(split.assignmentStatus().getStatusCode());
+            }
+            out.writeBoolean(initialDiscoveryFinished);
+            out.writeInt(trackedTopicIdsByName.size());
+            for (Map.Entry<String, String> entry : 
trackedTopicIdsByName.entrySet()) {
+                out.writeUTF(entry.getKey());
+                out.writeUTF(entry.getValue());
+            }
+            out.flush();
+            return baos.toByteArray();
+        }
+    }
+
+    private static KafkaSourceEnumState deserializeVersion4(byte[] serialized) 
throws IOException {
+
+        final KafkaPartitionSplitSerializer splitSerializer = new 
KafkaPartitionSplitSerializer();
+
+        try (ByteArrayInputStream bais = new ByteArrayInputStream(serialized);
+                DataInputStream in = new DataInputStream(bais)) {
+
+            final int numPartitions = in.readInt();
+            final int splitVersion = in.readInt();
+            Set<SplitAndAssignmentStatus> partitions = new 
HashSet<>(numPartitions);
+
+            for (int i = 0; i < numPartitions; i++) {
+                final KafkaPartitionSplit split =
+                        splitSerializer.deserialize(splitVersion, 
in.readNBytes(in.readInt()));
+                final int statusCode = in.readInt();
+                partitions.add(
+                        new SplitAndAssignmentStatus(
+                                split, 
AssignmentStatus.ofStatusCode(statusCode)));
+            }
+            final boolean initialDiscoveryFinished = in.readBoolean();
+            final int trackedTopicIdsByNameSize = in.readInt();
+            final Map<String, String> trackedTopicIdsByName =
+                    new HashMap<>(trackedTopicIdsByNameSize);
+            for (int i = 0; i < trackedTopicIdsByNameSize; i++) {
+                trackedTopicIdsByName.put(in.readUTF(), in.readUTF());
+            }
+            if (in.available() > 0) {
+                throw new IOException("Unexpected trailing bytes in serialized 
topic partitions");
+            }
+
+            return new KafkaSourceEnumState(
+                    partitions, initialDiscoveryFinished, 
trackedTopicIdsByName);
+        }
+    }
 }
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java
index 4eb67181..2f936706 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java
@@ -26,6 +26,8 @@ import 
org.apache.flink.api.connector.source.SplitEnumeratorContext;
 import org.apache.flink.api.connector.source.SplitsAssignment;
 import org.apache.flink.connector.kafka.source.KafkaSourceOptions;
 import 
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicIntegrityProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataSettable;
 import 
org.apache.flink.connector.kafka.source.enumerator.subscriber.KafkaSubscriber;
 import org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit;
 import org.apache.flink.util.FlinkRuntimeException;
@@ -145,6 +147,10 @@ public class KafkaSourceEnumerator
     // this flag will be marked as true if initial partitions are discovered 
after enumerator starts
     private boolean initialDiscoveryFinished;
 
+    private final boolean topicIntegrityCheckEnabled;
+
+    private final TopicIntegrityProvider topicIntegrityProvider;
+
     public KafkaSourceEnumerator(
             KafkaSubscriber subscriber,
             OffsetsInitializer startingOffsetInitializer,
@@ -213,6 +219,24 @@ public class KafkaSourceEnumerator
         this.initialDiscoveryFinished = 
kafkaSourceEnumState.initialDiscoveryFinished();
         this.assignedSplits = 
indexByPartition(kafkaSourceEnumState.assignedSplits());
         this.unassignedSplits = 
indexByPartition(kafkaSourceEnumState.unassignedSplits());
+        this.topicIntegrityCheckEnabled =
+                KafkaSourceOptions.getOption(
+                        properties,
+                        KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED,
+                        Boolean::parseBoolean);
+        this.topicIntegrityProvider =
+                new TopicIntegrityProvider(
+                        topicIntegrityCheckEnabled
+                                ? kafkaSourceEnumState.trackedTopicIdsByName()
+                                : Collections.emptyMap());
+        LOG.debug(
+                "KafkaSourceEnumerator initialized with assignedSplits: {}, 
unassignedSplits: {}, "
+                        + "initialDiscoveryFinished: {}, 
topicIntegrityCheckEnabled: {}, trackedTopicIdsByName: {}",
+                assignedSplits.keySet(),
+                unassignedSplits.keySet(),
+                initialDiscoveryFinished,
+                topicIntegrityCheckEnabled,
+                topicIntegrityProvider.getTrackedTopicIdsByName());
     }
 
     private static Map<TopicPartition, KafkaPartitionSplit> indexByPartition(
@@ -241,6 +265,10 @@ public class KafkaSourceEnumerator
             addPartitionSplitChangeToPendingAssignments(preinitializedSplits);
         }
 
+        if (topicIntegrityCheckEnabled && subscriber instanceof 
TopicMetadataSettable) {
+            ((TopicMetadataSettable) 
subscriber).setTopicMetadataProvider(topicIntegrityProvider);
+        }
+
         if (partitionDiscoveryIntervalMs > 0) {
             LOG.info(
                     "Starting the KafkaSourceEnumerator for consumer group {} "
@@ -292,7 +320,10 @@ public class KafkaSourceEnumerator
     @Override
     public KafkaSourceEnumState snapshotState(long checkpointId) throws 
Exception {
         return new KafkaSourceEnumState(
-                assignedSplits.values(), unassignedSplits.values(), 
initialDiscoveryFinished);
+                assignedSplits.values(),
+                unassignedSplits.values(),
+                initialDiscoveryFinished,
+                topicIntegrityProvider.getTrackedTopicIdsByName());
     }
 
     @Override
@@ -810,4 +841,9 @@ public class KafkaSourceEnumerator
                                                     
entry.getValue().leaderEpoch())));
         }
     }
+
+    @VisibleForTesting
+    boolean topicIntegrityCheckEnabled() {
+        return topicIntegrityCheckEnabled;
+    }
 }
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/PartitionSetSubscriber.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/PartitionSetSubscriber.java
index 6ddf7c57..c27c9317 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/PartitionSetSubscriber.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/PartitionSetSubscriber.java
@@ -20,6 +20,8 @@ package 
org.apache.flink.connector.kafka.source.enumerator.subscriber;
 
 import org.apache.flink.connector.kafka.lineage.DefaultKafkaDatasetIdentifier;
 import org.apache.flink.connector.kafka.lineage.KafkaDatasetIdentifierProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataSettable;
 
 import org.apache.kafka.clients.admin.AdminClient;
 import org.apache.kafka.clients.admin.TopicDescription;
@@ -33,18 +35,30 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.stream.Collectors;
 
-import static 
org.apache.flink.connector.kafka.util.AdminUtils.getTopicMetadata;
-
 /** A subscriber for a partition set. */
-class PartitionSetSubscriber implements KafkaSubscriber, 
KafkaDatasetIdentifierProvider {
+class PartitionSetSubscriber
+        implements KafkaSubscriber, KafkaDatasetIdentifierProvider, 
TopicMetadataSettable {
     private static final long serialVersionUID = 390970375272146036L;
     private static final Logger LOG = 
LoggerFactory.getLogger(PartitionSetSubscriber.class);
     private final Set<TopicPartition> subscribedPartitions;
+    private transient TopicMetadataProvider topicMetadataProvider;
 
     PartitionSetSubscriber(Set<TopicPartition> partitions) {
         this.subscribedPartitions = partitions;
     }
 
+    @Override
+    public void setTopicMetadataProvider(TopicMetadataProvider 
topicMetadataProvider) {
+        this.topicMetadataProvider = topicMetadataProvider;
+    }
+
+    private TopicMetadataProvider getTopicMetadataProvider() {
+        if (topicMetadataProvider == null) {
+            topicMetadataProvider = TopicMetadataProvider.createDefault();
+        }
+        return topicMetadataProvider;
+    }
+
     @Override
     public Set<TopicPartition> getSubscribedTopicPartitions(AdminClient 
adminClient) {
         final Set<String> topicNames =
@@ -54,8 +68,7 @@ class PartitionSetSubscriber implements KafkaSubscriber, 
KafkaDatasetIdentifierP
 
         LOG.debug("Fetching descriptions for topics: {}", topicNames);
         final Map<String, TopicDescription> topicMetadata =
-                getTopicMetadata(adminClient, topicNames);
-
+                getTopicMetadataProvider().getTopicMetadata(adminClient, 
topicNames);
         Set<TopicPartition> existingSubscribedPartitions = new HashSet<>();
 
         for (TopicPartition subscribedPartition : this.subscribedPartitions) {
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicListSubscriber.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicListSubscriber.java
index 20234bc1..8f5f97ba 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicListSubscriber.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicListSubscriber.java
@@ -20,6 +20,8 @@ package 
org.apache.flink.connector.kafka.source.enumerator.subscriber;
 
 import org.apache.flink.connector.kafka.lineage.DefaultKafkaDatasetIdentifier;
 import org.apache.flink.connector.kafka.lineage.KafkaDatasetIdentifierProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataSettable;
 
 import org.apache.kafka.clients.admin.AdminClient;
 import org.apache.kafka.clients.admin.TopicDescription;
@@ -34,27 +36,38 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
 
-import static 
org.apache.flink.connector.kafka.util.AdminUtils.getTopicMetadata;
-
 /**
  * A subscriber to a fixed list of topics. The subscribed topics must have 
existed in the Kafka
  * cluster, otherwise an exception will be thrown.
  */
-class TopicListSubscriber implements KafkaSubscriber, 
KafkaDatasetIdentifierProvider {
+class TopicListSubscriber
+        implements KafkaSubscriber, KafkaDatasetIdentifierProvider, 
TopicMetadataSettable {
     private static final long serialVersionUID = -6917603843104947866L;
     private static final Logger LOG = 
LoggerFactory.getLogger(TopicListSubscriber.class);
     private final List<String> topics;
+    private transient TopicMetadataProvider topicMetadataProvider;
 
     TopicListSubscriber(List<String> topics) {
         this.topics = topics;
     }
 
+    @Override
+    public void setTopicMetadataProvider(TopicMetadataProvider 
topicMetadataProvider) {
+        this.topicMetadataProvider = topicMetadataProvider;
+    }
+
+    private TopicMetadataProvider getTopicMetadataProvider() {
+        if (topicMetadataProvider == null) {
+            topicMetadataProvider = TopicMetadataProvider.createDefault();
+        }
+        return topicMetadataProvider;
+    }
+
     @Override
     public Set<TopicPartition> getSubscribedTopicPartitions(AdminClient 
adminClient) {
         LOG.debug("Fetching descriptions for topics: {}", topics);
         final Map<String, TopicDescription> topicMetadata =
-                getTopicMetadata(adminClient, new HashSet<>(topics));
-
+                getTopicMetadataProvider().getTopicMetadata(adminClient, 
topics);
         Set<TopicPartition> subscribedPartitions = new HashSet<>();
         for (TopicDescription topic : topicMetadata.values()) {
             for (TopicPartitionInfo partition : topic.partitions()) {
diff --git 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicPatternSubscriber.java
 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicPatternSubscriber.java
index e525cc59..9cf16054 100644
--- 
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicPatternSubscriber.java
+++ 
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicPatternSubscriber.java
@@ -20,6 +20,8 @@ package 
org.apache.flink.connector.kafka.source.enumerator.subscriber;
 
 import org.apache.flink.connector.kafka.lineage.DefaultKafkaDatasetIdentifier;
 import org.apache.flink.connector.kafka.lineage.KafkaDatasetIdentifierProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataSettable;
 
 import org.apache.kafka.clients.admin.AdminClient;
 import org.apache.kafka.clients.admin.TopicDescription;
@@ -34,24 +36,35 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.regex.Pattern;
 
-import static 
org.apache.flink.connector.kafka.util.AdminUtils.getTopicMetadata;
-
 /** A subscriber to a topic pattern. */
-class TopicPatternSubscriber implements KafkaSubscriber, 
KafkaDatasetIdentifierProvider {
+class TopicPatternSubscriber
+        implements KafkaSubscriber, KafkaDatasetIdentifierProvider, 
TopicMetadataSettable {
     private static final long serialVersionUID = -7471048577725467797L;
     private static final Logger LOG = 
LoggerFactory.getLogger(TopicPatternSubscriber.class);
     private final Pattern topicPattern;
+    private transient TopicMetadataProvider topicMetadataProvider;
 
     TopicPatternSubscriber(Pattern topicPattern) {
         this.topicPattern = topicPattern;
     }
 
+    @Override
+    public void setTopicMetadataProvider(TopicMetadataProvider 
topicMetadataProvider) {
+        this.topicMetadataProvider = topicMetadataProvider;
+    }
+
+    private TopicMetadataProvider getTopicMetadataProvider() {
+        if (topicMetadataProvider == null) {
+            topicMetadataProvider = TopicMetadataProvider.createDefault();
+        }
+        return topicMetadataProvider;
+    }
+
     @Override
     public Set<TopicPartition> getSubscribedTopicPartitions(AdminClient 
adminClient) {
         LOG.debug("Fetching descriptions for {} topics on Kafka cluster", 
topicPattern.pattern());
         final Map<String, TopicDescription> matchedTopicMetadata =
-                getTopicMetadata(adminClient, topicPattern);
-
+                getTopicMetadataProvider().getTopicMetadata(adminClient, 
topicPattern);
         Set<TopicPartition> subscribedTopicPartitions = new HashSet<>();
 
         matchedTopicMetadata.forEach(
diff --git 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
index bb7d71c0..6a0612fe 100644
--- 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
+++ 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
@@ -273,6 +273,57 @@ public class KafkaSourceBuilderTest {
                 .isEqualTo(-1L);
     }
 
+    @Test
+    public void testDefaultCheckSourceIntegrity() {
+        final KafkaSource<String> kafkaSource = getBasicBuilder().build();
+        assertThat(
+                        kafkaSource
+                                .getConfiguration()
+                                
.get(KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED))
+                
.isEqualTo(KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED.defaultValue());
+    }
+
+    @Test
+    public void testCheckSourceIntegrityEnabled() {
+        boolean checkSourceIntegrity = true;
+        final KafkaSource<String> kafkaSource =
+                getBasicBuilder()
+                        .setProperty(
+                                
KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED.key(),
+                                Boolean.toString(checkSourceIntegrity))
+                        .build();
+        assertThat(
+                        kafkaSource
+                                .getConfiguration()
+                                
.get(KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED))
+                .isEqualTo(checkSourceIntegrity);
+    }
+
+    @Test
+    public void 
testTopicIntegrityCheckEnabledWithoutTopicIntegrityAwareSubscriber() {
+        assertThatThrownBy(
+                        () ->
+                                new KafkaSourceBuilder<String>()
+                                        .setBootstrapServers("testServer")
+                                        .setDeserializer(
+                                                
KafkaRecordDeserializationSchema.valueOnly(
+                                                        
StringDeserializer.class))
+                                        .enableTopicIntegrityCheck()
+                                        .setKafkaSubscriber(
+                                                new KafkaSubscriber() {
+                                                    @Override
+                                                    public Set<TopicPartition>
+                                                            
getSubscribedTopicPartitions(
+                                                                    
AdminClient adminClient) {
+                                                        return null;
+                                                    }
+                                                })
+                                        .build())
+                .isInstanceOf(IllegalStateException.class)
+                .hasMessageContaining(
+                        "Topic integrity check is not supported for non 
TopicMetadataSettable subscriber");
+    }
+
     private KafkaSourceBuilder<String> getBasicBuilder() {
         return new KafkaSourceBuilder<String>()
                 .setBootstrapServers("testServer")
diff --git 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializerTest.java
 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializerTest.java
index 7cc66a8c..43a0c5ca 100644
--- 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializerTest.java
+++ 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumStateSerializerTest.java
@@ -23,6 +23,7 @@ import 
org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit;
 import 
org.apache.flink.connector.kafka.source.split.KafkaPartitionSplitSerializer;
 
 import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.Uuid;
 import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
@@ -77,7 +78,20 @@ class KafkaSourceEnumStateSerializerTest {
                                         new SplitAndAssignmentStatus(
                                                 split, 
getAssignmentStatus(split)))
                         .collect(Collectors.toList());
-
+        final Set<SplitAndAssignmentStatus> splitAndAssignmentStatusesSet =
+                splits.stream()
+                        .map(
+                                split ->
+                                        new SplitAndAssignmentStatus(
+                                                split, 
getAssignmentStatus(split)))
+                        .collect(Collectors.toSet());
+        final Map<String, String> trackedTopicIdsByName =
+                splits.stream()
+                        .map(KafkaPartitionSplit::getTopic)
+                        .distinct()
+                        .collect(
+                                Collectors.toMap(
+                                        topic -> topic, topic -> 
Uuid.randomUuid().toString()));
         // Create bytes in the way of KafkaEnumStateSerializer version 0 doing 
serialization
         final byte[] bytesV0 =
                 SerdeUtils.serializeSplitAssignments(
@@ -86,6 +100,13 @@ class KafkaSourceEnumStateSerializerTest {
         final byte[] bytesV1 = 
KafkaSourceEnumStateSerializer.serializeV1(splits);
         final byte[] bytesV2 =
                 
KafkaSourceEnumStateSerializer.serializeV2(splitAndAssignmentStatuses, false);
+        final byte[] bytesV3 =
+                KafkaSourceEnumStateSerializer.serializeV3(
+                        new 
KafkaSourceEnumState(splitAndAssignmentStatusesSet, false));
+        final byte[] bytesV4 =
+                KafkaSourceEnumStateSerializer.serializeV4(
+                        new KafkaSourceEnumState(
+                                splitAndAssignmentStatusesSet, false, 
trackedTopicIdsByName));
 
         // Deserialize above bytes with KafkaEnumStateSerializer version 2 to 
check backward
         // compatibility
@@ -95,6 +116,10 @@ class KafkaSourceEnumStateSerializerTest {
                 new KafkaSourceEnumStateSerializer().deserialize(1, bytesV1);
         final KafkaSourceEnumState kafkaSourceEnumStateV2 =
                 new KafkaSourceEnumStateSerializer().deserialize(2, bytesV2);
+        final KafkaSourceEnumState kafkaSourceEnumStateV3 =
+                new KafkaSourceEnumStateSerializer().deserialize(3, bytesV3);
+        final KafkaSourceEnumState kafkaSourceEnumStateV4 =
+                new KafkaSourceEnumStateSerializer().deserialize(4, bytesV4);
 
         assertThat(kafkaSourceEnumStateV0.assignedSplits())
                 .containsExactlyInAnyOrderElementsOf(splits);
@@ -120,6 +145,21 @@ class KafkaSourceEnumStateSerializerTest {
                 .containsExactlyInAnyOrderElementsOf(
                         splitsByStatus.get(AssignmentStatus.UNASSIGNED));
         
assertThat(kafkaSourceEnumStateV2.initialDiscoveryFinished()).isFalse();
+
+        assertThat(kafkaSourceEnumStateV3.assignedSplits())
+                
.containsExactlyInAnyOrderElementsOf(splitsByStatus.get(AssignmentStatus.ASSIGNED));
+        assertThat(kafkaSourceEnumStateV3.unassignedSplits())
+                .containsExactlyInAnyOrderElementsOf(
+                        splitsByStatus.get(AssignmentStatus.UNASSIGNED));
+        
assertThat(kafkaSourceEnumStateV3.initialDiscoveryFinished()).isFalse();
+
+        assertThat(kafkaSourceEnumStateV4.assignedSplits())
+                
.containsExactlyInAnyOrderElementsOf(splitsByStatus.get(AssignmentStatus.ASSIGNED));
+        assertThat(kafkaSourceEnumStateV4.unassignedSplits())
+                .containsExactlyInAnyOrderElementsOf(
+                        splitsByStatus.get(AssignmentStatus.UNASSIGNED));
+        
assertThat(kafkaSourceEnumStateV4.initialDiscoveryFinished()).isFalse();
+        
assertThat(kafkaSourceEnumStateV4.trackedTopicIdsByName()).isEqualTo(trackedTopicIdsByName);
     }
 
     private static AssignmentStatus getAssignmentStatus(KafkaPartitionSplit 
split) {
diff --git 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java
 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java
index 1a8a0427..cb216aee 100644
--- 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java
+++ 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumeratorTest.java
@@ -947,4 +947,45 @@ public class KafkaSourceEnumeratorTest {
             context.runNextOneTimeCallable();
         }
     }
+
+    @Test
+    public void testCheckSourceIntegrityFromProperties() throws Exception {
+        // Test that properties are used when job configuration doesn't have 
the setting
+        final boolean propertiesCheckSourceIntegrity = true;
+
+        Properties properties = new Properties();
+        properties.setProperty(
+                KafkaSourceOptions.TOPIC_INTEGRITY_CHECK_ENABLED.key(),
+                String.valueOf(propertiesCheckSourceIntegrity));
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(
+                                context,
+                                ENABLE_PERIODIC_PARTITION_DISCOVERY ? 1 : -1,
+                                OffsetsInitializer.earliest(),
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                Collections.emptySet(),
+                                true,
+                                properties)) {
+
+            // Verify that the properties value is used
+            assertThat(propertiesCheckSourceIntegrity)
+                    .isEqualTo(enumerator.topicIntegrityCheckEnabled());
+        }
+    }
+
+    @Test
+    public void testCheckSourceIntegrityDefaultValue() throws Exception {
+        final boolean defaultCheckSourceIntegrity = false;
+        try (MockSplitEnumeratorContext<KafkaPartitionSplit> context =
+                        new MockSplitEnumeratorContext<>(NUM_SUBTASKS);
+                KafkaSourceEnumerator enumerator =
+                        createEnumerator(context, 
DISABLE_PERIODIC_PARTITION_DISCOVERY)) {
+
+            assertThat(defaultCheckSourceIntegrity)
+                    .isEqualTo(enumerator.topicIntegrityCheckEnabled());
+        }
+    }
 }
diff --git 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberTest.java
 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberTest.java
index 79f35889..59680df3 100644
--- 
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberTest.java
+++ 
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberTest.java
@@ -20,9 +20,14 @@ package 
org.apache.flink.connector.kafka.source.enumerator.subscriber;
 
 import org.apache.flink.connector.kafka.lineage.DefaultKafkaDatasetIdentifier;
 import org.apache.flink.connector.kafka.lineage.KafkaDatasetIdentifierProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicIntegrityException;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicIntegrityProvider;
+import 
org.apache.flink.connector.kafka.source.enumerator.metadata.TopicMetadataSettable;
 import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
 
 import org.apache.kafka.clients.admin.AdminClient;
+import org.apache.kafka.clients.admin.TopicDescription;
+import org.apache.kafka.common.KafkaFuture;
 import org.apache.kafka.common.TopicPartition;
 import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
 import org.junit.jupiter.api.AfterAll;
@@ -32,9 +37,11 @@ import org.junit.jupiter.api.parallel.ResourceLock;
 
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
+import java.util.concurrent.ExecutionException;
 import java.util.regex.Pattern;
 
 import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
@@ -46,6 +53,8 @@ import static 
org.assertj.core.api.Assertions.assertThatThrownBy;
 class KafkaSubscriberTest {
     private static final String TOPIC1 = "topic1";
     private static final String TOPIC2 = "pattern-topic";
+    private static final String TOPIC3 = "topic3";
+    private static final List<String> TOPICS = Arrays.asList(TOPIC1, TOPIC3);
     private static final TopicPartition NON_EXISTING_TOPIC = new 
TopicPartition("removed", 0);
     private static AdminClient adminClient;
 
@@ -54,6 +63,7 @@ class KafkaSubscriberTest {
         KafkaSourceTestEnv.setup();
         KafkaSourceTestEnv.createTestTopic(TOPIC1);
         KafkaSourceTestEnv.createTestTopic(TOPIC2);
+        KafkaSourceTestEnv.createTestTopic(TOPIC3);
         adminClient = KafkaSourceTestEnv.getAdminClient();
     }
 
@@ -90,6 +100,69 @@ class KafkaSubscriberTest {
                 
.satisfies(anyCauseMatches(UnknownTopicOrPartitionException.class));
     }
 
+    @Test
+    void testNonExistingTopicWithTopicIntegrity() {
+        final KafkaSubscriber subscriber =
+                KafkaSubscriber.getTopicListSubscriber(
+                        Collections.singletonList(NON_EXISTING_TOPIC.topic()));
+        enableTopicIntegrityCheck(subscriber);
+        assertThatThrownBy(() -> 
subscriber.getSubscribedTopicPartitions(adminClient))
+                .isInstanceOf(TopicIntegrityException.class)
+                .hasMessage("Topic " + NON_EXISTING_TOPIC.topic() + " is 
missing");
+    }
+
+    private static void enableTopicIntegrityCheck(KafkaSubscriber subscriber) {
+        ((TopicMetadataSettable) subscriber)
+                .setTopicMetadataProvider(new TopicIntegrityProvider(new 
HashMap<>()));
+    }
+
+    static String getTopicId(String topicName, AdminClient adminClient)
+            throws ExecutionException, InterruptedException {
+        KafkaFuture<TopicDescription> topicDescriptionFuture =
+                adminClient
+                        .describeTopics(Arrays.asList(topicName))
+                        .topicNameValues()
+                        .get(topicName);
+        return topicDescriptionFuture.get().topicId().toString();
+    }
+
+    public void testSubscriberWithCheckTopicIntegrityEnabled(KafkaSubscriber 
subscriber)
+            throws Throwable {
+        enableTopicIntegrityCheck(subscriber);
+        final Set<TopicPartition> subscribedPartitions =
+                subscriber.getSubscribedTopicPartitions(adminClient);
+        final Set<TopicPartition> expectedSubscribedPartitions =
+                new 
HashSet<>(KafkaSourceTestEnv.getPartitionsForTopics(TOPICS));
+        
assertThat(subscribedPartitions).isEqualTo(expectedSubscribedPartitions);
+
+        // Recreate the environment to simulate topic recreation
+        // (simple recreation won't do because topics are only marked for 
deletion)
+        tearDown();
+        setup();
+        assertThatThrownBy(() -> 
subscriber.getSubscribedTopicPartitions(adminClient))
+                .isInstanceOf(TopicIntegrityException.class)
+                .hasMessageMatching("Topic (" + String.join("|", TOPICS) + ") 
was recreated");
+    }
+
+    @Test
+    public void testTopicSubscriberWithCheckTopicIntegrityEnabled() throws 
Throwable {
+        testSubscriberWithCheckTopicIntegrityEnabled(
+                KafkaSubscriber.getTopicListSubscriber(TOPICS));
+    }
+
+    @Test
+    public void testPatternSubscriberWithCheckTopicIntegrityEnabled() throws 
Throwable {
+        testSubscriberWithCheckTopicIntegrityEnabled(
+                
KafkaSubscriber.getTopicPatternSubscriber(Pattern.compile("topic.*")));
+    }
+
+    @Test
+    public void testPartitionSubscriberWithCheckTopicIntegrityEnabled() throws 
Throwable {
+        testSubscriberWithCheckTopicIntegrityEnabled(
+                KafkaSubscriber.getPartitionSetSubscriber(
+                        new 
HashSet<>(KafkaSourceTestEnv.getPartitionsForTopics(TOPICS))));
+    }
+
     @Test
     void testTopicPatternSubscriber() {
         Pattern pattern = Pattern.compile("pattern.*");
@@ -114,10 +187,10 @@ class KafkaSubscriberTest {
         partitions.remove(new TopicPartition(TOPIC1, 1));
 
         KafkaSubscriber subscriber = 
KafkaSubscriber.getPartitionSetSubscriber(partitions);
+        enableTopicIntegrityCheck(subscriber);
 
         final Set<TopicPartition> subscribedPartitions =
                 subscriber.getSubscribedTopicPartitions(adminClient);
-
         assertThat(subscribedPartitions).isEqualTo(partitions);
         assertThat(((KafkaDatasetIdentifierProvider) 
subscriber).getDatasetIdentifier().get())
                 .isEqualTo(DefaultKafkaDatasetIdentifier.ofTopics(topics));
@@ -129,6 +202,7 @@ class KafkaSubscriberTest {
         final KafkaSubscriber subscriber =
                 KafkaSubscriber.getPartitionSetSubscriber(
                         Collections.singleton(nonExistingPartition));
+        enableTopicIntegrityCheck(subscriber);
 
         assertThatThrownBy(() -> 
subscriber.getSubscribedTopicPartitions(adminClient))
                 .isInstanceOf(RuntimeException.class)

Reply via email to