flink-commits mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From dwysakow...@apache.org
Subject [flink] branch master updated (2ed3b55 -> 6f21603)
Date Fri, 02 Jul 2021 16:01:07 GMT
This is an automated email from the ASF dual-hosted git repository.

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


    from 2ed3b55  [hotfix][streaming] Fix the boundary condition when fetch limited records
from streaming job
     add ce68c4d  [FLINK-21086][runtime][checkpoint] Add EndOfUserRecordsEvent
     add af72dc3  [FLINK-21086][runtime][checkpoint] Let result partition supports waiting
for the downstream tasks to processed all the records
     add 8539005  [FLINK-21086][runtime][checkpoint] StreamTask waits for all the records
get processed by downstream tasks
     add ec74c46  [FLINK-21086][runtime][checkpoint] Downstream tasks response with Ack message
when received EndOfUserRecordsEvent
     add 6f21603  [FLINK-21086][runtime][checkpoint] Make CheckpointBarrierHandler supports
alignment with finished channels

No new revisions were added by this update.

Summary of changes:
 .../flink/state/api/output/BoundedStreamTask.java  |   3 +-
 .../runtime/io/network/NetworkClientHandler.java   |   7 +
 .../io/network/NetworkSequenceViewReader.java      |   3 +
 .../runtime/io/network/PartitionRequestClient.java |   7 +
 ...titionEvent.java => EndOfUserRecordsEvent.java} |  30 ++--
 .../network/api/serialization/EventSerializer.java |   7 +
 .../network/api/writer/ResultPartitionWriter.java  |  13 ++
 .../CreditBasedPartitionRequestClientHandler.java  |  25 +++
 .../CreditBasedSequenceNumberingViewReader.java    |   5 +
 .../runtime/io/network/netty/NettyMessage.java     |  35 ++++
 .../network/netty/NettyPartitionRequestClient.java |   5 +
 .../io/network/netty/PartitionRequestQueue.java    |  14 ++
 .../netty/PartitionRequestServerHandler.java       |  10 +-
 ...edBlockingSubpartitionDirectTransferReader.java |   5 +
 .../BoundedBlockingSubpartitionReader.java         |   5 +
 .../partition/NoOpResultSubpartitionView.java      |   3 +
 .../partition/PipelinedResultPartition.java        |  71 +++++++-
 .../network/partition/PipelinedSubpartition.java   |   4 +
 .../partition/PipelinedSubpartitionView.java       |   5 +
 .../io/network/partition/ResultPartition.java      |  19 +++
 .../network/partition/ResultSubpartitionView.java  |   2 +
 .../partition/SortMergeSubpartitionReader.java     |   5 +
 .../network/partition/consumer/InputChannel.java   |   6 +
 .../io/network/partition/consumer/InputGate.java   |   3 +
 .../partition/consumer/LocalInputChannel.java      |   7 +
 .../partition/consumer/RecoveredInputChannel.java  |  10 ++
 .../partition/consumer/RemoteInputChannel.java     |   8 +
 .../partition/consumer/SingleInputGate.java        |   6 +
 .../network/partition/consumer/UnionInputGate.java |   7 +
 .../partition/consumer/UnknownInputChannel.java    |   6 +
 ...bleNotifyingResultPartitionWriterDecorator.java |  10 ++
 .../runtime/taskmanager/InputGateWithMetrics.java  |   5 +
 .../io/network/TestingPartitionRequestClient.java  |   3 +
 .../api/serialization/EventSerializerTest.java     |   2 +
 .../network/netty/CancelPartitionRequestTest.java  |   3 +
 .../NettyMessageServerSideSerializationTest.java   |   9 +
 .../netty/NettyPartitionRequestClientTest.java     |  38 +++++
 .../netty/PartitionRequestServerHandlerTest.java   |  47 +++++
 .../partition/MockResultPartitionWriter.java       |   8 +
 .../io/network/partition/ResultPartitionTest.java  |  35 ++++
 .../partition/consumer/InputChannelTest.java       |   3 +
 .../partition/consumer/TestInputChannel.java       |   6 +
 .../AbstractAlignedBarrierHandlerState.java        |  13 +-
 ...tractAlternatingAlignedBarrierHandlerState.java |   2 +-
 .../AlternatingCollectingBarriers.java             |  29 ++++
 .../AlternatingCollectingBarriersUnaligned.java    |  17 ++
 .../AlternatingWaitingForFirstBarrier.java         |  10 ++
 ...AlternatingWaitingForFirstBarrierUnaligned.java |  10 ++
 .../io/checkpointing/BarrierHandlerState.java      |   8 +
 .../runtime/io/checkpointing/ChannelState.java     |   5 +
 .../io/checkpointing/CheckpointBarrierHandler.java |  16 +-
 .../io/checkpointing/CheckpointBarrierTracker.java | 185 +++++++++++++++-----
 .../io/checkpointing/CheckpointedInputGate.java    |   8 +-
 .../io/checkpointing/CollectingBarriers.java       |  20 +++
 .../SingleCheckpointBarrierHandler.java            |  98 ++++++++---
 .../io/checkpointing/WaitingForFirstBarrier.java   |  12 ++
 .../streaming/runtime/tasks/SourceStreamTask.java  |   2 +-
 .../runtime/tasks/StreamIterationHead.java         |   3 +-
 .../flink/streaming/runtime/tasks/StreamTask.java  |  41 ++++-
 .../runtime/tasks/mailbox/MailboxProcessor.java    |   8 +
 .../streaming/runtime/io/MockIndexedInputGate.java |   7 +
 .../flink/streaming/runtime/io/MockInputGate.java  |   7 +
 .../AlignedCheckpointsMassiveRandomTest.java       |   3 +
 .../io/checkpointing/AlignedCheckpointsTest.java   |  71 ++++++++
 .../checkpointing/AlternatingCheckpointsTest.java  |  50 ++++++
 .../CheckpointBarrierTrackerTest.java              | 181 ++++++++++++++++++++
 .../UnalignedCheckpointsCancellationTest.java      |   3 +-
 .../io/checkpointing/UnalignedCheckpointsTest.java |  71 +++++++-
 .../checkpointing/ValidatingCheckpointHandler.java |   4 +-
 ...tStreamTaskChainedSourcesCheckpointingTest.java | 139 ++++++++++-----
 .../runtime/tasks/MultipleInputStreamTaskTest.java | 182 +++++++++-----------
 .../runtime/tasks/SourceStreamTaskTest.java        |  71 ++++++++
 .../runtime/tasks/StreamConfigChainer.java         |  73 +++++---
 .../runtime/tasks/StreamMockEnvironment.java       |   7 +-
 .../tasks/StreamTaskMailboxTestHarness.java        |   4 +-
 .../tasks/StreamTaskMailboxTestHarnessBuilder.java |  21 ++-
 .../runtime/tasks/StreamTaskTerminationTest.java   |   3 +-
 .../streaming/runtime/tasks/StreamTaskTest.java    | 190 +++++++++++++--------
 .../runtime/tasks/StreamTaskTestHarness.java       |   2 +-
 .../runtime/tasks/SynchronousCheckpointITCase.java |   3 +-
 .../runtime/tasks/SynchronousCheckpointTest.java   |   3 +-
 .../tasks/TaskCheckpointingBehaviourTest.java      |   3 +-
 .../JobMasterStopWithSavepointITCase.java          |   6 +-
 83 files changed, 1732 insertions(+), 364 deletions(-)
 copy flink-runtime/src/main/java/org/apache/flink/runtime/io/network/api/{EndOfPartitionEvent.java
=> EndOfUserRecordsEvent.java} (62%)

Mime
View raw message