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

Reply via email to