flink-commits mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From ar...@apache.org
Subject [flink] branch master updated (e050a48 -> e396f4a)
Date Wed, 01 Sep 2021 06:28:46 GMT
This is an automated email from the ASF dual-hosted git repository.

arvid pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git.


    from e050a48  [FLINK-23468][test] Do not hide original exception if SSL is not available
     add f01c636  [FLINK-23854][datastream] Expose the restored checkpoint id in ManagedInitializationContext.
     add 809cf1c  [FLINK-23854][datastream] Pass checkpoint id to SinkWriter#snapshotState
and restoredCheckpointId to Sink#InitContext.
     add dab4204  [FLINK-23896][streaming] Implement retrying for failed committables for
Sinks.
     add 39b9bad  [hotfix][connectors/kafka] Add TestLogger to Kafka tests.
     add 9a51c12  [FLINK-23678][tests] Re-enable KafkaSinkITCase and optimize it.
     add beed3a6  [FLINK-23854][connectors/kafka] Abort transactions on close in KafkaWriter.
     add ac604d2  [FLINK-23854][connectors/kafka] Transfer KafkaProducer from writer to committer.
     add a9a1db8  [FLINK-23896][connectors/kafka] Retry committing KafkaCommittables on transient
failures.
     add c3e6cec  [FLINK-23854][connectors/kafka] Add FlinkKafkaInternalProducer#setTransactionalId.
     add b203b49  [FLINK-23854][connectors/kafka] Reliably abort lingering transactions in
Kafka.
     add e396f4a  [FLINK-23854][kafka/connectors] Adding pooling

No new revisions were added by this update.

