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)
