You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@flink.apache.org by dw...@apache.org on 2020/07/08 12:23:20 UTC
[flink] branch release-1.11 updated: [FLINK-15414] catch the right
KafkaException for binding errors
This is an automated email from the ASF dual-hosted git repository.
dwysakowicz pushed a commit to branch release-1.11
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/release-1.11 by this push:
new fcd4c58 [FLINK-15414] catch the right KafkaException for binding errors
fcd4c58 is described below
commit fcd4c58c966b3df4150ed9f761273ae608055475
Author: Yuan Mei <yu...@gmail.com>
AuthorDate: Thu Jun 25 19:15:19 2020 +0800
[FLINK-15414] catch the right KafkaException for binding errors
The right exception should be `org.apache.kafka.common.KafkaException` instead of `kafka.common.KafkaException`
At least the "numTries" exception should be thrown instead of the binding exception.
---
.../flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java | 2 +-
1 file changed, 1 insertion(+), 1 deletion(-)
diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java
index f2c1767..40fefff 100644
--- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java
+++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java
@@ -26,7 +26,6 @@ import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkKafkaPartiti
import org.apache.flink.streaming.util.serialization.KeyedSerializationSchema;
import org.apache.flink.util.NetUtils;
-import kafka.common.KafkaException;
import kafka.metrics.KafkaMetricsReporter;
import kafka.server.KafkaConfig;
import kafka.server.KafkaServer;
@@ -39,6 +38,7 @@ import org.apache.kafka.clients.admin.TopicDescription;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
+import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.network.ListenerName;
import org.apache.kafka.common.security.auth.SecurityProtocol;