[ 
https://issues.apache.org/jira/browse/FLINK-40836?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40836:
-----------------------------------
    Labels: pull-request-available  (was: )

> 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: 1.14.0, 2.0.2, 2.3.0, 2.2.1, 1.20.5, 2.1.3
>            Reporter: Spoorthi Basu
>            Priority: Major
>              Labels: pull-request-available
>
> 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)

Reply via email to