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

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

> Periodic materialization permanently stops after taking a native-format 
> savepoint when the changelog state backend   is enabled
> -------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40302
>                 URL: https://issues.apache.org/jira/browse/FLINK-40302
>             Project: Flink
>          Issue Type: Bug
>          Components: Runtime / State Backends
>    Affects Versions: 2.0.0, 1.20.0, 1.20.1, 2.1.0, 2.2.0, 2.3.0
>            Reporter: 최원용
>            Priority: Critical
>              Labels: pull-request-available
>
> h3. Summary
> With the changelog state backend enabled ({{{}state.changelog.enabled: 
> true{}}}), triggering a *single intermediate savepoint in NATIVE format* 
> permanently stops periodic materialization.
> The job continues running and checkpointing normally, and no exception is 
> reported. However, the delegated RocksDB backend is never flushed again: 
> MemTables grow to their write-buffer limit, tombstones are never removed by 
> compaction, and the changelog on which the job depends for recovery continues 
> to grow. The only way to recover is to restart and restore the job.
> The only visible indication is the following INFO log, which repeats on every 
> materialization attempt:
> {code:java}
> INFO  o.a.f.state.changelog.ChangelogKeyedStateBackend - materialization:819 
> not confirmed or failed or cancelled,
> skip trigger new one.
> {code}
> This was observed in production on Flink 1.20.1. A job remained in this state 
> for more than seven hours, with materialization stopped and RocksDB 
> flush/compaction stalled, until the job was restarted.
> h3. Reproduction
>  # Run any streaming job using RocksDB with {{{}state.changelog.enabled: 
> true{}}}. The default periodic materialization interval of 10 minutes is 
> sufficient.
>  # Let the job complete at least one successful materialization cycle.
>  # Trigger an intermediate savepoint in NATIVE format while leaving the job 
> running:
> {code:bash}
> flink savepoint <jobId> --type native
> {code}
>  # After the savepoint completes, every subsequent materialization attempt 
> logs the following message indefinitely:
> {code:java}
> materialization:N not confirmed or failed or cancelled, skip trigger new one.
> {code}
> The delegated {{RocksDBKeyedStateBackend}} is never flushed again.
> h3. Root cause
> Three individually correct mechanisms interact to cause this issue. All code 
> references below point to master at 
> [69cdb54|https://github.com/apache/flink/commit/69cdb54e953e]. The relevant 
> code paths are unchanged on the current master branch; the {{release-1.20}} 
> and {{release-2.2}} branches are functionally identical in these paths.
> *1. {{nativeSavepoint()}} consumes a materialization ID and relies on 
> {{notifyCheckpointComplete()}} to confirm it.*
> See [ChangelogKeyedStateBackend#nativeSavepoint(), 
> L447-L496|https://github.com/apache/flink/blob/69cdb54e953e/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java#L447-L496]:
> {code:java}
> long materializationID = materializedId++; // consumes an ID
> ...
> materializationIdByCheckpointId.put(
>         checkpointId, materializationID);  // waits for a notification
> {code}
> *2. However, intermediate savepoints never receive 
> {{{}notifyCheckpointComplete(){}}}.*
> This is intentional behavior introduced by 
> [FLIP-203|https://cwiki.apache.org/confluence/display/FLINK/FLIP-203%3A+Incremental+savepoints]
>  in Flink 1.15: intermediate savepoints do not commit any side effects.
> [CheckpointCoordinator#cleanupAfterCompletedCheckpoint(), 
> L1416-L1441|https://github.com/apache/flink/blob/69cdb54e953e/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java#L1416-L1441]
>  sends acknowledge messages only under the following condition:
> {code:java}
> if (!props.isSavepoint() || props.isSynchronous())
> {code}
> Therefore, a savepoint that does not stop the job never sends a completion 
> notification to the tasks.
> *3. The skip guard in {{initMaterialization()}} then blocks all subsequent 
> materializations.*
> The guard was added by hotfix 
> [85f32d6bcb|https://github.com/apache/flink/commit/85f32d6bcb31] as a 
> follow-up to FLINK-29913 / [PR 
> #22669|https://github.com/apache/flink/pull/22669].
> See [ChangelogKeyedStateBackend#initMaterialization(), 
> L853-L864|https://github.com/apache/flink/blob/69cdb54e953e/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java#L853-L864]:
> {code:java}
> if (lastConfirmedMaterializationId < materializedId - 1
>         && lastFailedMaterializationId < materializedId - 1) {
>     LOG.info(
>             "materialization:{} not confirmed or failed or cancelled, "
>                     + "skip trigger new one.",
>             materializedId - 1);
>     return Optional.empty();
> }
> {code}
> The following sequence demonstrates the issue, starting from a steady state 
> in which {{lastConfirmedMaterializationId = 7}} and {{{}materializedId = 
> 8{}}}:
> ||Event||{{lastConfirmedMaterializationId}}||{{materializedId}}||Guard: 
> {{lastConfirmed < materializedId -
> 1}}||
> |Steady state|7|8|{{7 < 7}} = false → materialize|
> |NATIVE savepoint consumes ID 8|7|9|Completion notification never arrives|
> |Every subsequent {{initMaterialization()}} call|7|9|{{7 < 8}} = *true* → 
> skip indefinitely|
> *Why it never recovers automatically:*
>  * The savepoint succeeds, so 
> {{handleMaterializationFailureOrCancellation()}} is never called. Therefore, 
> {{lastFailedMaterializationId}} remains {{{}-1{}}}.
>  * Because materialization has stopped, the backend's base remains at ID 7. 
> Every subsequent checkpoint maps its checkpoint ID to materialization ID 7, 
> so {{notifyCheckpointComplete()}} can never advance 
> {{lastConfirmedMaterializationId}} beyond 7.
>  * The savepoint's pending entry is silently removed. The next checkpoint 
> notification executes the following code at L655:
> {code:java}
> materializationIdByCheckpointId.headMap(checkpointId, true).clear();
> {code}
> This removes the savepoint's mapping without confirming it. After that, no 
> state remains that could ever confirm materialization ID 8.
>  * The only code path that resets these counters is 
> {{{}completeRestore(){}}}, which requires a full job restore.
> Regular checkpoints cannot cause this issue. {{snapshot()}} only references 
> the current materialized base and never increments {{{}materializedId{}}}.
> The only two producers of materialization IDs are:
>  * Periodic materialization, for which the confirmation loop is complete.
>  * {{{}nativeSavepoint(){}}}, for which the confirmation path is severed.
> h3. Affected versions
> This deadlock requires both of the following:
>  # The skip guard introduced by 
> [85f32d6bcb|https://github.com/apache/flink/commit/85f32d6bcb31], which is 
> present in Flink 1.16.3, 1.17.2, 1.18.0, and later versions.
>  # Native savepoint support in the changelog state backend, introduced by 
> FLINK-29802 / [b62de02f|https://github.com/apache/flink/commit/b62de02f], 
> which is present starting with Flink 1.20.0.
> Therefore, all releases starting with {*}Flink 1.20.0{*}, including all 2.x 
> releases, are affected. The relevant code on the current master branch 
> remains unchanged.
> h3. Proposed fix
> The guard protects {{SharedStateRegistry}} continuity for shared files, as 
> discussed in [PR #22669|https://github.com/apache/flink/pull/22669]. However, 
> a native savepoint uses {{NO_SHARING}} and does not participate in shared 
> state. Furthermore, by design under FLIP-203, an intermediate savepoint can 
> never be confirmed through {{{}notifyCheckpointComplete(){}}}.
> Therefore, a materialization ID consumed by {{nativeSavepoint()}} should not 
> be able to block all subsequent materializations indefinitely.
> Two possible approaches are:
> *Option A:* When the native savepoint's snapshot future completes, whether 
> successfully or exceptionally, mark its materialization ID as terminated for 
> the purpose of the guard. For example:
> {code:java}
> lastFailedMaterializationId =
>         Math.max(lastFailedMaterializationId, materializationID);
> {code}
> While the savepoint is in progress, the guard would continue to block 
> periodic materialization, preserving the no-concurrent-materialization 
> property. Restore behavior would remain unaffected because 
> {{completeRestore()}} reads the materialization ID from the savepoint handle 
> itself.
> *Option B:* Do not register the savepoint's checkpoint ID in 
> {{{}materializationIdByCheckpointId{}}}, because the notification it waits 
> for can never arrive. Instead, account for materialization IDs consumed by 
> native savepoints separately in the guard.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to