This is an automated email from the ASF dual-hosted git repository. aljoscha pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 7a5ec023f3cb2d9ac619916bf920009d43f5c449 Author: yanghua <[email protected]> AuthorDate: Sat Oct 13 14:16:56 2018 +0800 [FLINK-9697] Rename KafkaConsumerCallBridge to KafkaConsumerCallBridge09 --- .../connectors/kafka/internal/KafkaConsumerCallBridge010.java | 2 +- .../flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java | 4 ++-- .../{KafkaConsumerCallBridge.java => KafkaConsumerCallBridge09.java} | 2 +- .../streaming/connectors/kafka/internal/KafkaConsumerThread.java | 4 ++-- .../streaming/connectors/kafka/internal/KafkaConsumerThreadTest.java | 2 +- 5 files changed, 7 insertions(+), 7 deletions(-) diff --git a/flink-connectors/flink-connector-kafka-0.10/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge010.java b/flink-connectors/flink-connector-kafka-0.10/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge010.java index 5815bfa..579e9d1 100644 --- a/flink-connectors/flink-connector-kafka-0.10/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge010.java +++ b/flink-connectors/flink-connector-kafka-0.10/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge010.java @@ -35,7 +35,7 @@ import java.util.List; * <p>Because of that, we need two versions whose compiled code goes against different method signatures. */ @Internal -public class KafkaConsumerCallBridge010 extends KafkaConsumerCallBridge { +public class KafkaConsumerCallBridge010 extends KafkaConsumerCallBridge09 { @Override public void assignPartitions(KafkaConsumer<?, ?> consumer, List<TopicPartition> topicPartitions) throws Exception { diff --git a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java index dcc67d5..852d56f 100644 --- a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java +++ b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java @@ -202,8 +202,8 @@ public class Kafka09Fetcher<T> extends AbstractFetcher<T, TopicPartition> { return "Kafka 0.9 Fetcher"; } - protected KafkaConsumerCallBridge createCallBridge() { - return new KafkaConsumerCallBridge(); + protected KafkaConsumerCallBridge09 createCallBridge() { + return new KafkaConsumerCallBridge09(); } // ------------------------------------------------------------------------ diff --git a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge.java b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge09.java similarity index 98% rename from flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge.java rename to flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge09.java index b789633..775b7d7 100644 --- a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge.java +++ b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerCallBridge09.java @@ -36,7 +36,7 @@ import java.util.List; * are compiled against different dependencies. */ @Internal -public class KafkaConsumerCallBridge { +public class KafkaConsumerCallBridge09 { public void assignPartitions(KafkaConsumer<?, ?> consumer, List<TopicPartition> topicPartitions) throws Exception { consumer.assign(topicPartitions); diff --git a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java index 38e8a41..7003c3f 100644 --- a/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java +++ b/flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java @@ -81,7 +81,7 @@ public class KafkaConsumerThread extends Thread { private final ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>> unassignedPartitionsQueue; /** The indirections on KafkaConsumer methods, for cases where KafkaConsumer compatibility is broken. */ - private final KafkaConsumerCallBridge consumerCallBridge; + private final KafkaConsumerCallBridge09 consumerCallBridge; /** The maximum number of milliseconds to wait for a fetch batch. */ private final long pollTimeout; @@ -125,7 +125,7 @@ public class KafkaConsumerThread extends Thread { Handover handover, Properties kafkaProperties, ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>> unassignedPartitionsQueue, - KafkaConsumerCallBridge consumerCallBridge, + KafkaConsumerCallBridge09 consumerCallBridge, String threadName, long pollTimeout, boolean useMetrics, diff --git a/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThreadTest.java b/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThreadTest.java index 8dbf871..5da4976 100644 --- a/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThreadTest.java +++ b/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThreadTest.java @@ -716,7 +716,7 @@ public class KafkaConsumerThreadTest { handover, new Properties(), unassignedPartitionsQueue, - new KafkaConsumerCallBridge(), + new KafkaConsumerCallBridge09(), "test-kafka-consumer-thread", 0, false,
