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

Reply via email to