diff --git a/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java b/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java index 4182023900b37a..3e4b5db5beb739 100644 --- a/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java +++ b/flink-state-backends/flink-statebackend-changelog/src/main/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackend.java @@ -199,6 +199,9 @@ public class ChangelogKeyedStateBackend /** last failed or cancelled materialization. */ private long lastFailedMaterializationId = -1L; + /** ID of the last materialization actually triggered by {@link #initMaterialization()}. */ + private long lastTriggeredMaterializationId = -1L; + private final ChangelogTruncateHelper changelogTruncateHelper; /** @@ -820,6 +823,7 @@ private ChangelogSnapshotState completeRestore( } } this.lastConfirmedMaterializationId = materializationId; + this.lastTriggeredMaterializationId = materializationId; this.materializedId = materializationId + 1; if (!isRescaling @@ -851,15 +855,15 @@ private ChangelogSnapshotState completeRestore( */ @Override public Optional initMaterialization() throws Exception { - if (lastConfirmedMaterializationId < materializedId - 1 - && lastFailedMaterializationId < materializedId - 1) { + if (lastConfirmedMaterializationId < lastTriggeredMaterializationId + && lastFailedMaterializationId < lastTriggeredMaterializationId) { // SharedStateRegistry potentially requires that the checkpoint's dependency on the // shared file be continuous, it will be broken if we trigger a new materialization // before the previous one has either confirmed or failed. See discussion in // https://github.com/apache/flink/pull/22669#issuecomment-1593370772 . LOG.info( "materialization:{} not confirmed or failed or cancelled, skip trigger new one.", - materializedId - 1); + lastTriggeredMaterializationId); return Optional.empty(); } @@ -878,6 +882,7 @@ public Optional initMaterialization() throws Exception // streamFactory that is designed for state backend snapshot, which requires unique // checkpoint ID. A faked materialized Id is provided here. long materializationID = materializedId++; + lastTriggeredMaterializationId = materializationID; MaterializationRunnable materializationRunnable = new MaterializationRunnable( diff --git a/flink-state-backends/flink-statebackend-changelog/src/test/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackendTest.java b/flink-state-backends/flink-statebackend-changelog/src/test/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackendTest.java index 10a13a9b3d4762..c0f01b5d0ea86f 100644 --- a/flink-state-backends/flink-statebackend-changelog/src/test/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackendTest.java +++ b/flink-state-backends/flink-statebackend-changelog/src/test/java/org/apache/flink/state/changelog/ChangelogKeyedStateBackendTest.java @@ -20,11 +20,14 @@ import org.apache.flink.api.common.ExecutionConfig; import org.apache.flink.api.common.JobID; import org.apache.flink.api.common.typeutils.base.IntSerializer; +import org.apache.flink.core.execution.SavepointFormatType; import org.apache.flink.core.fs.CloseableRegistry; import org.apache.flink.runtime.checkpoint.CheckpointOptions; +import org.apache.flink.runtime.checkpoint.SavepointType; import org.apache.flink.runtime.jobgraph.JobVertexID; import org.apache.flink.runtime.metrics.groups.UnregisteredMetricGroups; import org.apache.flink.runtime.query.KvStateRegistry; +import org.apache.flink.runtime.state.CheckpointStorageLocationReference; import org.apache.flink.runtime.state.KeyGroupRange; import org.apache.flink.runtime.state.KeyedStateHandle; import org.apache.flink.runtime.state.SnapshotResult; @@ -182,4 +185,28 @@ private RunnableFuture> checkpoint( new MemCheckpointStreamFactory(1000), CheckpointOptions.forCheckpointWithDefaultLocation()); } + + @Test + public void testInitMaterializationAfterAbortedNativeSavepoint() throws Exception { + ChangelogKeyedStateBackend backend = createChangelog(createMock()); + try { + long savepointId = 1L; + backend.snapshot( + savepointId, + 0L, + new MemCheckpointStreamFactory(1000), + new CheckpointOptions( + SavepointType.savepoint(SavepointFormatType.NATIVE), + CheckpointStorageLocationReference.getDefault())); + backend.notifyCheckpointAborted(savepointId); + appendMockStateChange(backend); + + assertTrue( + "materialization should not be blocked by an unconfirmed native savepoint", + backend.initMaterialization().isPresent()); + } finally { + backend.close(); + backend.dispose(); + } + } }