rkhachatryan commented on code in PR #28709:
URL: https://github.com/apache/flink/pull/28709#discussion_r3587955490


##########
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/metadata/MetadataV2V3SerializerBase.java:
##########
@@ -760,13 +793,43 @@ static void serializeStreamStateHandle(StreamStateHandle 
stateHandle, DataOutput
             for (int keyGroup : keyGroupsStateHandle.getKeyGroupRange()) {
                 
dos.writeLong(keyGroupsStateHandle.getOffsetForKeyGroup(keyGroup));
             }
-            
serializeStreamStateHandle(keyGroupsStateHandle.getDelegateStateHandle(), dos);
+            
serializeStreamStateHandle(keyGroupsStateHandle.getDelegateStateHandle(), dos, 
context);
         } else {
             throw new IOException(
                     "Unknown implementation of StreamStateHandle: " + 
stateHandle.getClass());
         }
     }
 
+    /**
+     * Decides whether a {@link RelativeFileStateHandle} can be written with 
the relative encoding,
+     * or has to be written with its absolute path like a plain {@link 
FileStateHandle}.
+     *
+     * <p>The relative encoding stores only the file name. Whoever reads the 
metadata later rebuilds
+     * the full path as {@code <directory containing the metadata>/<file 
name>} (see {@link
+     * #resolveRelativeHandle}). This only works for files that actually live 
in the directory the
+     * metadata is written to. Files elsewhere (for example a claimed 
savepoint's SSTs that a later
+     * incremental checkpoint still references) would be looked up at a path 
where they do not
+     * exist, so their handles must keep the absolute path. If the directory 
being written to is
+     * unknown ({@code null}), the legacy behavior applies and the relative 
encoding is kept.
+     */
+    private static boolean canKeepRelativeEncoding(
+            RelativeFileStateHandle handle, @Nullable SerializationContext 
context) {

Review Comment:
   Could you elaborate what those legacy paths are?
   
   I think keeping this nullable makes it a bit more error-prone without giving 
much value: if this change happens to be problematic there will be no way to 
disable it anyways



##########
flink-tests/src/test/java/org/apache/flink/test/checkpointing/ResumeCheckpointAfterClaimedNativeSavepointITCase.java:
##########
@@ -0,0 +1,238 @@
+/*
+ * 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.test.checkpointing;
+
+import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.api.common.functions.OpenContext;
+import org.apache.flink.api.common.functions.RichMapFunction;
+import org.apache.flink.api.common.state.ValueState;
+import org.apache.flink.api.common.state.ValueStateDescriptor;
+import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
+import org.apache.flink.client.program.ClusterClient;
+import org.apache.flink.configuration.CheckpointingOptions;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.ExternalizedCheckpointRetention;
+import org.apache.flink.configuration.MemorySize;
+import org.apache.flink.configuration.RestartStrategyOptions;
+import org.apache.flink.configuration.StateBackendOptions;
+import org.apache.flink.configuration.StateChangelogOptions;
+import org.apache.flink.core.execution.RecoveryClaimMode;
+import org.apache.flink.core.execution.SavepointFormatType;
+import org.apache.flink.runtime.checkpoint.metadata.CheckpointMetadata;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.jobgraph.SavepointRestoreSettings;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import 
org.apache.flink.runtime.state.IncrementalKeyedStateHandle.HandleAndLocalPath;
+import org.apache.flink.runtime.state.IncrementalRemoteKeyedStateHandle;
+import org.apache.flink.runtime.state.filesystem.FileStateHandle;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.state.rocksdb.RocksDBConfigurableOptions;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.v2.DiscardingSink;
+import org.apache.flink.test.util.MiniClusterWithClientResource;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+import java.util.List;
+import java.util.Optional;
+import java.util.stream.Collectors;
+
+import static 
org.apache.flink.runtime.testutils.CommonTestUtils.waitForAllTaskRunning;
+import static org.apache.flink.test.util.TestUtils.loadCheckpointMetadata;
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Verifies, through a real cluster, that an incremental checkpoint taken 
after a CLAIM-mode restore
+ * from a NATIVE savepoint remains restorable.
+ *
+ * <p>A NATIVE RocksDB savepoint references its SST files relative to the 
savepoint directory. A job
+ * restored from it in CLAIM mode reuses those files, so its first incremental 
checkpoint references
+ * files that live outside the checkpoint's own exclusive directory. The 
checkpoint metadata must
+ * record such references by their absolute location: a relative encoding 
would be resolved against
+ * the checkpoint directory when the metadata is read back, and restoring from 
the checkpoint would
+ * fail looking for files that were never there.
+ *
+ * <p>{@code RocksDBClaimSavepointThenCheckpointRestoreTest} pins the same 
property at the
+ * backend/serializer level, in-process, without the JobManager finalization 
wiring. This test
+ * instead drives the whole path on a real cluster, so it also covers what the 
component test
+ * cannot: in particular that the metadata finalization in {@code 
PendingCheckpoint} binds to the
+ * {@code CheckpointMetadataOutputStream} overload of {@code 
Checkpoints#storeCheckpointMetadata},
+ * which hands the checkpoint's exclusive directory to the serializer.
+ */
+class ResumeCheckpointAfterClaimedNativeSavepointITCase {
+
+    private static final int PARALLELISM = 2;
+    private static final int NUM_KEYS = 128;
+
+    @TempDir private Path checkpointsDir;
+    @TempDir private Path savepointsDir;
+
+    @Test
+    void testFirstCheckpointAfterClaimIsRestorable() throws Exception {
+        final Configuration config = new Configuration();
+        config.set(StateBackendOptions.STATE_BACKEND, "rocksdb");
+        config.set(CheckpointingOptions.INCREMENTAL_CHECKPOINTS, true);
+        config.set(StateChangelogOptions.ENABLE_STATE_CHANGE_LOG, false);
+        config.set(CheckpointingOptions.FILE_MERGING_ENABLED, false);
+        config.set(RocksDBConfigurableOptions.USE_INGEST_DB_RESTORE_MODE, 
false);
+        config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, 
checkpointsDir.toUri().toString());
+        config.set(
+                CheckpointingOptions.EXTERNALIZED_CHECKPOINT_RETENTION,
+                ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
+        // Keep every state file a real file so the shared-file assertion sees 
actual paths.
+        config.set(CheckpointingOptions.FS_SMALL_FILE_THRESHOLD, 
MemorySize.ZERO);
+        // Fail fast on a broken restore instead of looping through restarts.
+        config.set(RestartStrategyOptions.RESTART_STRATEGY, "none");
+
+        // Manual lifecycle because the configuration depends on per-test 
@TempDirs.
+        final MiniClusterWithClientResource cluster =
+                new MiniClusterWithClientResource(
+                        new MiniClusterResourceConfiguration.Builder()
+                                .setConfiguration(config)
+                                .setNumberTaskManagers(1)
+                                .setNumberSlotsPerTaskManager(PARALLELISM)
+                                .build());
+        cluster.before();
+        try {
+            final ClusterClient<?> client = cluster.getClusterClient();
+            final MiniCluster miniCluster = cluster.getMiniCluster();
+
+            // Checkpoint once so the savepoint covers a mix of reused and 
freshly flushed SSTs,
+            // then stop with a NATIVE savepoint.
+            final JobGraph initialJob = createJobGraph(config);
+            client.submitJob(initialJob).get();
+            waitForAllTaskRunning(miniCluster, initialJob.getJobID(), false);
+            miniCluster.triggerCheckpoint(initialJob.getJobID()).get();
+            final String savepointPath =
+                    client.stopWithSavepoint(
+                                    initialJob.getJobID(),
+                                    false,
+                                    savepointsDir.toUri().toString(),
+                                    SavepointFormatType.NATIVE)
+                            .get();
+
+            // Restore in CLAIM mode, the first incremental checkpoint reuses 
the savepoint's SSTs.
+            final JobGraph restoredFromSavepoint = createJobGraph(config);
+            restoredFromSavepoint.setSavepointRestoreSettings(
+                    SavepointRestoreSettings.forPath(
+                            savepointPath, false, RecoveryClaimMode.CLAIM));
+            client.submitJob(restoredFromSavepoint).get();
+            waitForAllTaskRunning(miniCluster, 
restoredFromSavepoint.getJobID(), false);
+            final String checkpointPath =
+                    
miniCluster.triggerCheckpoint(restoredFromSavepoint.getJobID()).get();
+            client.cancel(restoredFromSavepoint.getJobID()).get();
+            
miniCluster.requestJobResult(restoredFromSavepoint.getJobID()).get();
+
+            // The regression guard: reused savepoint files must be recorded 
by their absolute
+            // location. The trailing separator keeps a sibling directory 
prefix from matching.
+            final String savepointPrefix =
+                    savepointPath.endsWith("/") ? savepointPath : 
savepointPath + "/";
+            assertThat(sharedStateFilePaths(checkpointPath))

Review Comment:
   This might be flaky: 
   - the 1st checkpoint might not create any SSTs
   - and the 2nd checkpoint might compact all the inherited SSTs out
   
   The former can be avoided by re-taking the checkpoint; the latter can be 
avoided by disabling compactions.
   But probably it's enough to just retry the checkpoint-restore for several 
iterations if this happens?



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