Spoorthi Basu created FLINK-40837:
-------------------------------------
Summary: SavepointWriter#changeOperatorIdentifier drops the
operator coordinator state
Key: FLINK-40837
URL: https://issues.apache.org/jira/browse/FLINK-40837
Project: Flink
Issue Type: Bug
Components: API / State Processor
Affects Versions: 2.1.3, 1.20.5, 2.2.1, 2.3.0, 2.0.2, 1.17.0
Reporter: Spoorthi Basu
When a savepoint is rewritten with
{{SavepointWriter#changeOperatorIdentifier}}, the operator's subtask state is
moved to the new identifier but its coordinator state is not. For a source
built on the {{Source}} API this is the split enumerator state. A job restored
from the rewritten savepoint starts its source with {{Source#createEnumerator}}
instead of {{Source#restoreEnumerator}}, as if it were a new job, while the
readers restore their splits as usual. The restore succeeds with
{{allowNonRestoredState=false}}. With {{NumberSequenceSource}}, the restored
job emits the whole sequence a second time.
The code is the same in every release since 1.17.0 (checked at each release
tag) and on current master, where I reproduced it (1690c6bed94).
Cause, at master:
* {{SavepointWriter.CheckpointMetadataCheckpointMetadataMapFunction}}
(SavepointWriter.java:390) applies the identifier changes by calling
{{OperatorState#copyWithNewIDs}} (line 423).
* {{OperatorState#copyWithNewIDs}} (OperatorState.java:178) builds a new
{{OperatorState}} and copies only the subtask states. {{coordinatorState}} is
not copied. FLINK-40836 is the same omission in
{{OperatorState#copyAndDiscardInFlightData}}.
For comparison, when a savepoint is restored into a job that no longer has the
operator, {{Checkpoints#loadAndValidateCheckpoint}} rejects it if that operator
has coordinator state, unless {{allowNonRestoredState}} is set
(Checkpoints.java:266). The rewrite drops the same state without any such check.
The method was added as {{copyWithNewOperatorID}} by FLINK-29457 (1.17.0),
which introduced {{changeOperatorIdentifier}}, and renamed to
{{copyWithNewIDs}} by FLINK-36001 (2.0.0). Neither version copies the
coordinator state.
*Reproduction*
MiniCluster, parallelism 1. All operators have explicit uids.
# Run {{env.fromSequence(0, 9_999).uid("source-v1")}} followed by a map that
sleeps 1 ms per record, and stop it with a savepoint after about 2,000 records.
# Rewrite the savepoint twice with {{SavepointWriter.fromExistingSavepoint(env,
path)}}: once unchanged, once with
{{changeOperatorIdentifier(OperatorIdentifier.forUid("source-v1"),
OperatorIdentifier.forUid("source-v2"))}}.
# Restore the unchanged rewrite into the job with uid {{source-v1}}, and the
changed rewrite into the same job with uid {{source-v2}}, both with
{{allowNonRestoredState=false}} and without the sleep. Let each finish and
count the records.
*Result*
Coordinator state of the source operator in each savepoint:
||Savepoint||NumberSequenceSource||Test source, 5,000 enumerator entries||
|original (uid source-v1)|16 bytes|520,016 bytes|
|rewritten, uid unchanged|16 bytes|520,016 bytes|
|rewritten, uid source-v1 to source-v2|none|none|
The subtask state was present under the new operator ID in the changed rewrite.
Records emitted by the restored {{fromSequence(0, 9_999)}} job:
||Restored from||Records||
|unchanged rewrite (uid source-v1)|7,988 (2,012 to 9,999), each once|
|changed rewrite (uid source-v2)|17,988: the same 7,988 plus all 10,000 values
again|
A second run gave 7,991 and 17,991, again exactly 10,000 extra records. With
the test source, the changed rewrite led to {{createEnumerator}} instead of
{{restoreEnumerator}}, as in FLINK-40836.
*Fix*
{{copyWithNewIDs}} should carry the coordinator state over to the new
{{OperatorState}} along with the subtask states. Coordinators are matched by
operator ID on restore, so the state then reaches the coordinator of the
renamed operator.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)