Summary of changes:
 .../connector/file/sink/writer/FileWriter.java     |   2 +-
 .../connector/file/sink/writer/FileWriterTest.java |   8 +-
 .../connector/jdbc/xa/JdbcXaSinkTestBase.java      |   2 +-
 flink-connectors/flink-connector-kafka/pom.xml     |   6 +
 .../kafka/sink/FlinkKafkaInternalProducer.java     | 112 +++++++-
 .../connector/kafka/sink/KafkaCommittable.java     |  24 +-
 .../kafka/sink/KafkaCommittableSerializer.java     |   2 +-
 .../flink/connector/kafka/sink/KafkaCommitter.java |  80 ++++--
 .../connector/kafka/sink/KafkaTransactionLog.java  | 276 --------------------
 .../flink/connector/kafka/sink/KafkaWriter.java    | 288 ++++++++-------------
 .../connector/kafka/sink/KafkaWriterState.java     |  23 +-
 .../kafka/sink/KafkaWriterStateSerializer.java     |   6 +-
 .../flink/connector/kafka/sink/Recyclable.java}    |  33 ++-
 .../connector/kafka/sink/TransactionAborter.java   | 131 ++++++++++
 .../kafka/sink/TransactionalIdFactory.java         |  42 ---
 .../kafka/table/ReducingUpsertWriter.java          |   4 +-
 .../sink/FlinkKafkaInternalProducerITCase.java     | 123 +++++++++
 .../kafka/sink/KafkaCommittableSerializerTest.java |   6 +-
 ...SerializerTest.java => KafkaCommitterTest.java} |  24 +-
 .../KafkaRecordSerializationSchemaBuilderTest.java |   3 +-
 .../connector/kafka/sink/KafkaSinkITCase.java      | 104 ++++----
 .../connector/kafka/sink/KafkaTransactionLog.java  | 247 ++++++++++++++++++
 .../kafka/sink/KafkaTransactionLogITCase.java      | 138 ++++------
 .../connector/kafka/sink/KafkaWriterITCase.java    | 174 +++++++++++--
 .../kafka/sink/KafkaWriterStateSerializerTest.java |   6 +-
 .../kafka/sink/TransactionIdFactoryTest.java       |  44 +---
 .../kafka/sink/TransactionToAbortCheckerTest.java  |   4 +-
 .../kafka/FlinkKafkaConsumerBaseTest.java          |   6 +
 .../org/apache/flink/api/connector/sink/Sink.java  |  15 +-
 .../flink/api/connector/sink/SinkWriter.java       |  14 +-
 .../PrioritizedOperatorSubtaskState.java           |  44 +++-
 .../checkpoint/StateAssignmentOperation.java       |  12 +-
 .../flink/runtime/executiongraph/Execution.java    |   2 +-
 .../runtime/executiongraph/ExecutionVertex.java    |   4 +-
 .../state/ManagedInitializationContext.java        |  13 +-
 .../state/StateInitializationContextImpl.java      |  19 +-
 .../flink/runtime/state/TaskStateManagerImpl.java  |   7 +-
 .../runtime/state/TaskStateManagerImplTest.java    |   5 +-
 .../flink/runtime/state/TestTaskStateManager.java  |   3 +-
 .../api/operators/StreamOperatorStateContext.java  |  14 +-
 .../api/operators/StreamOperatorStateHandler.java  |   5 +-
 .../operators/StreamTaskStateInitializerImpl.java  |  15 +-
 .../operators/sink/AbstractCommitterHandler.java   |  42 ++-
 .../sink/AbstractStreamingCommitterHandler.java    |  49 ++--
 .../operators/sink/BatchCommitterHandler.java      |  22 +-
 .../runtime/operators/sink/CommitRetrier.java      |  84 ++++++
 .../runtime/operators/sink/CommitterHandler.java   |  11 +
 .../runtime/operators/sink/CommitterOperator.java  |  43 +--
 .../operators/sink/CommitterOperatorFactory.java   |   3 +-
 .../operators/sink/ForwardCommittingHandler.java   |  12 +-
 .../sink/GlobalBatchCommitterHandler.java          |  21 +-
 .../sink/GlobalStreamingCommitterHandler.java      |  27 +-
 .../operators/sink/NoopCommitterHandler.java       |   9 +
 .../runtime/operators/sink/SinkOperator.java       |  46 +++-
 .../operators/sink/SinkWriterStateHandler.java     |  13 +-
 .../sink/StatefulSinkWriterStateHandler.java       |   8 +-
 .../sink/StatelessSinkWriterStateHandler.java      |   5 +-
 .../operators/sink/StreamingCommitterHandler.java  |  16 +-
 .../runtime/tasks/TestProcessingTimeService.java   |   4 +
 .../api/operators/SourceOperatorTest.java          |   2 +-
 .../StateInitializationContextImplTest.java        |   4 +-
 .../StreamTaskStateInitializerImplTest.java        |   4 +-
 .../utils/MockFunctionInitializationContext.java   |   7 +
 .../source/SourceOperatorEventTimeTest.java        |   2 +-
 .../operators/sink/BatchCommitterHandlerTest.java  |  10 +-
 .../runtime/operators/sink/CommitRetrierTest.java  | 122 +++++++++
 .../sink/GlobalBatchCommitterHandlerTest.java      |  11 +-
 .../sink/GlobalStreamingCommitterHandlerTest.java  |  35 ++-
 .../operators/sink/SinkWriterOperatorTest.java     |  11 +-
 .../sink/StreamingCommitterHandlerTest.java        |  16 +-
 .../streaming/runtime/operators/sink/TestSink.java |  38 ++-
 .../runtime/tasks/RestoreStreamTaskTest.java       |  41 ++-
 .../streaming/runtime/tasks/StreamTaskTest.java    |   6 +
 73 files changed, 1815 insertions(+), 999 deletions(-)
 delete mode 100644 flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLog.java
 copy flink-connectors/flink-connector-kafka/src/{test/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializerTest.java
=> main/java/org/apache/flink/connector/kafka/sink/Recyclable.java} (56%)
 create mode 100644 flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionAborter.java
 create mode 100644 flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/FlinkKafkaInternalProducerITCase.java
 copy flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/{KafkaWriterStateSerializerTest.java
=> KafkaCommitterTest.java} (61%)
 create mode 100644 flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLog.java
 create mode 100644 flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/operators/sink/CommitRetrier.java
 create mode 100644 flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/operators/sink/CommitRetrierTest.java

Mime
View raw message