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]
