rob-9 opened a new pull request, #1175: URL: https://github.com/apache/flink-agents/pull/1175
Linked issue: #1174 ### Purpose of change Recovery should reuse a durable call result saved before a checkpoint when its action is still pending. Fluss previously saved bucket end offsets as the recovery marker, which could place that result before the replay window and cause the external call to run again. Closes #1174. Related Kafka fix: #1161; this PR does not depend on it. #### Runtime flow When the operator snapshots its recovery marker, `FlussActionStateStore` first reads the bucket end offsets, then synchronously appends the latest cached action states. It returns the captured offsets only after every append succeeds. Recovery reads those refreshed states and any subsequent writes through the existing replay path. Checkpoint completion still prunes completed input sequences from the cache. #### Design and comparison with Kafka The Kafka implementation in #1161, reviewed at `14ee1fe`, tracks exact offsets returned by producer writes and uses per-key completed-sequence boundaries to decide which records need replay. Fluss 0.9.0's `AppendResult` does not expose the written record's offset, so this PR refreshes states after the marker instead. | Approach | Benefit | Cost | | --- | --- | --- | | Kafka #1161: track offsets and rewind | No checkpoint refresh writes; completed sequences need not hold the replay window back | Additional offset and sequence bookkeeping, including recovery reconstruction and eviction | | This PR: capture offsets and refresh cached Fluss states | Uses the existing cache and recovery-marker format; no offset guesses or new per-key maps | One synchronous append per cached state per snapshot, including completed states awaiting pruning; more checkpoint latency and log traffic | The review decision for this draft is whether that checkpoint write cost is acceptable. Performance under large caches and frequent checkpoints has not been measured. Filtering solely on `ActionState.isCompleted()` would be unsafe: a completed action can still belong to an input sequence whose other actions are pending. ### Behavioral Semantics | State at snapshot | Behavior | | --- | --- | | Pending action with saved call results | Refresh into the recovery window | | Completed action still in the cache | Refresh until its input sequence is pruned | | State reconstructed during recovery | Refresh again at the next snapshot | | Pruned or divergence-evicted entry | No refresh write | | Empty cache | Capture offsets without appending records | A metadata lookup, serialization, or append failure prevents the marker from being returned and fails the snapshot. Partial refresh writes do not evict cached state; a retry refreshes all remaining entries again. Interrupted offset lookups and appends preserve the interrupt flag. Fluss retention must still preserve the required replay window. ### Tests | Contract | Coverage | | --- | --- | | Pending results survive repeated checkpoints and replay completes without repeating the call | `FlussActionStateRecoveryIntegrationTest#pendingCallSurvivesRepeatedCheckpointRestoreAndCompletes` | | Pre-checkpoint states remain reachable across buckets | `FlussActionStateStoreIntegrationTest#testRebuildStateWithRecoveryMarkers`, `#testMultiBucketRecovery` | | Refresh uses latest cached states, including completed actions | `FlussActionStateStoreTest#testCheckpointRefreshesLatestStatesAfterCapturingOffsets` | | Pruning and divergence stop refresh writes | `#testCheckpointDoesNotRefreshPrunedOrDivergedStates` | | Partial refresh failure retains state for retry | `#testFailedCheckpointRefreshRetainsCacheForRetry` | | Cancellation preserves interruption | `#testInterruptedRefreshRestoresInterruptFlag`, `#testInterruptedOffsetLookupDoesNotRefresh` | All 99 focused tests passed on JDK 17 with an embedded Fluss 0.9.0-incubating cluster. A review rerun of the 23 store and operator recovery tests also passed. The operator regression fails against the original implementation at `2e97add7` because the saved state is missing after restore, and passes with the fix. The Maven reactor rebuilt the affected modules from source; Spotless and Apache RAT checks passed. The operator test uses an embedded Fluss cluster and real keyed-state snapshots/restores. A test task deterministically yields after executing a real durable call, so it runs on every supported JDK; subsequent recovery executes the normal Java action through completion. It does not exercise JDK continuation scheduling, a full distributed Flink restart, or Python-specific suspension. Retention expiry and checkpoint performance have not been tested. Before marking this draft ready, benchmark checkpoint cost with a representative cache size and add a regression that restores the previous successful checkpoint after a later checkpoint's refresh partially succeeds and then fails. The current failure test verifies cache preservation and retry, but not that recovery scenario. <details> <summary>Kafka review considerations and verification command</summary> The Kafka review identified interrupted writes clearing the interrupt flag and completed-sequence bookkeeping growing after state eviction. Both were addressed in its latest reviewed revision. This implementation preserves interruption and adds no extra per-key bookkeeping. The Kafka operator test restores the saved call result but does not run the pending action through completion. This PR's regression test additionally restores twice, executes the pending action, verifies its output and a call count of one, and checks recovery after completion and pruning. - Interrupt handling: https://github.com/apache/flink-agents/pull/1161#discussion_r4108563928 - Bookkeeping cleanup: https://github.com/apache/flink-agents/pull/1161#discussion_r4108564256 - Replay test coverage: https://github.com/apache/flink-agents/pull/1161#discussion_r4108564565 ```sh mvn -pl runtime -am test -Dtest=FlussActionStateStoreTest,FlussActionStateStoreIntegrationTest,FlussActionStateRecoveryIntegrationTest,DurableExecutionManagerTest,ActionExecutionOperatorTest -Dsurefire.failIfNoSpecifiedTests=false ``` </details> ### API No new configuration, public API, table schema, or serialized marker format. Fluss marker capture now performs synchronous writes. Existing markers remain readable, but this change cannot repair results already excluded by a checkpoint taken before the fix. ### Documentation - [ ] `doc-needed` - [ ] `doc-not-needed` - [x] `doc-included` ### Was this patch authored or co-authored using generative AI tooling? - [x] Yes - [ ] No -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
