This is an automated email from the ASF dual-hosted git repository.
arvid pushed a change to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
from 45df794b [hotfix] Remove unused method
new d43fcbd1 [FLINK-34554] Introduce transaction strategies
new fb280151 [FLINK-34554] Adding listing abort strategy
new 2d4f402f [FLINK-34554] Adding pooling name strategy
new 91b048db [FLINK-34554] Adding strategies to table API
new 2b199b0c [hotfix] Harden KafkaWriterFaultToleranceITCase
The 5 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:
.../flink/tests/util/kafka/KafkaSinkE2ECase.java | 25 +-
.../c0d94764-76a0-4c50-b617-70b1754c4612 | 19 +-
.../kafka/sink/ExactlyOnceKafkaWriter.java | 184 ++++++++++--
.../flink/connector/kafka/sink/KafkaCommitter.java | 51 +++-
.../flink/connector/kafka/sink/KafkaSink.java | 8 +-
.../connector/kafka/sink/KafkaSinkBuilder.java | 34 ++-
.../connector/kafka/sink/KafkaWriterState.java | 75 ++++-
.../kafka/sink/KafkaWriterStateSerializer.java | 44 ++-
.../connector/kafka/sink/TransactionAborter.java | 128 --------
.../kafka/sink/TransactionNamingStrategy.java | 93 ++++++
.../connector/kafka/sink/internal/Backchannel.java | 3 +
.../sink/internal/FlinkKafkaInternalProducer.java | 7 +-
.../kafka/sink/internal/ProducerPool.java | 8 +
.../kafka/sink/internal/ProducerPoolImpl.java | 106 +++++--
.../kafka/sink/internal/ReadableBackchannel.java | 3 +
.../TransactionAbortStrategyContextImpl.java | 124 ++++++++
.../internal/TransactionAbortStrategyImpl.java | 212 +++++++++++++
.../TransactionNamingStrategyContextImpl.java | 93 ++++++
.../internal/TransactionNamingStrategyImpl.java | 113 +++++++
.../kafka/sink/internal/TransactionOwnership.java | 192 ++++++++++++
.../sink/internal/TransactionalIdFactory.java | 13 +
.../kafka/sink/internal/WritableBackchannel.java | 3 +
.../subscriber/KafkaSubscriberUtils.java | 59 ----
.../subscriber/PartitionSetSubscriber.java | 2 +-
.../enumerator/subscriber/TopicListSubscriber.java | 2 +-
.../subscriber/TopicPatternSubscriber.java | 2 +-
.../flink/connector/kafka/util/AdminUtils.java | 137 +++++++++
.../DynamicKafkaRecordSerializationSchema.java | 52 +++-
.../kafka/table/KafkaConnectorOptions.java | 25 ++
.../connectors/kafka/table/KafkaDynamicSink.java | 11 +-
.../kafka/table/KafkaDynamicTableFactory.java | 15 +-
.../table/UpsertKafkaDynamicTableFactory.java | 8 +-
.../kafka/sink/ExactlyOnceKafkaWriterITCase.java | 101 ++++++-
.../connector/kafka/sink/KafkaCommitterTest.java | 15 +-
.../connector/kafka/sink/KafkaSinkBuilderTest.java | 4 +-
.../connector/kafka/sink/KafkaSinkITCase.java | 331 +++++++++++++++++----
.../flink/connector/kafka/sink/KafkaSinkTest.java | 40 ++-
.../sink/KafkaWriterFaultToleranceITCase.java | 11 +-
.../kafka/sink/KafkaWriterStateSerializerTest.java | 15 +-
.../connector/kafka/sink/KafkaWriterTestBase.java | 33 +-
.../sink/internal/ProducerPoolImplITCase.java | 53 +++-
.../internal/TransactionAbortStrategyImplTest.java | 235 +++++++++++++++
.../sink/internal/TransactionIdFactoryTest.java} | 39 ++-
.../sink/internal/TransactionOwnershipTest.java | 165 ++++++++++
.../sink/testutils/KafkaSinkExternalContext.java | 12 +-
.../testutils/KafkaSinkExternalContextFactory.java | 11 +-
.../kafka/table/KafkaDynamicTableFactoryTest.java | 59 +++-
.../connectors/kafka/table/KafkaTableITCase.java | 51 ++++
.../table/UpsertKafkaDynamicTableFactoryTest.java | 54 +++-
49 files changed, 2649 insertions(+), 431 deletions(-)
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionAborter.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionNamingStrategy.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionAbortStrategyContextImpl.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionAbortStrategyImpl.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionNamingStrategyContextImpl.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionNamingStrategyImpl.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionOwnership.java
delete mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberUtils.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/util/AdminUtils.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/internal/TransactionAbortStrategyImplTest.java
copy
flink-connector-kafka/src/{main/java/org/apache/flink/connector/kafka/sink/internal/ReadableBackchannel.java
=>
test/java/org/apache/flink/connector/kafka/sink/internal/TransactionIdFactoryTest.java}
(50%)
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/internal/TransactionOwnershipTest.java