Package | Description |
---|---|
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.io.network.partition.consumer | |
org.apache.flink.runtime.jobgraph.tasks | |
org.apache.flink.runtime.messages.checkpoint |
This package contains the messages that are sent between
JobMaster and TaskExecutor to coordinate the checkpoint snapshots of the
distributed dataflow. |
org.apache.flink.runtime.taskexecutor.rpc | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.state.api.runtime | |
org.apache.flink.streaming.api.operators | |
org.apache.flink.streaming.api.operators.sort | |
org.apache.flink.streaming.runtime.io | |
org.apache.flink.streaming.runtime.io.checkpointing | |
org.apache.flink.streaming.runtime.io.recovery | |
org.apache.flink.streaming.runtime.tasks |
This package contains classes that realize streaming tasks.
|
Modifier and Type | Method and Description |
---|---|
CheckpointException |
PendingCheckpoint.getFailureCause() |
Modifier and Type | Method and Description |
---|---|
void |
CheckpointCoordinator.abortPendingCheckpoints(CheckpointException exception)
Aborts all the pending checkpoints due to en exception.
|
void |
CheckpointFailureManager.checkFailureCounter(CheckpointException exception,
long checkpointId) |
void |
CheckpointFailureManager.handleJobLevelCheckpointException(CheckpointException exception,
long checkpointId)
Handle job level checkpoint exception with a handler callback.
|
void |
CheckpointFailureManager.handleTaskLevelCheckpointException(CheckpointException exception,
long checkpointId,
ExecutionAttemptID executionAttemptID)
Handle task level checkpoint exception with a handler callback.
|
Modifier and Type | Method and Description |
---|---|
boolean |
CheckpointCoordinator.receiveAcknowledgeMessage(AcknowledgeCheckpoint message,
String taskManagerLocationInfo)
Receives an AcknowledgeCheckpoint message and returns whether the message was associated with
a pending checkpoint.
|
void |
CheckpointCoordinator.reportStats(long id,
ExecutionAttemptID attemptId,
CheckpointMetrics metrics) |
Modifier and Type | Method and Description |
---|---|
void |
Environment.declineCheckpoint(long checkpointId,
CheckpointException checkpointException)
Declines a checkpoint.
|
Modifier and Type | Method and Description |
---|---|
void |
RecoveredInputChannel.checkpointStarted(CheckpointBarrier barrier) |
void |
RemoteInputChannel.checkpointStarted(CheckpointBarrier barrier)
Spills all queued buffers on checkpoint start.
|
void |
LocalInputChannel.checkpointStarted(CheckpointBarrier barrier) |
void |
InputChannel.checkpointStarted(CheckpointBarrier barrier)
Called by task thread when checkpointing is started (e.g., any input channel received
barrier).
|
void |
CheckpointableInput.checkpointStarted(CheckpointBarrier barrier) |
void |
IndexedInputGate.checkpointStarted(CheckpointBarrier barrier) |
protected void |
ChannelStatePersister.startPersisting(long barrierId,
List<Buffer> knownBuffers) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractInvokable.abortCheckpointOnBarrier(long checkpointId,
CheckpointException cause)
Aborts a checkpoint as the result of receiving possibly some checkpoint barriers, but at
least one
CancelCheckpointMarker . |
Modifier and Type | Method and Description |
---|---|
CheckpointException |
SerializedCheckpointException.unwrap() |
Constructor and Description |
---|
DeclineCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
CheckpointException checkpointException) |
SerializedCheckpointException(CheckpointException checkpointException) |
Modifier and Type | Method and Description |
---|---|
void |
RpcCheckpointResponder.declineCheckpoint(JobID jobID,
ExecutionAttemptID executionAttemptID,
long checkpointId,
CheckpointException checkpointException) |
Modifier and Type | Method and Description |
---|---|
void |
CheckpointResponder.declineCheckpoint(JobID jobID,
ExecutionAttemptID executionAttemptID,
long checkpointId,
CheckpointException checkpointException)
Declines the given checkpoint.
|
void |
RuntimeEnvironment.declineCheckpoint(long checkpointId,
CheckpointException checkpointException) |
Modifier and Type | Method and Description |
---|---|
void |
SavepointEnvironment.declineCheckpoint(long checkpointId,
CheckpointException checkpointException) |
Modifier and Type | Method and Description |
---|---|
OperatorSnapshotFutures |
StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.CheckpointedStreamOperator streamOperator,
Optional<InternalTimeServiceManager<?>> timeServiceManager,
String operatorName,
long checkpointId,
long timestamp,
CheckpointOptions checkpointOptions,
CheckpointStreamFactory factory,
boolean isUsingCustomRawKeyedState) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Void> |
SortingDataInput.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Void> |
StreamInputProcessor.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
CompletableFuture<Void> |
StreamTaskNetworkInput.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
CompletableFuture<Void> |
StreamOneInputProcessor.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
CompletableFuture<Void> |
StreamMultipleInputProcessor.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
CompletableFuture<Void> |
StreamTaskInput.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId)
Prepares to spill the in-flight input buffers as checkpoint snapshot.
|
CompletableFuture<Void> |
StreamTaskSourceInput.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
CompletableFuture<Void> |
StreamTwoInputProcessor.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
protected void |
CheckpointBarrierHandler.notifyAbort(long checkpointId,
CheckpointException cause) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Void> |
RescalingStreamTaskNetworkInput.prepareSnapshot(ChannelStateWriter channelStateWriter,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
void |
StreamTask.abortCheckpointOnBarrier(long checkpointId,
CheckpointException cause) |
void |
MultipleInputStreamTask.abortCheckpointOnBarrier(long checkpointId,
CheckpointException cause) |
void |
SubtaskCheckpointCoordinator.abortCheckpointOnBarrier(long checkpointId,
CheckpointException cause,
OperatorChain<?,?> operatorChain) |
Modifier and Type | Method and Description |
---|---|
void |
SubtaskCheckpointCoordinator.initInputsCheckpoint(long id,
CheckpointOptions checkpointOptions)
Initialize new checkpoint.
|
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.