From d27599b02a9425a2f423700af00514311b3aaaf1 Mon Sep 17 00:00:00 2001 From: johnjcasey <95318300+johnjcasey@users.noreply.github.com> Date: Fri, 1 May 2026 15:56:48 -0400 Subject: [PATCH] Revert "Fix unhandled exception in KafkaIO SDF (#37449) (#37553)" This reverts commit a4cb67621c6a24387c6699564638ee38effd5119. --- .../org/apache/beam/sdk/io/kafka/KafkaIO.java | 12 ++++-- .../kafka/WatchForKafkaTopicPartitions.java | 15 +++---- .../WatchForKafkaTopicPartitionsTest.java | 43 ------------------- 3 files changed, 13 insertions(+), 57 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java index f9bfdd87aaac..518319a38e32 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java @@ -2130,12 +2130,16 @@ public void processElement(OutputReceiver receiver) { } else { for (String topic : topics) { List partitionInfoList = consumer.partitionsFor(topic); - if (partitionInfoList == null || partitionInfoList.isEmpty()) { + if (logTopicVerification == null || !logTopicVerification) { + checkState( + partitionInfoList != null && !partitionInfoList.isEmpty(), + "Could not find any partitions info for topic %s. Please check Kafka configuration and make sure that provided topics exist.", + topic); + } else { LOG.warn( - "Could not find any partitions info for topic {}. Please check Kafka " - + "configuration and make sure that the provided topics exist.", + "Could not find any partitions info for topic {}. Please check Kafka configuration " + + "and make sure that the provided topics exist.", topic); - continue; } for (PartitionInfo p : partitionInfoList) { diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitions.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitions.java index 3184b18267b2..490faafb22fa 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitions.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitions.java @@ -18,6 +18,7 @@ package org.apache.beam.sdk.io.kafka; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects.firstNonNull; +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; import java.util.ArrayList; import java.util.List; @@ -44,8 +45,6 @@ import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; import org.joda.time.Instant; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * A {@link PTransform} for continuously querying Kafka for new partitions, and emitting those @@ -58,7 +57,6 @@ */ class WatchForKafkaTopicPartitions extends PTransform> { - private static final Logger LOG = LoggerFactory.getLogger(WatchForKafkaTopicPartitions.class); private static final Duration DEFAULT_CHECK_DURATION = Duration.standardHours(1); private static final String COUNTER_NAMESPACE = "watch_kafka_topic_partition"; @@ -193,13 +191,10 @@ static List getAllTopicPartitions( if (topics != null && !topics.isEmpty()) { for (String topic : topics) { List partitionInfoList = kafkaConsumer.partitionsFor(topic); - if (partitionInfoList == null || partitionInfoList.isEmpty()) { - LOG.warn( - "Could not find any partitions info for topic {}. Please check Kafka " - + "configuration and make sure that the provided topics exist.", - topic); - continue; - } + checkState( + partitionInfoList != null && !partitionInfoList.isEmpty(), + "Could not find any partitions info for topic %s. Please check Kafka configuration and make sure that provided topics exist.", + topic); for (PartitionInfo partition : partitionInfoList) { current.add(new TopicPartition(topic, partition.partition())); } diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitionsTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitionsTest.java index 30ace6cd86d0..595d040bf403 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitionsTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitionsTest.java @@ -18,12 +18,10 @@ package org.apache.beam.sdk.io.kafka; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -import java.util.Collections; import java.util.Set; import java.util.regex.Pattern; import org.apache.beam.sdk.io.kafka.KafkaMocks.PartitionGrowthMockConsumer; @@ -110,47 +108,6 @@ public void testGetAllTopicPartitionsWithGivenTopics() throws Exception { (input) -> mockConsumer, null, givenTopics, null)); } - @Test - public void testGetAllTopicPartitionsWithNullPartitionInfo() throws Exception { - Set givenTopics = ImmutableSet.of("topic1"); - - Consumer mockConsumer = Mockito.mock(Consumer.class); - when(mockConsumer.partitionsFor("topic1")).thenReturn(null); - assertTrue( - WatchForKafkaTopicPartitions.getAllTopicPartitions( - (input) -> mockConsumer, null, givenTopics, null) - .isEmpty()); - } - - @Test - public void testGetAllTopicPartitionsWithEmptyPartitionInfo() throws Exception { - Set givenTopics = ImmutableSet.of("topic1"); - - Consumer mockConsumer = Mockito.mock(Consumer.class); - when(mockConsumer.partitionsFor("topic1")).thenReturn(Collections.emptyList()); - assertTrue( - WatchForKafkaTopicPartitions.getAllTopicPartitions( - (input) -> mockConsumer, null, givenTopics, null) - .isEmpty()); - } - - @Test - public void testGetAllTopicPartitionsSkipsMissingTopics() throws Exception { - Set givenTopics = ImmutableSet.of("topic1", "topic2"); - - Consumer mockConsumer = Mockito.mock(Consumer.class); - when(mockConsumer.partitionsFor("topic1")).thenReturn(null); - when(mockConsumer.partitionsFor("topic2")) - .thenReturn( - ImmutableList.of( - new PartitionInfo("topic2", 0, null, null, null), - new PartitionInfo("topic2", 1, null, null, null))); - assertEquals( - ImmutableList.of(new TopicPartition("topic2", 0), new TopicPartition("topic2", 1)), - WatchForKafkaTopicPartitions.getAllTopicPartitions( - (input) -> mockConsumer, null, givenTopics, null)); - } - @Test public void testGetAllTopicPartitionsWithGivenPattern() throws Exception { Consumer mockConsumer = Mockito.mock(Consumer.class);