diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java index aa69752f88e92..1f4e255331f1a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java @@ -662,4 +662,10 @@ public TopicPartition registeredChangelogPartitionFor(final String storeName) { public String changelogFor(final String storeName) { return storeToChangelogTopic.get(storeName); } + + public void deleteCheckPointFileIfEOSEnabled() throws IOException { + if (eosEnabled) { + checkpointFile.delete(); + } + } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 57aaf9767e283..38fdab86a983b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -332,6 +332,15 @@ public void resume() { case SUSPENDED: // just transit the state without any logical changes: suspended and restoring states // are not actually any different for inner modules + + // Deleting checkpoint file before transition to RESTORING state (KAFKA-10362) + try { + stateMgr.deleteCheckPointFileIfEOSEnabled(); + log.debug("Deleted check point file upon resuming with EOS enabled"); + } catch (final IOException ioe) { + log.error("Encountered error while deleting the checkpoint file due to this exception", ioe); + } + transitionTo(State.RESTORING); log.info("Resumed to restoring state"); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java index 448c2b1b4c715..30b2766523559 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java @@ -988,6 +988,36 @@ public void shouldBeAbleToCloseWithoutRegisteringAnyStores() { stateMgr.close(); } + @Test + public void shouldDeleteCheckPointFileIfEosEnabled() throws IOException { + final long checkpointOffset = 10L; + final Map offsets = mkMap( + mkEntry(persistentStorePartition, checkpointOffset), + mkEntry(nonPersistentStorePartition, checkpointOffset), + mkEntry(irrelevantPartition, 999L) + ); + checkpoint.write(offsets); + final ProcessorStateManager stateMgr = getStateManager(Task.TaskType.ACTIVE, true); + stateMgr.deleteCheckPointFileIfEOSEnabled(); + stateMgr.close(); + assertFalse(checkpointFile.exists()); + } + + @Test + public void shouldNotDeleteCheckPointFileIfEosNotEnabled() throws IOException { + final long checkpointOffset = 10L; + final Map offsets = mkMap( + mkEntry(persistentStorePartition, checkpointOffset), + mkEntry(nonPersistentStorePartition, checkpointOffset), + mkEntry(irrelevantPartition, 999L) + ); + checkpoint.write(offsets); + final ProcessorStateManager stateMgr = getStateManager(Task.TaskType.ACTIVE, false); + stateMgr.deleteCheckPointFileIfEOSEnabled(); + stateMgr.close(); + assertTrue(checkpointFile.exists()); + } + private ProcessorStateManager getStateManager(final Task.TaskType taskType, final boolean eosEnabled) { return new ProcessorStateManager( taskId,