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]