Skip to content

Commit c782e1c

Browse files
authored
Revert "Fix unhandled exception in KafkaIO SDF (#37449) (#37553)" (#38361)
This reverts commit a4cb676.
1 parent 7920213 commit c782e1c

3 files changed

Lines changed: 13 additions & 57 deletions

File tree

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2130,12 +2130,16 @@ public void processElement(OutputReceiver<KafkaSourceDescriptor> receiver) {
21302130
} else {
21312131
for (String topic : topics) {
21322132
List<PartitionInfo> partitionInfoList = consumer.partitionsFor(topic);
2133-
if (partitionInfoList == null || partitionInfoList.isEmpty()) {
2133+
if (logTopicVerification == null || !logTopicVerification) {
2134+
checkState(
2135+
partitionInfoList != null && !partitionInfoList.isEmpty(),
2136+
"Could not find any partitions info for topic %s. Please check Kafka configuration and make sure that provided topics exist.",
2137+
topic);
2138+
} else {
21342139
LOG.warn(
2135-
"Could not find any partitions info for topic {}. Please check Kafka "
2136-
+ "configuration and make sure that the provided topics exist.",
2140+
"Could not find any partitions info for topic {}. Please check Kafka configuration "
2141+
+ "and make sure that the provided topics exist.",
21372142
topic);
2138-
continue;
21392143
}
21402144

21412145
for (PartitionInfo p : partitionInfoList) {

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitions.java

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
package org.apache.beam.sdk.io.kafka;
1919

2020
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects.firstNonNull;
21+
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
2122

2223
import java.util.ArrayList;
2324
import java.util.List;
@@ -44,8 +45,6 @@
4445
import org.checkerframework.checker.nullness.qual.Nullable;
4546
import org.joda.time.Duration;
4647
import org.joda.time.Instant;
47-
import org.slf4j.Logger;
48-
import org.slf4j.LoggerFactory;
4948

5049
/**
5150
* A {@link PTransform} for continuously querying Kafka for new partitions, and emitting those
@@ -58,7 +57,6 @@
5857
*/
5958
class WatchForKafkaTopicPartitions extends PTransform<PBegin, PCollection<KafkaSourceDescriptor>> {
6059

61-
private static final Logger LOG = LoggerFactory.getLogger(WatchForKafkaTopicPartitions.class);
6260
private static final Duration DEFAULT_CHECK_DURATION = Duration.standardHours(1);
6361
private static final String COUNTER_NAMESPACE = "watch_kafka_topic_partition";
6462

@@ -193,13 +191,10 @@ static List<TopicPartition> getAllTopicPartitions(
193191
if (topics != null && !topics.isEmpty()) {
194192
for (String topic : topics) {
195193
List<PartitionInfo> partitionInfoList = kafkaConsumer.partitionsFor(topic);
196-
if (partitionInfoList == null || partitionInfoList.isEmpty()) {
197-
LOG.warn(
198-
"Could not find any partitions info for topic {}. Please check Kafka "
199-
+ "configuration and make sure that the provided topics exist.",
200-
topic);
201-
continue;
202-
}
194+
checkState(
195+
partitionInfoList != null && !partitionInfoList.isEmpty(),
196+
"Could not find any partitions info for topic %s. Please check Kafka configuration and make sure that provided topics exist.",
197+
topic);
203198
for (PartitionInfo partition : partitionInfoList) {
204199
current.add(new TopicPartition(topic, partition.partition()));
205200
}

sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/WatchForKafkaTopicPartitionsTest.java

Lines changed: 0 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,10 @@
1818
package org.apache.beam.sdk.io.kafka;
1919

2020
import static org.junit.Assert.assertEquals;
21-
import static org.junit.Assert.assertTrue;
2221
import static org.mockito.Mockito.never;
2322
import static org.mockito.Mockito.verify;
2423
import static org.mockito.Mockito.when;
2524

26-
import java.util.Collections;
2725
import java.util.Set;
2826
import java.util.regex.Pattern;
2927
import org.apache.beam.sdk.io.kafka.KafkaMocks.PartitionGrowthMockConsumer;
@@ -110,47 +108,6 @@ public void testGetAllTopicPartitionsWithGivenTopics() throws Exception {
110108
(input) -> mockConsumer, null, givenTopics, null));
111109
}
112110

113-
@Test
114-
public void testGetAllTopicPartitionsWithNullPartitionInfo() throws Exception {
115-
Set<String> givenTopics = ImmutableSet.of("topic1");
116-
117-
Consumer<byte[], byte[]> mockConsumer = Mockito.mock(Consumer.class);
118-
when(mockConsumer.partitionsFor("topic1")).thenReturn(null);
119-
assertTrue(
120-
WatchForKafkaTopicPartitions.getAllTopicPartitions(
121-
(input) -> mockConsumer, null, givenTopics, null)
122-
.isEmpty());
123-
}
124-
125-
@Test
126-
public void testGetAllTopicPartitionsWithEmptyPartitionInfo() throws Exception {
127-
Set<String> givenTopics = ImmutableSet.of("topic1");
128-
129-
Consumer<byte[], byte[]> mockConsumer = Mockito.mock(Consumer.class);
130-
when(mockConsumer.partitionsFor("topic1")).thenReturn(Collections.emptyList());
131-
assertTrue(
132-
WatchForKafkaTopicPartitions.getAllTopicPartitions(
133-
(input) -> mockConsumer, null, givenTopics, null)
134-
.isEmpty());
135-
}
136-
137-
@Test
138-
public void testGetAllTopicPartitionsSkipsMissingTopics() throws Exception {
139-
Set<String> givenTopics = ImmutableSet.of("topic1", "topic2");
140-
141-
Consumer<byte[], byte[]> mockConsumer = Mockito.mock(Consumer.class);
142-
when(mockConsumer.partitionsFor("topic1")).thenReturn(null);
143-
when(mockConsumer.partitionsFor("topic2"))
144-
.thenReturn(
145-
ImmutableList.of(
146-
new PartitionInfo("topic2", 0, null, null, null),
147-
new PartitionInfo("topic2", 1, null, null, null)));
148-
assertEquals(
149-
ImmutableList.of(new TopicPartition("topic2", 0), new TopicPartition("topic2", 1)),
150-
WatchForKafkaTopicPartitions.getAllTopicPartitions(
151-
(input) -> mockConsumer, null, givenTopics, null));
152-
}
153-
154111
@Test
155112
public void testGetAllTopicPartitionsWithGivenPattern() throws Exception {
156113
Consumer<byte[], byte[]> mockConsumer = Mockito.mock(Consumer.class);

0 commit comments

Comments
 (0)