This is an automated email from the ASF dual-hosted git repository.
1996fanrui pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
from 12bd86b104c [FLINK-40165][ci] Add docs check to Github Actions CI stage
new c71751b5da6 [FLINK-39523][checkpoint] RecoveryCheckpointTrigger with
NOT_READY/NO_OP and a barrier-inserting in-memory implementation
new 77ecb4dddc0 [FLINK-39523][checkpoint] ChannelState: dispatch
checkpoint start through the recovery trigger
new ed201228225 [FLINK-39523] Thread the trigger through barrier-handler
construction and the StreamTask lifecycle
new b9cadc24b7d [FLINK-39523][test] Restore during-recovery flag
randomization; re-enable filtering ITCase
The 4 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:
.../checkpoint/channel/FetchedChannelState.java | 67 ++++++++++
.../channel/FetchedChannelStateDrainer.java | 84 ++++++++++++
.../channel/RecoveryCheckpointTrigger.java | 57 ++++++++
.../channel/SequentialChannelStateReader.java | 12 +-
.../channel/SequentialChannelStateReaderImpl.java | 47 +++++--
.../partition/consumer/RecoveredInputChannel.java | 19 +--
.../AlternatingCollectingBarriers.java | 5 +-
...AlternatingWaitingForFirstBarrierUnaligned.java | 4 +-
.../runtime/io/checkpointing/ChannelState.java | 27 ++++
.../io/checkpointing/InputProcessorUtil.java | 32 ++++-
.../SingleCheckpointBarrierHandler.java | 10 +-
.../runtime/tasks/MultipleInputStreamTask.java | 3 +-
.../runtime/tasks/OneInputStreamTask.java | 3 +-
.../flink/streaming/runtime/tasks/StreamTask.java | 125 +++++++++++++++---
.../runtime/tasks/TwoInputStreamTask.java | 3 +-
.../consumer/RecoveredInputChannelTest.java | 15 +--
.../runtime/io/checkpointing/ChannelStateTest.java | 147 +++++++++++++++++++++
.../checkpointing/TestBarrierHandlerFactory.java | 2 +
.../streaming/util/TestStreamEnvironment.java | 6 +-
.../RecoveredStateFilteringLargeRecordITCase.java | 5 -
20 files changed, 602 insertions(+), 71 deletions(-)
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/FetchedChannelState.java
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/FetchedChannelStateDrainer.java
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/RecoveryCheckpointTrigger.java
create mode 100644
flink-runtime/src/test/java/org/apache/flink/streaming/runtime/io/checkpointing/ChannelStateTest.java