This is an automated email from the ASF dual-hosted git repository.
aljoscha pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git.
from f195a6e [FLINK-4052] Use non-local unreachable ip:port in
testReturnLocalHostAddressUsingHeuristics
new 3a727c0 [FLINK-9697] Rename KafkaTableSource to KafkaTableSourceBase
new c4e74b5 [FLINK-9697] Rename KafkaTableSink to KafkaTableSinkBase
new 7a5ec02 [FLINK-9697] Rename KafkaConsumerCallBridge to
KafkaConsumerCallBridge09
new b4a47c3 [FLINK-9697] Upgrade slf4j from 1.7.7 to 1.7.15
new 328c005 [FLINK-9697] Remove usage of deprecated code in
KafkaConsumerTestBase
new d7a2234 [FLINK-9697] Add connector version value for kafka 2.0
new 2e3e820 [FLINK-9697] Add new kafka connector module
The 7 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../connectors/kafka/Kafka010JsonTableSink.java | 6 +-
.../connectors/kafka/Kafka010TableSink.java | 2 +-
.../connectors/kafka/Kafka010TableSource.java | 2 +-
.../kafka/Kafka010TableSourceSinkFactory.java | 4 +-
.../kafka/internal/KafkaConsumerCallBridge010.java | 2 +-
.../kafka/Kafka010AvroTableSourceTest.java | 2 +-
.../kafka/Kafka010JsonTableSinkTest.java | 4 +-
.../kafka/Kafka010JsonTableSourceTest.java | 2 +-
.../kafka/Kafka010TableSourceSinkFactoryTest.java | 4 +-
.../connectors/kafka/FlinkKafkaProducer011.java | 8 +-
.../connectors/kafka/Kafka011TableSink.java | 2 +-
.../connectors/kafka/Kafka011TableSource.java | 2 +-
.../kafka/Kafka011TableSourceSinkFactory.java | 4 +-
...Wrapper.java => KafkaMetricMutableWrapper.java} | 4 +-
.../kafka/Kafka011AvroTableSourceTest.java | 2 +-
.../kafka/Kafka011JsonTableSourceTest.java | 2 +-
.../kafka/Kafka011TableSourceSinkFactoryTest.java | 4 +-
.../connectors/kafka/Kafka08JsonTableSink.java | 8 +-
.../connectors/kafka/Kafka08TableSink.java | 2 +-
.../connectors/kafka/Kafka08TableSource.java | 2 +-
.../kafka/Kafka08TableSourceSinkFactory.java | 4 +-
.../kafka/Kafka08AvroTableSourceTest.java | 2 +-
.../connectors/kafka/Kafka08JsonTableSinkTest.java | 4 +-
.../kafka/Kafka08JsonTableSourceTest.java | 2 +-
.../kafka/Kafka08TableSourceSinkFactoryTest.java | 4 +-
.../connectors/kafka/Kafka09JsonTableSink.java | 8 +-
.../connectors/kafka/Kafka09TableSink.java | 2 +-
.../connectors/kafka/Kafka09TableSource.java | 2 +-
.../kafka/Kafka09TableSourceSinkFactory.java | 4 +-
.../connectors/kafka/internal/Kafka09Fetcher.java | 4 +-
...lBridge.java => KafkaConsumerCallBridge09.java} | 2 +-
.../kafka/internal/KafkaConsumerThread.java | 4 +-
.../kafka/Kafka09AvroTableSourceTest.java | 2 +-
.../connectors/kafka/Kafka09JsonTableSinkTest.java | 4 +-
.../kafka/Kafka09JsonTableSourceTest.java | 2 +-
.../kafka/Kafka09TableSourceSinkFactoryTest.java | 4 +-
.../kafka/internal/KafkaConsumerThreadTest.java | 2 +-
.../connectors/kafka/KafkaAvroTableSource.java | 4 +-
.../connectors/kafka/KafkaJsonTableSink.java | 4 +-
.../connectors/kafka/KafkaJsonTableSource.java | 4 +-
...KafkaTableSink.java => KafkaTableSinkBase.java} | 16 +-
...aTableSource.java => KafkaTableSourceBase.java} | 26 +-
.../kafka/KafkaTableSourceSinkFactoryBase.java | 6 +-
.../flink/table/descriptors/KafkaValidator.java | 4 +-
.../kafka/KafkaAvroTableSourceTestBase.java | 2 +-
.../connectors/kafka/KafkaConsumerTestBase.java | 46 +-
.../kafka/KafkaJsonTableSourceFactoryTestBase.java | 2 +-
...stBase.java => KafkaTableSinkBaseTestBase.java} | 14 +-
.../kafka/KafkaTableSourceBuilderTestBase.java | 32 +-
.../kafka/KafkaTableSourceSinkFactoryTestBase.java | 12 +-
flink-connectors/flink-connector-kafka/pom.xml | 318 ++++++++++
.../connectors/kafka/FlinkKafkaConsumer.java | 321 ++++++++++
.../connectors/kafka/FlinkKafkaErrorCode.java | 29 +
.../connectors/kafka/FlinkKafkaException.java | 46 ++
.../connectors/kafka/FlinkKafkaProducer.java} | 318 +++++-----
.../connectors/kafka/KafkaTableSink.java} | 49 +-
.../connectors/kafka/KafkaTableSource.java} | 61 +-
.../kafka/KafkaTableSourceSinkFactory.java} | 57 +-
.../kafka/internal/FlinkKafkaInternalProducer.java | 270 +++++++++
.../connectors/kafka/internal/Handover.java | 0
.../kafka/internal/KafkaConsumerThread.java | 15 +-
.../connectors/kafka/internal/KafkaFetcher.java} | 122 ++--
.../kafka/internal/KafkaPartitionDiscoverer.java | 106 ++++
.../kafka/internal/TransactionalIdsGenerator.java | 0
.../metrics/KafkaMetricMutableWrapper.java} | 6 +-
.../org.apache.flink.table.factories.TableFactory | 16 +
.../src/main/resources/log4j.properties | 0
.../kafka/FlinkKafkaInternalProducerITCase.java | 114 ++++
.../connectors/kafka/FlinkKafkaProducerITCase.java | 647 +++++++++++++++++++++
.../FlinkKafkaProducerStateSerializerTest.java | 108 ++++
.../streaming/connectors/kafka/KafkaITCase.java | 353 +++++++++++
.../kafka/KafkaProducerAtLeastOnceITCase.java} | 31 +-
.../kafka/KafkaProducerExactlyOnceITCase.java | 57 ++
.../kafka/KafkaTableSourceSinkFactoryTest.java | 95 +++
.../connectors/kafka/KafkaTestEnvironmentImpl.java | 440 ++++++++++++++
.../src/test/resources/log4j-test.properties | 0
flink-connectors/pom.xml | 1 +
pom.xml | 2 +-
78 files changed, 3374 insertions(+), 504 deletions(-)
copy
flink-connectors/flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/metrics/{KafkaMetricMuttableWrapper.java
=> KafkaMetricMutableWrapper.java} (90%)
rename
flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/{KafkaConsumerCallBridge.java
=> KafkaConsumerCallBridge09.java} (98%)
rename
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/{KafkaTableSink.java
=> KafkaTableSinkBase.java} (94%)
rename
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/{KafkaTableSource.java
=> KafkaTableSourceBase.java} (97%)
rename
flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/{KafkaTableSinkTestBase.java
=> KafkaTableSinkBaseTestBase.java} (90%)
create mode 100644 flink-connectors/flink-connector-kafka/pom.xml
create mode 100644
flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumer.java
create mode 100644
flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaErrorCode.java
create mode 100644
flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaException.java
copy
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducer011.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducer.java}
(77%)
copy
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/Kafka011TableSink.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaTableSink.java}
(50%)
copy
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/Kafka011TableSource.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaTableSource.java}
(63%)
copy
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/Kafka011TableSourceSinkFactory.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaTableSourceSinkFactory.java}
(52%)
create mode 100644
flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/FlinkKafkaInternalProducer.java
copy flink-connectors/{flink-connector-kafka-0.9 =>
flink-connector-kafka}/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Handover.java
(100%)
copy flink-connectors/{flink-connector-kafka-0.9 =>
flink-connector-kafka}/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java
(96%)
copy
flink-connectors/{flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaFetcher.java}
(69%)
create mode 100644
flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaPartitionDiscoverer.java
copy flink-connectors/{flink-connector-kafka-0.11 =>
flink-connector-kafka}/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/TransactionalIdsGenerator.java
(100%)
copy
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/metrics/KafkaMetricMuttableWrapper.java
=>
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/metrics/KafkaMetricMutableWrapper.java}
(86%)
create mode 100644
flink-connectors/flink-connector-kafka/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
copy flink-connectors/{flink-connector-kafka-0.9 =>
flink-connector-kafka}/src/main/resources/log4j.properties (100%)
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaInternalProducerITCase.java
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerITCase.java
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerStateSerializerTest.java
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaITCase.java
rename
flink-connectors/{flink-connector-kafka-0.11/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/metrics/KafkaMetricMuttableWrapper.java
=>
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaProducerAtLeastOnceITCase.java}
(52%)
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaProducerExactlyOnceITCase.java
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTableSourceSinkFactoryTest.java
create mode 100644
flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java
copy flink-connectors/{flink-connector-kafka-0.8 =>
flink-connector-kafka}/src/test/resources/log4j-test.properties (100%)