From 2ba81ca14a736d2606eb61ba700aa10280f32811 Mon Sep 17 00:00:00 2001 From: Sharath Date: Wed, 2 Sep 2020 23:56:06 +0530 Subject: [PATCH 1/3] When resuming Streams active task with EOS, the checkpoint file is deleted --- .../processor/internals/ProcessorStateManager.java | 6 ++++++ .../kafka/streams/processor/internals/StreamTask.java | 9 +++++++++ 2 files changed, 15 insertions(+) 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..088b2808c151d 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 deleteCheckPointFile() 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..7c1015c9ea8e3 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.deleteCheckPointFile(); + log.debug("Deleted check point file"); + } catch (final IOException error) { + log.error("Check point file not found"); + } + transitionTo(State.RESTORING); log.info("Resumed to restoring state"); From aae073ba496f1b00c586fe804688892129980199 Mon Sep 17 00:00:00 2001 From: Sharath Date: Thu, 3 Sep 2020 00:13:43 +0530 Subject: [PATCH 2/3] changed error message of IO exception --- .../apache/kafka/streams/processor/internals/StreamTask.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 7c1015c9ea8e3..19ad965ecba20 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 @@ -337,8 +337,8 @@ public void resume() { try { stateMgr.deleteCheckPointFile(); log.debug("Deleted check point file"); - } catch (final IOException error) { - log.error("Check point file not found"); + } catch (final IOException ioe) { + log.error("Encountered error while deleting the checkpoint file due to this exception", ioe); } transitionTo(State.RESTORING); From 489b034a8d34146d6f3de3471f44baf814cbf970 Mon Sep 17 00:00:00 2001 From: Sharath Date: Wed, 16 Sep 2020 20:03:28 +0530 Subject: [PATCH 3/3] unit test cases added --- .../internals/ProcessorStateManager.java | 2 +- .../processor/internals/StreamTask.java | 4 +-- .../internals/ProcessorStateManagerTest.java | 30 +++++++++++++++++++ 3 files changed, 33 insertions(+), 3 deletions(-) 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 088b2808c151d..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 @@ -663,7 +663,7 @@ public String changelogFor(final String storeName) { return storeToChangelogTopic.get(storeName); } - public void deleteCheckPointFile() throws IOException { + 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 19ad965ecba20..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 @@ -335,8 +335,8 @@ public void resume() { // Deleting checkpoint file before transition to RESTORING state (KAFKA-10362) try { - stateMgr.deleteCheckPointFile(); - log.debug("Deleted check point file"); + 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); } 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,