sangkyoonnam commented on code in PR #1175:
URL: https://github.com/apache/flink-agents/pull/1175#discussion_r4235327035


##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/FlussActionStateStore.java:
##########
@@ -483,19 +493,41 @@ private Map<Integer, Long> getBucketOffsets(OffsetSpec 
offsetSpec) {
                 buckets.add(b);
             }
             return admin.listOffsets(tablePath, buckets, 
offsetSpec).all().get();
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            throw new RuntimeException(
+                    "Interrupted getting offsets for Fluss table: " + 
tablePath, e);
         } catch (Exception e) {
             throw new RuntimeException("Failed to get offsets for Fluss table: 
" + tablePath, e);
         }
     }
 
     /**
-     * Returns the end offsets of each bucket as a recovery marker. Similar to 
Kafka's
-     * implementation, this captures the current log position so that {@link 
#rebuildState} can
-     * resume from these offsets instead of scanning from the beginning.
+     * Captures bucket end offsets, then synchronously appends the cached 
states after those offsets
+     * before returning the recovery marker.
+     *
+     * <p>A pending action may already have persisted call results before the 
checkpoint. Starting
+     * recovery at the unmodified log end would lose those results. Fluss 
append acknowledgements do
+     * not expose record offsets, so refreshing the cache into the recovery 
window avoids guessing
+     * where the earlier writes landed. Completed actions are also retained 
until the enclosing
+     * input sequence is pruned after a completed checkpoint.
+     *
+     * <p>Called on the operator thread, serially with puts and pruning. A 
failed refresh fails the
+     * checkpoint; partial appends remain safe for earlier recovery markers. 
This adds one append
+     * per cached state per snapshot and requires Fluss retention to preserve 
the recovery window.
      */
     @Override
     public Object getRecoveryMarker() {
-        return getBucketEndOffsets();
+        Map<Integer, Long> offsets = getBucketEndOffsets();
+        try {
+            for (Map.Entry<String, ActionState> entry : 
actionStates.entrySet()) {

Review Comment:
   The PR already documents this cost and explains why `isCompleted()` isn't a 
safe filter. A narrower filter is safe: skip states whose input sequence is at 
or below the key's `lastCompletedSequenceNumber` in this checkpoint. No replay 
after restoring this checkpoint needs them. Without it, refreshes include every 
completed input still awaiting pruning, not only in-flight ones. With timely 
checkpoint completion, a subtask completing 10 inputs/s, each leaving two 
cached action states, accumulates about 1,200 states over a 60 s interval. Each 
is an acknowledged append on the mailbox thread during the synchronous snapshot.
   
   `snapshotLastCompletedSequenceNumbers` already collects those boundaries, 
immediately after `snapshotRecoveryMarker()`. Collecting them before the 
refresh and registering them in `checkpointIdToSeqNums` after it succeeds would 
let this loop use them. Keys without a boundary keep all their states. That's 
the `needsReplay` rule in #1161 and [joeyutong's suggestion 
there](https://github.com/apache/flink-agents/pull/1161#discussion_r4226684419),
 so I'd add it as a follow-up shared by both stores once #1161 settles, unless 
you'd rather fold it in here. A test with states at, below and above the 
boundary, including a completed action in an unfinished input, would pin it 
down.
   
   Does a shared hook after #1161 work for you?



-- 
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