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

최원용 updated FLINK-40302:
------------------------
    Priority: Critical  (was: Major)

> 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
>
> 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}
> 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}
> 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|https://issues.apache.org/jira/browse/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|https://issues.apache.org/jira/browse/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.
> I would like to work on this issue. Could a committer please assign it to me?



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

Reply via email to