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]

Reply via email to