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


##########
runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreRecoveryTest.java:
##########
@@ -0,0 +1,253 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.flink.agents.runtime.actionstate;
+
+import org.apache.flink.agents.api.Event;
+import org.apache.flink.agents.api.InputEvent;
+import org.apache.flink.agents.api.agents.AgentExecutionOptions;
+import org.apache.flink.agents.api.context.DurableCallable;
+import org.apache.flink.agents.api.context.RunnerContext;
+import org.apache.flink.agents.plan.AgentConfiguration;
+import org.apache.flink.agents.plan.AgentPlan;
+import org.apache.flink.agents.plan.JavaFunction;
+import org.apache.flink.agents.plan.actions.Action;
+import org.apache.flink.agents.runtime.operator.ActionExecutionOperator;
+import org.apache.flink.agents.runtime.operator.ActionExecutionOperatorFactory;
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.api.java.functions.KeySelector;
+import org.apache.flink.runtime.checkpoint.OperatorSubtaskState;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness;
+import org.apache.flink.util.ExceptionUtils;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.consumer.MockConsumer;
+import org.apache.kafka.clients.producer.MockProducer;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.PartitionInfo;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static 
org.apache.flink.agents.runtime.actionstate.ActionStateTestUtils.createKeyEncoder;
+import static 
org.apache.kafka.clients.consumer.internals.AutoOffsetResetStrategy.EARLIEST;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.catchThrowable;
+
+/**
+ * End-to-end recovery test for a pending durable call result with a 
Kafka-backed store.
+ *
+ * <p>An action records a durable call, takes a real Flink checkpoint while 
the action is still
+ * unfinished, then fails. Recovery must rebuild the pending result from the 
checkpoint's recovery
+ * marker instead of skipping it, otherwise the durable supplier runs a second 
time. See
+ * https://github.com/apache/flink-agents/issues/1158.
+ */
+public class KafkaActionStateStoreRecoveryTest {
+
+    private static final String TOPIC = "test-action-state-recovery";
+    private static final int MAX_PARALLELISM = 128;
+
+    private static final AtomicInteger EXTERNAL_CALLS = new AtomicInteger();
+    private static final 
AtomicReference<KeyedOneInputStreamOperatorTestHarness<Long, Long, Object>>
+            HARNESS = new AtomicReference<>();
+    private static final AtomicReference<OperatorSubtaskState> CHECKPOINT = 
new AtomicReference<>();
+
+    /**
+     * Runs a real durable call, takes a real checkpoint while the action is 
still in flight, then
+     * fails without completing.
+     */
+    public static void durableCallThenCheckpointThenFail(Event event, 
RunnerContext context)
+            throws Exception {
+        Long input = (Long) InputEvent.fromEvent(event).getInput();
+        context.durableExecute(
+                new DurableCallable<Long>() {
+                    @Override
+                    public String getId() {
+                        return "multiply";
+                    }
+
+                    @Override
+                    public Class<Long> getResultClass() {
+                        return Long.class;
+                    }
+
+                    @Override
+                    public Long call() {
+                        EXTERNAL_CALLS.incrementAndGet();
+                        return input * 10;
+                    }
+                });
+        if (CHECKPOINT.get() == null) {
+            CHECKPOINT.set(HARNESS.get().snapshot(1L, 1L));

Review Comment:
   Could we take the checkpoint with the pending action still queued, then run 
it after restore and assert the result and supplier call count? The action is 
already dequeued here, so the restored queue is empty and the test only 
verifies store rebuild.



##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -93,9 +94,11 @@ public class KafkaActionStateStore implements 
ActionStateStore {
     // In memory action state for quick state retrieval
     private final Map<String, ActionState> actionStates;
 
-    // Record the lastest sequence number for each key that should be 
considered as valid
+    // Records the latest checkpointed sequence number for each business key 
identity.

Review Comment:
   Could we keep the completed-sequence tracking in the runtime and pass it to 
the store at snapshot time? `snapshotLastCompletedSequenceNumbers()` already 
collects this information, so reusing it could keep the store simpler and avoid 
maintaining a second copy.



##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/ActionStateStore.java:
##########
@@ -116,6 +116,12 @@ void put(Object key, long seqNum, Action action, Event 
event, ActionState state)
      */
     default void setOwnershipFilter(IntPredicate ownershipFilter) {}
 
+    /**
+     * Records the highest sequence number for {@code key} whose effects are 
already reflected in
+     * the current Flink checkpoint state.
+     */
+    default void markCheckpointedSequence(Object key, long seqNum) throws 
Exception {}

Review Comment:
   Could we also cover Fluss in this PR? Its marker still uses bucket end 
offsets and can skip saved results for pending actions in the same way. #1175 
may be useful as a reference.



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