Spoorthi Basu created FLINK-40836:
-------------------------------------
Summary: Restoring with
execution.state-recovery.without-channel-state.checkpoint-id drops operator
coordinator state
Key: FLINK-40836
URL: https://issues.apache.org/jira/browse/FLINK-40836
Project: Flink
Issue Type: Bug
Components: Runtime / Checkpointing
Affects Versions: 2.1.3, 1.20.5, 2.2.1, 2.3.0, 2.0.2, 1.14.0
Reporter: Spoorthi Basu
When a job is restored from the checkpoint whose id is set in
{{execution.state-recovery.without-channel-state.checkpoint-id}} (deprecated
key {{execution.checkpointing.recover-without-channel-state.checkpoint-id}}),
every operator coordinator is reset with no state. For a source built on the
{{Source}} API, {{Source#restoreEnumerator}} is not called and
{{Source#createEnumerator}} is called instead, so the split enumerator starts
as it would for a new job. The subtask state is restored as usual.
The option is documented as a last resort for corrupted in-flight data, and
dropping that data is its intended effect. Dropping the coordinator state is
not: with {{NumberSequenceSource}}, the restored job emits the whole sequence a
second time.
The code is the same in every release since 1.14.0 (checked at each release
tag) and on current master, where I reproduced it (1690c6bed94).
Cause, at master:
* {{CheckpointCoordinator#extractOperatorStates}}
(CheckpointCoordinator.java:1845) replaces every {{OperatorState}} with
{{OperatorState#copyAndDiscardInFlightData()}} when the restored checkpoint's
id equals the configured id. It is called at line 1807, and the result is
passed to {{StateAssignmentOperation}} at line 1819 and to
{{restoreStateToCoordinators}} at line 1838.
* {{OperatorState#copyAndDiscardInFlightData()}} (OperatorState.java:186)
builds a new {{OperatorState}} and copies only the subtask states.
{{coordinatorState}} is not copied.
* {{restoreStateToCoordinators}} (line 2146) therefore passes {{null}} to
{{resetToCheckpoint}}. {{SourceCoordinator#resetToCheckpoint}} treats {{null}}
as "no checkpoint" (SourceCoordinator.java:501), and {{start()}} then calls
{{Source#createEnumerator}} (line 252).
The check is only on the checkpoint id, so unaligned checkpoints do not have to
be enabled for this to happen.
The first version of {{extractOperatorStates}} (FLINK-22684) changed the
subtask states in place and kept the coordinator state. Commit bbc94d750db
("fixup: Create the new operator states instead of changing old one"), part of
the same work, switched to new {{OperatorState}} objects without the
coordinator state, and FLINK-23809 moved that code into
{{OperatorState#copyAndDiscardInFlightData()}}. All three are in 1.14.0.
{{IgnoreInFlightDataITCase}} uses a legacy {{SourceFunction}}, which has no
coordinator, so it does not cover this.
*Reproduction*
MiniCluster, parallelism 1, source chained to a map and a sink, so there is no
network exchange and no in-flight data in either restore. The first run has
unaligned checkpoints enabled and retains checkpoints on cancellation; both
restores run with unaligned checkpoints disabled.
# Run {{env.fromSequence(0, 9_999)}} followed by a map that sleeps 1 ms per
record, and stop it with a savepoint after about 2,000 records. Read the
savepoint's checkpoint id N, e.g.
{{SavepointLoader.loadSavepointMetadata(path).getCheckpointId()}}.
# Restore the same job from that savepoint, without the sleep, and let it
finish. Count the records.
# Restore it again from the same savepoint with
{{execution.state-recovery.without-channel-state.checkpoint-id: N}} and count
the records.
I repeated this with a retained checkpoint (triggered while the job ran, then
the job was cancelled) in place of the savepoint.
*Result*
||Snapshot||Restore without the option||Restore with the option set to N||
|savepoint, N = 2|7,984 records (2,016 to 9,999), each once|17,984 records: the
same 7,984 plus all 10,000 values again|
|retained checkpoint, N = 3|7,991 records (2,009 to 9,999), each once|17,991
records: the same 7,991 plus all 10,000 values again|
The enumerator state in the snapshot was 16 bytes (no splits left to assign).
The fresh enumerator created a split for the full range again, and the reader
read it in addition to its restored split. In a second run the snapshots fell
at slightly different records, and again exactly 10,000 extra records were
emitted.
With a test source whose enumerator checkpoints 5,000 entries (520,016 bytes of
coordinator state), a restore without the option called {{restoreEnumerator}}
with all 5,000 entries. With the option, {{restoreEnumerator}} was not called,
{{createEnumerator}} was called once, and the new enumerator assigned a second
split to a reader that had already restored its own. Savepoint and retained
checkpoint gave the same result.
*Fix*
{{copyAndDiscardInFlightData()}} should carry the coordinator state over to the
new {{OperatorState}}. In-flight data is part of the subtask state only.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)