1fanwang commented on code in PR #1161:
URL: https://github.com/apache/flink-agents/pull/1161#discussion_r4126049883
##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -161,9 +166,11 @@ 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:
Done in `f109972f`. `put()` now catches `InterruptedException` before the
generic catch, restores the flag with `Thread.currentThread().interrupt()`, and
rethrows the original exception instead of wrapping it.
`testPutRestoresInterruptFlagWhenSendIsInterrupted` covers it: the send fails
like a cancelled `FutureTask` (clears the flag, throws), and the test asserts
the flag is set again after `put()` returns.
##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -299,6 +310,11 @@ public void setOwnershipFilter(IntPredicate
ownershipFilter) {
this.ownershipFilter = ownershipFilter;
}
+ @Override
+ public void markCheckpointedSequence(Object key, long seqNum) {
+ latestKeySeqNum.merge(keyEncoder.generateBusinessKeyIdentity(key),
seqNum, Math::max);
Review Comment:
Done in `f109972f`. `pruneState()` now drops the key from `latestKeySeqNum`
once the identity has no live state left in the cache, so entries stop
accumulating per key for the job lifetime. Losing the boundary is harmless: any
later record for that identity is newer than the last completed sequence and
must be replayed either way.
`testPruneDropsCheckpointedBoundaryWhenNoLiveStateRemains` covers it.
##########
runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java:
##########
@@ -233,6 +233,87 @@ void testRecoveryMarker() throws Exception {
assertThat((Map<Integer, Long>) secondMarker).containsEntry(1, 3L);
}
+ @Test
+ void testRebuildStateKeepsPendingResultBeforeRecoveryMarker() throws
Exception {
+ ActionState pendingState = new ActionState(testEvent);
+ actionStateStore.put(TEST_KEY, 1L, testAction, testEvent,
pendingState);
Review Comment:
Replaced by `KafkaActionStateStoreRecoveryTest` in `b07f9bd5`, which now
executes a real durable call and takes a real checkpoint in the window the
issue describes:
- The action runs `context.durableExecute(...)` with a real supplier, then
snapshots the operator from inside the action body while the durable result is
still pending, then fails without completing.
- Recovery takes the checkpoint into a fresh `KafkaActionStateStore` via
`initializeState`, so it exercises the real `handleRecovery` -> `rebuildState`
path against the checkpoint's marker.
- The assertions check the pending result survived: the recovered
`ActionState` is incomplete, carries the one successful `CallResult`, and the
supplier ran exactly once.
A constraint: under `mvn test` the operator always runs the synchronous
executor, because the JDK 21 continuation executor lives only in the packaged
multi-release JAR and is never on the surefire classpath, so a public-API
checkpoint only observes post-invocation state. The checkpoint is therefore
taken from inside the action body to land in the pre-completion window. I kept
the original unit test because it still exercises the same marker rewind and
rebuild path at unit scope, cheaply.
--
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]