This is an automated email from the ASF dual-hosted git repository.

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


    from 8a2a3969701 [FLINK-27690][python][connector/pulsar][docs] Add pulsar 
example and documentation
     new 10b7afae742 [FLINK-27251][checkpoint] Refactor the barrier alignment 
timer and default priority sequence number
     new dd8f4e26033 [FLINK-27251][checkpoint] Timeout aligned to unaligned 
checkpoint barrier in the output buffers
     new e68d679e8ac [FLINK-27251][checkpoint] Refactor the close() and 
cancel() of SubtaskCheckpointCoordinator

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:
 .../runtime/checkpoint/CheckpointOptions.java      |   6 +-
 .../channel/ChannelStateWriteRequest.java          |  61 +++++-
 .../checkpoint/channel/ChannelStateWriter.java     |  24 +++
 .../checkpoint/channel/ChannelStateWriterImpl.java |  18 ++
 .../io/network/api/writer/RecordWriter.java        |   9 +
 .../network/api/writer/ResultPartitionWriter.java  |   7 +
 .../partition/BoundedBlockingSubpartition.java     |  11 ++
 .../partition/BufferWritingResultPartition.java    |  15 ++
 .../network/partition/PipelinedSubpartition.java   | 206 ++++++++++++++++++++-
 .../io/network/partition/ResultSubpartition.java   |   5 +
 .../partition/SortMergeResultPartition.java        |  11 ++
 .../ChannelStateWriteRequestDispatcherTest.java    |  19 ++
 .../checkpoint/channel/MockChannelStateWriter.java |  19 ++
 .../network/partition/InputChannelTestUtils.java   |   2 +
 .../partition/MockResultPartitionWriter.java       |   7 +
 .../partition/PipelinedSubpartitionTest.java       | 143 ++++++++++++++
 .../runtime/state/ChannelPersistenceITCase.java    |  34 +++-
 .../streaming/runtime/io/RecordWriterOutput.java   |   9 +
 ...tractAlternatingAlignedBarrierHandlerState.java |   1 +
 .../AlternatingCollectingBarriers.java             |   1 +
 .../io/checkpointing/BarrierAlignmentUtil.java     |  74 ++++++++
 .../io/checkpointing/CheckpointBarrierHandler.java |   5 -
 .../io/checkpointing/InputProcessorUtil.java       |  23 +--
 .../SingleCheckpointBarrierHandler.java            |  23 +--
 .../streaming/runtime/tasks/OperatorChain.java     |  13 ++
 .../flink/streaming/runtime/tasks/StreamTask.java  |  25 ++-
 .../tasks/SubtaskCheckpointCoordinator.java        |   3 +
 .../tasks/SubtaskCheckpointCoordinatorImpl.java    | 123 +++++++++---
 .../checkpointing/AlternatingCheckpointsTest.java  |  11 +-
 .../checkpointing/TestBarrierHandlerFactory.java   |  10 +-
 .../MockSubtaskCheckpointCoordinatorBuilder.java   |   6 +-
 .../tasks/SubtaskCheckpointCoordinatorTest.java    |  63 ++++++-
 .../tasks/TestSubtaskCheckpointCoordinator.java    |   3 +
 33 files changed, 876 insertions(+), 114 deletions(-)
 create mode 100644 
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/checkpointing/BarrierAlignmentUtil.java

Reply via email to