This is an automated email from the ASF dual-hosted git repository.
leonard pushed a change to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
from 33d00224 [FLINK-36630][Connectors/Kafka] Wrap consumer.position in
retryOnWakeup (#133)
new 2fdb66d5 [FLINK-36848][Connectors/Kafka] Remove Deprecated class and
relevant tests.
new 9996019c [FLINK-36648][Connectors/Kafka] Bump Flink version to 2.0.0
new c2805458 [hotfix][build] Add stable version to weekly CI
The 3 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:
.github/workflows/push_pr.yml | 11 +-
.github/workflows/weekly.yml | 7 +-
.../flink/tests/util/kafka/SmokeKafkaITCase.java | 13 +-
.../kafka/test/base/CustomWatermarkExtractor.java | 53 -
.../kafka/test/base/KafkaExampleUtil.java | 17 +-
.../kafka/test/base/RollingAdditionMapper.java | 55 -
.../flink/streaming/kafka/test/KafkaExample.java | 2 +-
.../4b58b35e-f9cd-43dc-a664-7af4fa8ec2d0 | 0
.../5b7ce6b8-e525-400c-935f-81a09bc7f0fe | 0
.../6182d789-a081-4f26-b3f4-24a22bc1f248 | 0
.../8511d84b-cbaa-4b54-9e3e-895926935dd7 | 0
.../86dfd459-67a9-4b26-9b5c-0b0bbf22681a | 50 +-
.../e0624cac-4ea1-4bf8-879a-ecedb41ce334 | 2 +-
.../f5cd467c-4694-4798-9e9a-cf7946b31265 | 0
flink-connector-kafka/pom.xml | 6 +
.../kafka/sink/ExactlyOnceKafkaWriter.java | 4 +-
.../KafkaRecordSerializationSchemaBuilder.java | 15 -
.../flink/connector/kafka/sink/KafkaSink.java | 6 +-
.../connector/kafka/sink/KafkaSinkBuilder.java | 14 -
.../flink/connector/kafka/sink/KafkaWriter.java | 8 +-
.../kafka/sink/TwoPhaseCommittingStatefulSink.java | 23 +-
.../kafka/source/reader/KafkaSourceReader.java | 2 +-
.../KafkaDeserializationSchemaWrapper.java | 67 -
.../KafkaRecordDeserializationSchema.java | 20 -
.../reader/fetcher/KafkaSourceFetcherManager.java | 2 +-
.../connectors/kafka/FlinkKafkaConsumer.java | 342 ---
.../connectors/kafka/FlinkKafkaConsumerBase.java | 1228 ---------
.../connectors/kafka/FlinkKafkaErrorCode.java | 33 -
.../connectors/kafka/FlinkKafkaException.java | 50 -
.../connectors/kafka/FlinkKafkaProducer.java | 1920 --------------
.../connectors/kafka/KafkaContextAware.java | 58 -
.../kafka/KafkaDeserializationSchema.java | 86 -
.../connectors/kafka/KafkaSerializationSchema.java | 63 -
.../connectors/kafka/config/OffsetCommitMode.java | 43 -
.../connectors/kafka/config/OffsetCommitModes.java | 53 -
.../kafka/internals/AbstractFetcher.java | 620 -----
.../internals/AbstractPartitionDiscoverer.java | 248 --
.../kafka/internals/ClosableBlockingQueue.java | 501 ----
.../connectors/kafka/internals/ExceptionProxy.java | 124 -
.../internals/FlinkKafkaInternalProducer.java | 423 ---
.../connectors/kafka/internals/Handover.java | 213 --
.../kafka/internals/KafkaCommitCallback.java | 45 -
.../kafka/internals/KafkaConsumerThread.java | 565 ----
.../KafkaDeserializationSchemaWrapper.java | 71 -
.../connectors/kafka/internals/KafkaFetcher.java | 268 --
.../kafka/internals/KafkaPartitionDiscoverer.java | 115 -
.../internals/KafkaSerializationSchemaWrapper.java | 113 -
.../kafka/internals/KafkaShuffleFetcher.java | 306 ---
.../kafka/internals/KafkaTopicPartition.java | 142 -
.../internals/KafkaTopicPartitionAssigner.java | 62 -
.../kafka/internals/KafkaTopicPartitionLeader.java | 108 -
.../kafka/internals/KafkaTopicPartitionState.java | 133 -
...aTopicPartitionStateWithWatermarkGenerator.java | 100 -
.../kafka/internals/KafkaTopicsDescriptor.java | 92 -
.../internals/KeyedSerializationSchemaWrapper.java | 59 -
.../SourceContextWatermarkOutputAdapter.java | 51 -
.../kafka/internals/TransactionalIdsGenerator.java | 101 -
.../metrics/KafkaConsumerMetricConstants.java | 57 -
.../internals/metrics/KafkaMetricWrapper.java | 43 -
.../kafka/partitioner/FlinkFixedPartitioner.java | 11 +-
.../kafka/partitioner/FlinkKafkaPartitioner.java | 35 -
.../kafka/shuffle/FlinkKafkaShuffle.java | 407 ---
.../kafka/shuffle/FlinkKafkaShuffleConsumer.java | 94 -
.../kafka/shuffle/FlinkKafkaShuffleProducer.java | 230 --
.../kafka/shuffle/StreamKafkaShuffleSink.java | 44 -
.../table/DynamicKafkaDeserializationSchema.java | 20 +-
.../kafka/table/KafkaConnectorOptionsUtil.java | 27 +-
.../connectors/kafka/table/KafkaDynamicSink.java | 2 +-
.../connectors/kafka/table/KafkaDynamicSource.java | 24 +-
.../kafka/table/KafkaDynamicTableFactory.java | 6 +-
.../connectors/kafka/table/ReducingUpsertSink.java | 21 +-
.../JSONKeyValueDeserializationSchema.java | 19 +-
.../serialization/KeyedDeserializationSchema.java | 56 -
.../serialization/KeyedSerializationSchema.java | 61 -
...TypeInformationKeyValueSerializationSchema.java | 206 --
.../dynamic/source/DynamicKafkaSourceITTest.java | 22 +-
.../kafka/sink/ExactlyOnceKafkaWriterITCase.java | 13 +
.../sink/FlinkKafkaInternalProducerITCase.java | 13 +
.../KafkaRecordSerializationSchemaBuilderTest.java | 6 +-
.../connector/kafka/sink/KafkaSinkITCase.java | 24 +-
.../kafka/sink/KafkaTransactionLogITCase.java | 13 +
.../sink/KafkaWriterFaultToleranceITCase.java | 13 +
.../connector/kafka/sink/KafkaWriterITCase.java | 14 +
.../connector/kafka/sink/KafkaWriterTestBase.java | 15 -
.../sink/internal/ProducerPoolImplITCase.java | 14 +
.../connector/kafka/source/KafkaSourceITCase.java | 10 +-
.../kafka/source/KafkaSourceLegacyITCase.java | 172 --
.../KafkaRecordDeserializationSchemaTest.java | 33 +-
.../kafka/testutils/SimpleCollector.java} | 30 +-
.../kafka/FlinkKafkaConsumerBaseMigrationTest.java | 436 ----
.../kafka/FlinkKafkaConsumerBaseTest.java | 1523 -----------
.../connectors/kafka/FlinkKafkaConsumerITCase.java | 129 -
.../kafka/FlinkKafkaInternalProducerITCase.java | 262 --
.../connectors/kafka/FlinkKafkaProducerITCase.java | 820 ------
.../kafka/FlinkKafkaProducerMigrationTest.java | 81 -
.../FlinkKafkaProducerStateSerializerTest.java | 148 --
.../connectors/kafka/FlinkKafkaProducerTest.java | 174 --
.../JSONKeyValueDeserializationSchemaTest.java | 25 +-
.../connectors/kafka/KafkaConsumerTestBase.java | 2734 --------------------
.../streaming/connectors/kafka/KafkaITCase.java | 380 ---
.../connectors/kafka/KafkaMigrationTestBase.java | 171 --
.../kafka/KafkaProducerAtLeastOnceITCase.java | 43 -
.../kafka/KafkaProducerExactlyOnceITCase.java | 38 -
.../connectors/kafka/KafkaProducerTestBase.java | 401 ---
.../kafka/KafkaSerializerUpgradeTest.java | 173 --
.../kafka/KafkaShortRetentionTestBase.java | 229 --
.../connectors/kafka/KafkaTestEnvironment.java | 81 -
.../connectors/kafka/KafkaTestEnvironmentImpl.java | 90 -
.../NextTransactionalIdHintSerializerTest.java | 56 -
.../kafka/internals/AbstractFetcherTest.java | 279 --
.../internals/AbstractFetcherWatermarksTest.java | 499 ----
.../internals/AbstractPartitionDiscovererTest.java | 561 ----
.../kafka/internals/ClosableBlockingQueueTest.java | 616 -----
.../kafka/internals/KafkaTopicPartitionTest.java | 57 -
.../kafka/internals/KafkaTopicsDescriptorTest.java | 66 -
.../metrics/KafkaMetricMutableWrapperTest.java | 5 -
.../shuffle/KafkaShuffleExactlyOnceITCase.java | 218 --
.../kafka/shuffle/KafkaShuffleITCase.java | 543 ----
.../kafka/shuffle/KafkaShuffleTestBase.java | 313 ---
.../kafka/table/KafkaChangelogTableITCase.java | 4 +-
.../kafka/table/KafkaDynamicTableFactoryTest.java | 33 +-
.../connectors/kafka/table/KafkaTableITCase.java | 14 +-
.../connectors/kafka/table/KafkaTableTestBase.java | 8 +-
.../kafka/table/ReducingUpsertWriterTest.java | 4 +-
.../table/UpsertKafkaDynamicTableFactoryTest.java | 155 +-
.../connectors/kafka/testutils/DataGenerators.java | 196 --
.../kafka/testutils/FailingIdentityMapper.java | 8 +-
.../connectors/kafka/testutils/IntegerSource.java | 9 +-
.../kafka/testutils/TestPartitionDiscoverer.java | 138 -
.../kafka/testutils/TestSourceContext.java | 4 +-
.../kafka/testutils/Tuple2FlinkPartitioner.java | 4 +-
.../testutils/TypeSerializerUpgradeTestBase.java | 584 -----
.../kafka/testutils/ValidatingExactlyOnceSink.java | 10 +-
.../datastream/connectors/tests/test_kafka.py | 47 -
pom.xml | 8 +-
135 files changed, 417 insertions(+), 22428 deletions(-)
delete mode 100644
flink-connector-kafka-e2e-tests/flink-streaming-kafka-test-base/src/main/java/org/apache/flink/streaming/kafka/test/base/CustomWatermarkExtractor.java
delete mode 100644
flink-connector-kafka-e2e-tests/flink-streaming-kafka-test-base/src/main/java/org/apache/flink/streaming/kafka/test/base/RollingAdditionMapper.java
delete mode 100644
flink-connector-kafka/archunit-violations/4b58b35e-f9cd-43dc-a664-7af4fa8ec2d0
delete mode 100644
flink-connector-kafka/archunit-violations/5b7ce6b8-e525-400c-935f-81a09bc7f0fe
delete mode 100644
flink-connector-kafka/archunit-violations/6182d789-a081-4f26-b3f4-24a22bc1f248
delete mode 100644
flink-connector-kafka/archunit-violations/8511d84b-cbaa-4b54-9e3e-895926935dd7
delete mode 100644
flink-connector-kafka/archunit-violations/f5cd467c-4694-4798-9e9a-cf7946b31265
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/deserializer/KafkaDeserializationSchemaWrapper.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumerBase.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaErrorCode.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaException.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaContextAware.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaDeserializationSchema.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaSerializationSchema.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/config/OffsetCommitMode.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/config/OffsetCommitModes.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcher.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractPartitionDiscoverer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/ClosableBlockingQueue.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/ExceptionProxy.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/FlinkKafkaInternalProducer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/Handover.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaCommitCallback.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaConsumerThread.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaDeserializationSchemaWrapper.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaFetcher.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaPartitionDiscoverer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaSerializationSchemaWrapper.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaShuffleFetcher.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartition.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartitionAssigner.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartitionLeader.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartitionState.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartitionStateWithWatermarkGenerator.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicsDescriptor.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/KeyedSerializationSchemaWrapper.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/SourceContextWatermarkOutputAdapter.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/TransactionalIdsGenerator.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/metrics/KafkaConsumerMetricConstants.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/metrics/KafkaMetricWrapper.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/partitioner/FlinkKafkaPartitioner.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffleConsumer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffleProducer.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/StreamKafkaShuffleSink.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/util/serialization/KeyedDeserializationSchema.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/util/serialization/KeyedSerializationSchema.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/streaming/util/serialization/TypeInformationKeyValueSerializationSchema.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceLegacyITCase.java
copy
flink-connector-kafka/src/test/java/org/apache/flink/{KafkaAssertjConfiguration.java
=> connector/kafka/testutils/SimpleCollector.java} (60%)
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumerBaseMigrationTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumerBaseTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumerITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaInternalProducerITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerMigrationTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerStateSerializerTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducerTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaConsumerTestBase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaMigrationTestBase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaProducerAtLeastOnceITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaProducerExactlyOnceITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaProducerTestBase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaSerializerUpgradeTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaShortRetentionTestBase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/NextTransactionalIdHintSerializerTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcherTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcherWatermarksTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractPartitionDiscovererTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/ClosableBlockingQueueTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicPartitionTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/KafkaTopicsDescriptorTest.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/shuffle/KafkaShuffleExactlyOnceITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/shuffle/KafkaShuffleITCase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/shuffle/KafkaShuffleTestBase.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/testutils/DataGenerators.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/testutils/TestPartitionDiscoverer.java
delete mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/testutils/TypeSerializerUpgradeTestBase.java