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 96dcf63f553a734940b5f76d161687c20efd3f91 Author: Aljoscha Krettek <[email protected]> AuthorDate: Wed Jun 5 11:19:14 2019 +0200 [hotfix] Remove unused method from Kafka Test Environments --- .../connectors/kafka/KafkaTestEnvironmentImpl.java | 8 -------- .../connectors/kafka/KafkaTestEnvironmentImpl.java | 15 --------------- .../connectors/kafka/KafkaTestEnvironmentImpl.java | 5 ----- .../connectors/kafka/KafkaTestEnvironmentImpl.java | 5 ----- .../streaming/connectors/kafka/KafkaTestEnvironment.java | 6 ------ .../connectors/kafka/KafkaTestEnvironmentImpl.java | 15 --------------- 6 files changed, 54 deletions(-) diff --git a/flink-connectors/flink-connector-kafka-0.10/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connectors/flink-connector-kafka-0.10/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java index 9e51aac..64649ee 100644 --- a/flink-connectors/flink-connector-kafka-0.10/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java +++ b/flink-connectors/flink-connector-kafka-0.10/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java @@ -166,14 +166,6 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { } @Override - public <T> DataStreamSink<T> writeToKafkaWithTimestamps(DataStream<T> stream, String topic, KeyedSerializationSchema<T> serSchema, Properties props) { - FlinkKafkaProducer010<T> prod = new FlinkKafkaProducer010<>(topic, serSchema, props); - prod.setFlushOnCheckpoint(true); - prod.setWriteTimestampToKafka(true); - return stream.addSink(prod); - } - - @Override public KafkaOffsetHandler createOffsetHandler() { return new KafkaOffsetHandlerImpl(); } diff --git a/flink-connectors/flink-connector-kafka-0.11/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connectors/flink-connector-kafka-0.11/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java index ae237f6..23cd57e 100644 --- a/flink-connectors/flink-connector-kafka-0.11/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java +++ b/flink-connectors/flink-connector-kafka-0.11/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java @@ -296,21 +296,6 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { } @Override - public <T> DataStreamSink<T> writeToKafkaWithTimestamps(DataStream<T> stream, String topic, KeyedSerializationSchema<T> serSchema, Properties props) { - FlinkKafkaProducer011<T> prod = new FlinkKafkaProducer011<>( - topic, - serSchema, - props, - Optional.of(new FlinkFixedPartitioner<>()), - producerSemantic, - FlinkKafkaProducer011.DEFAULT_KAFKA_PRODUCERS_POOL_SIZE); - - prod.setWriteTimestampToKafka(true); - - return stream.addSink(prod); - } - - @Override public KafkaOffsetHandler createOffsetHandler() { return new KafkaOffsetHandlerImpl(); } diff --git a/flink-connectors/flink-connector-kafka-0.8/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connectors/flink-connector-kafka-0.8/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java index 15c83e4..6836afb 100644 --- a/flink-connectors/flink-connector-kafka-0.8/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java +++ b/flink-connectors/flink-connector-kafka-0.8/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java @@ -160,11 +160,6 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { } @Override - public <T> DataStreamSink<T> writeToKafkaWithTimestamps(DataStream<T> stream, String topic, KeyedSerializationSchema<T> serSchema, Properties props) { - throw new UnsupportedOperationException(); - } - - @Override public KafkaOffsetHandler createOffsetHandler() { return new KafkaOffsetHandlerImpl(); } diff --git a/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java index 7a5c9e5..9d03319 100644 --- a/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java +++ b/flink-connectors/flink-connector-kafka-0.9/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java @@ -151,11 +151,6 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { } @Override - public <T> DataStreamSink<T> writeToKafkaWithTimestamps(DataStream<T> stream, String topic, KeyedSerializationSchema<T> serSchema, Properties props) { - throw new UnsupportedOperationException(); - } - - @Override public KafkaOffsetHandler createOffsetHandler() { return new KafkaOffsetHandlerImpl(); } diff --git a/flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironment.java b/flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironment.java index a787f27..336372d 100644 --- a/flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironment.java +++ b/flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironment.java @@ -154,12 +154,6 @@ public abstract class KafkaTestEnvironment { KeyedSerializationSchema<T> serSchema, Properties props, FlinkKafkaPartitioner<T> partitioner); - public abstract <T> DataStreamSink<T> writeToKafkaWithTimestamps( - DataStream<T> stream, - String topic, - KeyedSerializationSchema<T> serSchema, - Properties props); - // -- offset handlers /** 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 d6a7705..710c753 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 @@ -274,21 +274,6 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { } @Override - public <T> DataStreamSink<T> writeToKafkaWithTimestamps(DataStream<T> stream, String topic, KeyedSerializationSchema<T> serSchema, Properties props) { - FlinkKafkaProducer<T> prod = new FlinkKafkaProducer<T>( - topic, - serSchema, - props, - Optional.of(new FlinkFixedPartitioner<>()), - producerSemantic, - FlinkKafkaProducer.DEFAULT_KAFKA_PRODUCERS_POOL_SIZE); - - prod.setWriteTimestampToKafka(true); - - return stream.addSink(prod); - } - - @Override public KafkaOffsetHandler createOffsetHandler() { return new KafkaOffsetHandlerImpl(); }
