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


##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -224,13 +229,20 @@ public void put(Object key, long seqNum, Action action, 
Event event, ActionState
         try {
             ProducerRecord<String, ActionState> kafkaRecord =
                     new ProducerRecord<>(topic, stateKey, state);
-            producer.send(kafkaRecord);
+            RecordMetadata metadata = producer.send(kafkaRecord).get();

Review Comment:
   Waiting on `send(...).get()` before `flush()` is the ordering I flagged on 
#1177: with kafka-clients 4.0's 5 ms default `linger.ms`, a small record in an 
otherwise idle batch can wait out the linger delay, and puts are sequential 
within a subtask. In the #1177 plain-producer probe, `send.get; flush` took 
6.88 ms per send-loop iteration against 0.84 ms for `send; flush; get`; I 
haven't benchmarked this branch. This PR needs the `RecordMetadata`, so keep 
the future, flush, then await it before updating both caches; that keeps the 
metadata and the failure check without the delay.



##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -586,6 +630,12 @@ public Object getRecoveryMarker() {
             for (Map.Entry<TopicPartition, Long> entry : 
endOffsets.entrySet()) {
                 recoveryOffsets.put(entry.getKey().partition(), 
entry.getValue());
             }
+            latestStateOffsets.forEach(

Review Comment:
   One long-pending action holds the replay start for its whole partition, and 
the per-partition `min` merge makes every subtask replay from there. 
`latestStateOffsets` is rebuilt from Kafka after a restore, but 
`latestKeySeqNum` starts empty, so for keys with no newly completed input the 
already-checkpointed records keep the next marker rewound until a completed 
checkpoint prunes them. In a local probe with one pending key B followed by 
1000 checkpointed and pruned inputs on key A:
   
   ```
   marker={0=0}
   restored cacheSize=1001
   first marker after restore (B done, A untouched)={0=1}
   marker after first checkpoint-complete prune={0=1001} cache=0
   ```
   
   As a non-blocking follow-up to joeyutong's suggestion, collecting the keyed 
completed-sequence boundaries before building the marker and passing them to 
the store would drop the duplicate `latestKeySeqNum` map and keep restored 
completed inputs from becoming rewind candidates. A pending action would still 
pin its partition, so later completed records would still be scanned. For 
comparison, #1175 captures the Fluss end offsets and then rewrites all cached 
states, which avoids holding the old replay start at the cost of synchronous 
writes on every snapshot.



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