From 5b25a3af632cfd9ac4489ec1b692b46a9cdd73b9 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Wed, 3 Jun 2020 22:48:44 -0700 Subject: [PATCH 1/3] internalize checkpoint --- .../processor/internals/StandbyTask.java | 8 +-- .../processor/internals/StreamTask.java | 56 +++++++++++-------- .../streams/processor/internals/Task.java | 4 +- .../processor/internals/TaskManager.java | 33 +++++------ .../processor/internals/StandbyTaskTest.java | 9 ++- .../processor/internals/StreamTaskTest.java | 26 ++++----- .../processor/internals/TaskManagerTest.java | 28 +++++----- .../kafka/streams/TopologyTestDriver.java | 4 +- 8 files changed, 82 insertions(+), 86 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java index 8cba911e51c45..41598d6163201 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java @@ -33,7 +33,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.Map; -import java.util.Objects; import java.util.Set; /** @@ -153,12 +152,10 @@ public void postCommit() { } @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { prepareClose(true); log.info("Prepared clean close"); - - return Collections.emptyMap(); } @Override @@ -199,8 +196,7 @@ private void prepareClose(final boolean clean) { } @Override - public void closeClean(final Map checkpoint) { - Objects.requireNonNull(checkpoint); + public void closeClean() { close(true); log.info("Closed clean"); 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 d045144cb8414..59c26a08fa5bb 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 @@ -107,6 +107,9 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator, private boolean commitNeeded = false; private boolean commitRequested = false; + private boolean checkpointNeeded = false; + private Map checkpoint = Collections.emptyMap(); + public StreamTask(final TaskId id, final Set partitions, final ProcessorTopology topology, @@ -465,17 +468,15 @@ private Map extractPartitionTimes() { } @Override - public Map prepareCloseClean() { - final Map checkpoint = prepareClose(true); + public void prepareCloseClean() { + prepareClose(true); log.info("Prepared clean close"); - - return checkpoint; } @Override - public void closeClean(final Map checkpoint) { - close(true, checkpoint); + public void closeClean() { + close(true); log.info("Closed clean"); } @@ -489,7 +490,8 @@ public void prepareCloseDirty() { @Override public void closeDirty() { - close(false, null); + + close(false); log.info("Closed dirty"); } @@ -505,11 +507,10 @@ public void update(final Set topicPartitions, final ProcessorTop @Override public void closeAndRecycleState() { - final Map checkpoint = prepareClose(true); + prepareClose(true); + + writeCheckpointIfNeed(); - if (checkpoint != null) { - stateMgr.checkpoint(checkpoint); - } switch (state()) { case CREATED: case RUNNING: @@ -546,14 +547,14 @@ public void closeAndRecycleState() { * otherwise, just close open resources * @throws TaskMigratedException if the task producer got fenced (EOS) */ - private Map prepareClose(final boolean clean) { - final Map checkpoint; + private void prepareClose(final boolean clean) { + // Reset any previously scheduled checkpoint. + checkpointNeeded = false; switch (state()) { case CREATED: // the task is created and not initialized, just re-write the checkpoint file - checkpoint = Collections.emptyMap(); - + scheduleCheckpoint(emptyMap()); break; case RUNNING: @@ -562,9 +563,8 @@ private Map prepareClose(final boolean clean) { if (clean) { stateMgr.flush(); recordCollector.flush(); - checkpoint = checkpointableOffsets(); + scheduleCheckpoint(checkpointableOffsets()); } else { - checkpoint = null; // `null` indicates to not write a checkpoint executeAndMaybeSwallow(false, stateMgr::flush, "state manager flush", log); } @@ -572,22 +572,30 @@ private Map prepareClose(final boolean clean) { case RESTORING: executeAndMaybeSwallow(clean, stateMgr::flush, "state manager flush", log); - checkpoint = Collections.emptyMap(); + scheduleCheckpoint(emptyMap()); break; case SUSPENDED: case CLOSED: // not need to checkpoint, since when suspending we've already committed the state - checkpoint = null; // `null` indicates to not write a checkpoint - break; default: throw new IllegalStateException("Unknown state " + state() + " while prepare closing active task " + id); } + } + + private void scheduleCheckpoint(final Map checkpoint) { + this.checkpoint = checkpoint; + this.checkpointNeeded = true; + } - return checkpoint; + private void writeCheckpointIfNeed() { + if (checkpointNeeded) { + stateMgr.checkpoint(checkpoint); + checkpointNeeded = false; + } } /** @@ -598,9 +606,9 @@ private Map prepareClose(final boolean clean) { * 3. finally release the state manager lock * */ - private void close(final boolean clean, final Map checkpoint) { - if (clean && checkpoint != null) { - executeAndMaybeSwallow(clean, () -> stateMgr.checkpoint(checkpoint), "state manager checkpoint", log); + private void close(final boolean clean) { + if (clean && checkpointNeeded) { + executeAndMaybeSwallow(true, this::writeCheckpointIfNeed, "state manager checkpoint", log); } switch (state()) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java index eee290ad8e733..9283e861d714d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java @@ -145,12 +145,12 @@ enum TaskType { * * @throws StreamsException fatal error, should close the thread */ - Map prepareCloseClean(); + void prepareCloseClean(); /** * Must be idempotent. */ - void closeClean(final Map checkpoint); + void closeClean(); /** * Prepare to close a task that we may not own. Discard any uncommitted progress and close the task. diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index d1be8a3f3055e..4b290f56cee96 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -189,7 +189,7 @@ public void handleAssignment(final Map> activeTasks, // first rectify all existing tasks final LinkedHashMap taskCloseExceptions = new LinkedHashMap<>(); - final Map> checkpointPerTask = new HashMap<>(); + final Set tasksToClose = new HashSet<>(); final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); final Set additionalTasksForCommitting = new HashSet<>(); final Set dirtyTasks = new HashSet<>(); @@ -209,11 +209,11 @@ public void handleAssignment(final Map> activeTasks, tasksToRecycle.add(task); } else { try { - final Map checkpoint = task.prepareCloseClean(); + task.prepareCloseClean(); final Map committableOffsets = task .committableOffsetsAndMetadata(); - checkpointPerTask.put(task, checkpoint); + tasksToClose.add(task); if (!committableOffsets.isEmpty()) { consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); } @@ -249,20 +249,17 @@ public void handleAssignment(final Map> activeTasks, log.error("Failed to batch commit tasks, " + "will close all tasks involved in this commit as dirty by the end", e); dirtyTasks.addAll(additionalTasksForCommitting); - dirtyTasks.addAll(checkpointPerTask.keySet()); + dirtyTasks.addAll(tasksToClose); - checkpointPerTask.clear(); + tasksToClose.clear(); // Just add first taskId to re-throw by the end. taskCloseExceptions.put(consumedOffsetsAndMetadataPerTask.keySet().iterator().next(), e); } } - for (final Map.Entry> taskAndCheckpoint : checkpointPerTask.entrySet()) { - final Task task = taskAndCheckpoint.getKey(); - final Map checkpoint = taskAndCheckpoint.getValue(); - + for (final Task task : tasksToClose) { try { - completeTaskCloseClean(task, checkpoint); + completeTaskCloseClean(task); cleanUpTaskProducer(task, taskCloseExceptions); tasks.remove(task.id()); } catch (final RuntimeException e) { @@ -630,9 +627,9 @@ private void closeTaskDirty(final Task task) { task.closeDirty(); } - private void completeTaskCloseClean(final Task task, final Map checkpoint) { + private void completeTaskCloseClean(final Task task) { cleanupTask(task); - task.closeClean(checkpoint); + task.closeClean(); } // Note: this MUST be called *before* actually closing the task @@ -651,16 +648,16 @@ private void cleanupTask(final Task task) { void shutdown(final boolean clean) { final AtomicReference firstException = new AtomicReference<>(null); - final Map> checkpointPerTask = new HashMap<>(); + final Set tasksToClose = new HashSet<>(); final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); for (final Task task : tasks.values()) { if (clean) { try { - final Map checkpoint = task.prepareCloseClean(); + task.prepareCloseClean(); final Map committableOffsets = task.committableOffsetsAndMetadata(); - checkpointPerTask.put(task, checkpoint); + tasksToClose.add(task); if (!committableOffsets.isEmpty()) { consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); } @@ -680,11 +677,9 @@ void shutdown(final boolean clean) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); } - for (final Map.Entry> taskAndCheckpoint : checkpointPerTask.entrySet()) { - final Task task = taskAndCheckpoint.getKey(); - final Map checkpoint = taskAndCheckpoint.getValue(); + for (final Task task : tasksToClose) { try { - completeTaskCloseClean(task, checkpoint); + completeTaskCloseClean(task); } catch (final RuntimeException e) { firstException.compareAndSet(null, e); closeTaskDirty(task); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java index f868de4dff557..04ef95155895d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StandbyTaskTest.java @@ -52,7 +52,6 @@ import java.io.File; import java.io.IOException; import java.util.Collections; -import java.util.Map; import static java.util.Arrays.asList; import static org.apache.kafka.common.utils.Utils.mkEntry; @@ -270,8 +269,8 @@ public void shouldCommitOnCloseClean() { task = createStandbyTask(); task.initializeIfNeeded(); - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -323,8 +322,8 @@ public void shouldThrowOnCloseCleanError() { task = createStandbyTask(); task.initializeIfNeeded(); - final Map checkpoint = task.prepareCloseClean(); - assertThrows(RuntimeException.class, () -> task.closeClean(checkpoint)); + task.prepareCloseClean(); + assertThrows(RuntimeException.class, () -> task.closeClean()); final double expectedCloseTaskMetric = 0.0; verifyCloseTaskMetric(expectedCloseTaskMetric, streamsMetrics, metricName); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java index a4431a2491f30..27135f80bbd39 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java @@ -1567,8 +1567,8 @@ public void shouldCheckpointWithCreatedStateOnClose() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); assertFalse(source1.initialized); @@ -1642,8 +1642,8 @@ public void shouldNotCommitOnCloseRestoring() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); task.initializeIfNeeded(); - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -1668,8 +1668,8 @@ public void shouldCommitOnCloseClean() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); task.initializeIfNeeded(); task.completeRestoration(); - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -1696,8 +1696,8 @@ public void shouldSwallowExceptionOnCloseCleanError() { task.initializeIfNeeded(); task.completeRestoration(); - final Map checkpoint = task.prepareCloseClean(); - assertThrows(ProcessorStateException.class, () -> task.closeClean(checkpoint)); + task.prepareCloseClean(); + assertThrows(ProcessorStateException.class, () -> task.closeClean()); final double expectedCloseTaskMetric = 0.0; verifyCloseTaskMetric(expectedCloseTaskMetric, streamsMetrics, metricName); @@ -1760,8 +1760,8 @@ public void shouldThrowOnCloseCleanCheckpointError() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); task.initializeIfNeeded(); - final Map checkpoint = task.prepareCloseClean(); - assertThrows(ProcessorStateException.class, () -> task.closeClean(checkpoint)); + task.prepareCloseClean(); + assertThrows(ProcessorStateException.class, () -> task.closeClean()); assertEquals(Task.State.RESTORING, task.state()); @@ -1795,11 +1795,11 @@ public void closeShouldBeIdempotent() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); // close calls are idempotent since we are already in closed - task.closeClean(checkpoint); + task.closeClean(); task.closeDirty(); EasyMock.reset(stateManager); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java index 4d38b2f3d9f4d..5338194622cdf 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java @@ -604,7 +604,7 @@ public void shouldReviveCorruptTasksEvenIfTheyCannotCloseClean() { final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, true, stateManager) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new RuntimeException("oops"); } }; @@ -1116,7 +1116,7 @@ public Collection changelogPartitions() { final AtomicBoolean closedDirtyTask03 = new AtomicBoolean(false); final Task task01 = new StateMachineTask(taskId01, taskId01Partitions, true) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("migrated", new RuntimeException("cause")); } @@ -1134,7 +1134,7 @@ public void closeDirty() { }; final Task task02 = new StateMachineTask(taskId02, taskId02Partitions, true) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new RuntimeException("oops"); } @@ -1452,13 +1452,13 @@ public Collection changelogPartitions() { }; final Task task01 = new StateMachineTask(taskId01, taskId01Partitions, true) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("migrated", new RuntimeException("cause")); } }; final Task task02 = new StateMachineTask(taskId02, taskId02Partitions, true) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new RuntimeException("oops"); } }; @@ -2274,14 +2274,14 @@ public void shouldHaveRemainingPartitionsUncleared() { public void shouldThrowTaskMigratedWhenAllTaskCloseExceptionsAreTaskMigrated() { final StateMachineTask migratedTask01 = new StateMachineTask(taskId01, taskId01Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("t1 close exception", new RuntimeException()); } }; final StateMachineTask migratedTask02 = new StateMachineTask(taskId02, taskId02Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("t2 close exception", new RuntimeException()); } }; @@ -2304,14 +2304,14 @@ public Map prepareCloseClean() { public void shouldThrowRuntimeExceptionWhenEncounteredUnknownExceptionDuringTaskClose() { final StateMachineTask migratedTask01 = new StateMachineTask(taskId01, taskId01Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("t1 close exception", new RuntimeException()); } }; final StateMachineTask migratedTask02 = new StateMachineTask(taskId02, taskId02Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new IllegalStateException("t2 illegal state exception", new RuntimeException()); } }; @@ -2333,14 +2333,14 @@ public Map prepareCloseClean() { public void shouldThrowSameKafkaExceptionWhenEncounteredDuringTaskClose() { final StateMachineTask migratedTask01 = new StateMachineTask(taskId01, taskId01Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new TaskMigratedException("t1 close exception", new RuntimeException()); } }; final StateMachineTask migratedTask02 = new StateMachineTask(taskId02, taskId02Partitions, false) { @Override - public Map prepareCloseClean() { + public void prepareCloseClean() { throw new KafkaException("Kaboom for t2!", new RuntimeException()); } }; @@ -2704,15 +2704,13 @@ public void resume() { } @Override - public Map prepareCloseClean() { - return Collections.emptyMap(); - } + public void prepareCloseClean() {} @Override public void prepareCloseDirty() {} @Override - public void closeClean(final Map checkpoint) { + public void closeClean() { transitionTo(State.CLOSED); } diff --git a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java index db310128f411f..3416ec0c5fa42 100644 --- a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java +++ b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java @@ -1180,8 +1180,8 @@ public SessionStore getSessionStore(final String name) { */ public void close() { if (task != null) { - final Map checkpoint = task.prepareCloseClean(); - task.closeClean(checkpoint); + task.prepareCloseClean(); + task.closeClean(); } if (globalStateTask != null) { try { From 6b7bcd100f349698b658a1e4ad5b7735b847d977 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Fri, 5 Jun 2020 19:26:16 -0700 Subject: [PATCH 2/3] minor change --- .../apache/kafka/streams/processor/internals/StreamTask.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 59c26a08fa5bb..dc8ac6d50200d 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 @@ -607,7 +607,7 @@ private void writeCheckpointIfNeed() { * */ private void close(final boolean clean) { - if (clean && checkpointNeeded) { + if (clean) { executeAndMaybeSwallow(true, this::writeCheckpointIfNeed, "state manager checkpoint", log); } From bc9c9dacb0254f78b628d845f1c3e2d8e91fb2e0 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Sat, 6 Jun 2020 08:46:46 -0700 Subject: [PATCH 3/3] only checkpoint struct --- .../kafka/streams/processor/internals/StreamTask.java | 10 ++++------ 1 file changed, 4 insertions(+), 6 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 dc8ac6d50200d..b23a1a7325238 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 @@ -107,8 +107,7 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator, private boolean commitNeeded = false; private boolean commitRequested = false; - private boolean checkpointNeeded = false; - private Map checkpoint = Collections.emptyMap(); + private Map checkpoint = null; public StreamTask(final TaskId id, final Set partitions, @@ -549,7 +548,7 @@ public void closeAndRecycleState() { */ private void prepareClose(final boolean clean) { // Reset any previously scheduled checkpoint. - checkpointNeeded = false; + checkpoint = null; switch (state()) { case CREATED: @@ -588,13 +587,12 @@ private void prepareClose(final boolean clean) { private void scheduleCheckpoint(final Map checkpoint) { this.checkpoint = checkpoint; - this.checkpointNeeded = true; } private void writeCheckpointIfNeed() { - if (checkpointNeeded) { + if (checkpoint != null) { stateMgr.checkpoint(checkpoint); - checkpointNeeded = false; + checkpoint = null; } }