From 3789a7c1dc71b69ec4d6af621a3415210d731fa8 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jun 2020 12:25:03 -0700 Subject: [PATCH 01/14] enforce SUSPENDED and check commitNeeded --- .../processor/internals/StandbyTask.java | 19 +++++++++++---- .../processor/internals/StreamTask.java | 23 +++++++++---------- .../processor/internals/TaskManager.java | 11 +++++---- 3 files changed, 31 insertions(+), 22 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 1b069d6bfc03d..79a4e4b93c990 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 @@ -112,10 +112,19 @@ public void completeRestoration() { @Override public void suspend() { log.trace("No-op suspend with state {}", state()); - if (state() == State.RUNNING) { - transitionTo(State.SUSPENDED); - } else if (state() == State.RESTORING) { - throw new IllegalStateException("Illegal state " + state() + " while suspending standby task " + id); + switch (state()) { + case CREATED: + case RUNNING: + case SUSPENDED: + transitionTo(State.SUSPENDED); + break; + + case RESTORING: + case CLOSED: + throw new IllegalStateException("Illegal state " + state() + " while suspending standby task " + id); + + default: + throw new IllegalStateException("Unknown state " + state() + " while suspending standby task " + id); } } @@ -175,7 +184,7 @@ public void closeAndRecycleState() { suspend(); prepareCommit(); - if (state() == State.CREATED || state() == State.SUSPENDED) { + if (state() == State.SUSPENDED) { stateMgr.recycle(); } else { throw new IllegalStateException("Illegal state " + state() + " while closing standby task " + id); 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 7f08643648582..b1726d9ab8487 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 @@ -250,14 +250,10 @@ public void completeRestoration() { public void suspend() { switch (state()) { case CREATED: - case SUSPENDED: - log.info("Skip suspending since state is {}", state()); - - break; - case RESTORING: + case SUSPENDED: transitionTo(State.SUSPENDED); - log.info("Suspended restoring"); + log.info("Suspended {}", state()); break; @@ -475,17 +471,20 @@ public void update(final Set topicPartitions, final Map> activeTasks, } else { try { task.suspend(); - final Map committableOffsets = task.prepareCommit(); - - tasksToClose.add(task); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.commitNeeded()) { + final Map committableOffsets = task.prepareCommit(); + if (!committableOffsets.isEmpty()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + } } + tasksToClose.add(task); } catch (final RuntimeException e) { final String uncleanMessage = String.format( "Failed to close task %s cleanly. Attempting to close remaining tasks before re-throwing:", From f621c3a02d673a57bcfb32a653a257bfc4d48241 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jun 2020 14:07:34 -0700 Subject: [PATCH 02/14] check if commit needed in TM --- .../streams/processor/internals/Task.java | 10 ++++---- .../processor/internals/TaskManager.java | 24 +++++++++++-------- 2 files changed, 19 insertions(+), 15 deletions(-) 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 62332c729994f..70211cc59e75b 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 @@ -66,11 +66,11 @@ public interface Task { * */ enum State { - CREATED(1, 4), // 0 - RESTORING(2, 3, 4), // 1 - RUNNING(3), // 2 - SUSPENDED(1, 4), // 3 - CLOSED(0); // 4, we allow CLOSED to transit to CREATED to handle corrupted tasks + CREATED(1, 3, 4), // 0 + RESTORING(2, 3, 4), // 1 + RUNNING(3), // 2 + SUSPENDED(1, 3, 4), // 3 + CLOSED(0); // 4, we allow CLOSED to transit to CREATED to handle corrupted tasks private final Set validTransitions = new HashSet<>(); 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 e52863d4cfb8b..7289ac7b7cfb5 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 @@ -475,10 +475,11 @@ void handleRevocation(final Collection revokedPartitions) { for (final Task task : tasks.values()) { if (remainingPartitions.containsAll(task.inputPartitions())) { task.suspend(); - final Map committableOffsets = task.prepareCommit(); - - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.commitNeeded()) { + final Map committableOffsets = task.prepareCommit(); + if (!committableOffsets.isEmpty()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + } } } else if (task.isActive() && task.commitNeeded()) { final Map committableOffsets = task.prepareCommit(); @@ -697,12 +698,13 @@ void shutdown(final boolean clean) { if (clean) { try { task.suspend(); - final Map committableOffsets = task.prepareCommit(); - - tasksToClose.add(task); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.commitNeeded()) { + final Map committableOffsets = task.prepareCommit(); + if (!committableOffsets.isEmpty()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + } } + tasksToClose.add(task); } catch (final TaskMigratedException e) { // just ignore the exception as it doesn't matter during shutdown closeTaskDirty(task); @@ -721,7 +723,9 @@ void shutdown(final boolean clean) { for (final Task task : tasksToClose) { try { - task.postCommit(); + if (consumedOffsetsAndMetadataPerTask.containsKey(task.id())) { + task.postCommit(); + } completeTaskCloseClean(task); } catch (final RuntimeException e) { firstException.compareAndSet(null, e); From 0c3c080bf6c6007979404dcfaf45e75399711332 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jun 2020 14:15:41 -0700 Subject: [PATCH 03/14] just suspend in closeandRecycleState --- .../streams/processor/internals/StandbyTask.java | 2 -- .../streams/processor/internals/StreamTask.java | 13 ++++--------- .../streams/processor/internals/TaskManager.java | 3 +++ 3 files changed, 7 insertions(+), 11 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 79a4e4b93c990..cff52317ae009 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 @@ -182,8 +182,6 @@ public void closeDirty() { @Override public void closeAndRecycleState() { suspend(); - prepareCommit(); - if (state() == State.SUSPENDED) { stateMgr.recycle(); } else { 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 b1726d9ab8487..5dab5953a52eb 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 @@ -471,12 +471,6 @@ public void update(final Set topicPartitions, final Map> activeTasks, standbyTasksToCreate.remove(task.id()); // check for tasks that were owned previously but have changed active/standby status } else if (activeTasks.containsKey(task.id()) || standbyTasks.containsKey(task.id())) { + if (task.commitNeeded()) { + additionalTasksForCommitting.add(task); + } tasksToRecycle.add(task); } else { try { From a6805838f5c3fbbade7eaa82330b14d88e52f783 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jun 2020 16:20:18 -0700 Subject: [PATCH 04/14] clean up handleRevocation and handleAssignment --- .../processor/internals/StandbyTask.java | 1 - .../processor/internals/StreamTask.java | 2 +- .../processor/internals/TaskManager.java | 99 +++++-------------- 3 files changed, 25 insertions(+), 77 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 cff52317ae009..e6ce9bf3bc13a 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 @@ -181,7 +181,6 @@ public void closeDirty() { @Override public void closeAndRecycleState() { - suspend(); if (state() == State.SUSPENDED) { stateMgr.recycle(); } else { 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 5dab5953a52eb..b579d24bc760d 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 @@ -470,7 +470,7 @@ public void update(final Set topicPartitions, final Map> activeTasks, final Map> activeTasksToCreate = new HashMap<>(activeTasks); final Map> standbyTasksToCreate = new HashMap<>(standbyTasks); final Set tasksToRecycle = new HashSet<>(); + final Set dirtyTasks = new HashSet<>(); builder.addSubscribedTopicsFromAssignment( activeTasks.values().stream().flatMap(Collection::stream).collect(Collectors.toList()), @@ -227,37 +228,31 @@ public void handleAssignment(final Map> activeTasks, // first rectify all existing tasks final LinkedHashMap taskCloseExceptions = new LinkedHashMap<>(); - final Set tasksToClose = new HashSet<>(); - final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); - final Set additionalTasksForCommitting = new HashSet<>(); - final Set dirtyTasks = new HashSet<>(); - for (final Task task : tasks.values()) { if (activeTasks.containsKey(task.id()) && task.isActive()) { updateInputPartitionsAndResume(task, activeTasks.get(task.id())); - if (task.commitNeeded()) { - additionalTasksForCommitting.add(task); - } activeTasksToCreate.remove(task.id()); } else if (standbyTasks.containsKey(task.id()) && !task.isActive()) { updateInputPartitionsAndResume(task, standbyTasks.get(task.id())); standbyTasksToCreate.remove(task.id()); // check for tasks that were owned previously but have changed active/standby status } else if (activeTasks.containsKey(task.id()) || standbyTasks.containsKey(task.id())) { - if (task.commitNeeded()) { - additionalTasksForCommitting.add(task); - } tasksToRecycle.add(task); } else { try { - task.suspend(); if (task.commitNeeded()) { - final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.isActive()) { + log.error("Active task {} was revoked and should have already been committed", task.id()); + throw new IllegalStateException("Revoked active task was not committed during handleRevocation"); + } else { + task.suspend(); + task.prepareCommit(); + task.postCommit(); } } - tasksToClose.add(task); + completeTaskCloseClean(task); + cleanUpTaskProducer(task, taskCloseExceptions); + tasks.remove(task.id()); } catch (final RuntimeException e) { final String uncleanMessage = String.format( "Failed to close task %s cleanly. Attempting to close remaining tasks before re-throwing:", @@ -271,54 +266,15 @@ public void handleAssignment(final Map> activeTasks, } } - if (!consumedOffsetsAndMetadataPerTask.isEmpty()) { - try { - for (final Task task : additionalTasksForCommitting) { - final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); - } - } - - commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - - for (final Task task : additionalTasksForCommitting) { - task.postCommit(); - } - } catch (final RuntimeException e) { - 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(tasksToClose); - - tasksToClose.clear(); - // Just add first taskId to re-throw by the end. - taskCloseExceptions.put(consumedOffsetsAndMetadataPerTask.keySet().iterator().next(), e); - } - } - - for (final Task task : tasksToClose) { - try { - completeTaskCloseClean(task); - cleanUpTaskProducer(task, taskCloseExceptions); - tasks.remove(task.id()); - } catch (final RuntimeException e) { - final String uncleanMessage = String.format("Failed to close task %s cleanly. Attempting to close remaining tasks before re-throwing:", task.id()); - log.error(uncleanMessage, e); - taskCloseExceptions.put(task.id(), e); - // We've already recorded the exception (which is the point of clean). - // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. - dirtyTasks.add(task); - } - } - for (final Task oldTask : tasksToRecycle) { final Task newTask; try { if (oldTask.isActive()) { + // If active, the task should have already been suspended and committed during handleRevocation final Set partitions = standbyTasksToCreate.remove(oldTask.id()); newTask = standbyTaskCreator.createStandbyTaskFromActive((StreamTask) oldTask, partitions); } else { + oldTask.suspend(); final Set partitions = activeTasksToCreate.remove(oldTask.id()); newTask = activeTaskCreator.createActiveTaskFromStandby((StandbyTask) oldTask, partitions, mainConsumer); } @@ -472,42 +428,35 @@ boolean tryToCompleteRestoration() { * @throws TaskMigratedException if the task producer got fenced (EOS only) */ void handleRevocation(final Collection revokedPartitions) { - final Set remainingPartitions = new HashSet<>(revokedPartitions); + final Set remainingRevokedPartitions = new HashSet<>(revokedPartitions); final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); - for (final Task task : tasks.values()) { - if (remainingPartitions.containsAll(task.inputPartitions())) { + for (final Task task : activeTaskIterable()) { + if (remainingRevokedPartitions.containsAll(task.inputPartitions())) { task.suspend(); - if (task.commitNeeded()) { - final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); - } - } - } else if (task.isActive() && task.commitNeeded()) { + } + if (task.commitNeeded()) { final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); } } - remainingPartitions.removeAll(task.inputPartitions()); + remainingRevokedPartitions.removeAll(task.inputPartitions()); } if (!consumedOffsetsAndMetadataPerTask.isEmpty()) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); } - for (final Task task : tasks.values()) { - if (consumedOffsetsAndMetadataPerTask.containsKey(task.id())) { - task.postCommit(); - } + for (final TaskId committedTaskId : consumedOffsetsAndMetadataPerTask.keySet()) { + final Task task = tasks.get(committedTaskId); + task.postCommit(); } - if (!remainingPartitions.isEmpty()) { + if (!remainingRevokedPartitions.isEmpty()) { log.warn("The following partitions {} are missing from the task partitions. It could potentially " + "due to race condition of consumer detecting the heartbeat failure, or the tasks " + - "have been cleaned up by the handleAssignment callback.", remainingPartitions); + "have been cleaned up by the handleAssignment callback.", remainingRevokedPartitions); } } From fb64b9787153ae1da96db9694008139f4d299a4d Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jun 2020 16:54:43 -0700 Subject: [PATCH 05/14] Task and TaskManager tests --- .../processor/internals/StandbyTask.java | 2 +- .../processor/internals/StreamTask.java | 4 +- .../streams/processor/internals/Task.java | 16 +- .../processor/internals/TaskManager.java | 71 +++---- .../processor/internals/StandbyTaskTest.java | 57 +++++- .../processor/internals/StreamTaskTest.java | 64 +++++-- .../processor/internals/TaskManagerTest.java | 173 +++++------------- 7 files changed, 196 insertions(+), 191 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 e6ce9bf3bc13a..20ff5f6715531 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 @@ -195,7 +195,6 @@ public void closeAndRecycleState() { private void close(final boolean clean) { switch (state()) { - case CREATED: case SUSPENDED: executeAndMaybeSwallow( clean, @@ -218,6 +217,7 @@ private void close(final boolean clean) { log.trace("Skip closing since state is {}", state()); return; + case CREATED: case RESTORING: // a StandbyTask is never in RESTORING state case RUNNING: throw new IllegalStateException("Illegal state " + state() + " while closing standby task " + id); 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 b579d24bc760d..07c3c872151bb 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 @@ -521,8 +521,8 @@ private void maybeScheduleCheckpoint() { private void writeCheckpointIfNeed() { if (commitNeeded) { - throw new IllegalStateException("A checkpoint should only be written if the previous commit has completed" - + " and there is no new commit needed."); + log.error("Tried to write a checkpoint with pending uncommitted data, should complete the commit first."); + throw new IllegalStateException("A checkpoint should only be written if no commit is needed."); } if (checkpoint != null) { stateMgr.checkpoint(checkpoint); 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 70211cc59e75b..a48d4ae9c77cd 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 @@ -56,18 +56,18 @@ public interface Task { * | | | | * | v | | * | +------+--------+ | | - * | | Suspended (3) | <---+ | //TODO Suspended(3) could be removed after we've stable on KIP-429 - * | +------+--------+ | - * | | | - * | v | - * | +-----+-------+ | - * +----> | Closed (4) | -----------+ + * +---->| Suspended (3) | ----+ | //TODO Suspended(3) could be removed after we've stable on KIP-429 + * +------+--------+ | + * | | + * v | + * +-----+-------+ | + * | Closed (4) | -----------+ * +-------------+ * */ enum State { - CREATED(1, 3, 4), // 0 - RESTORING(2, 3, 4), // 1 + CREATED(1, 3), // 0 + RESTORING(2, 3), // 1 RUNNING(3), // 2 SUSPENDED(1, 3, 4), // 3 CLOSED(0); // 4, we allow CLOSED to transit to CREATED to handle corrupted tasks 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 e2d8796d852b9..1e843c41eccf6 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 @@ -215,19 +215,20 @@ public void handleAssignment(final Map> activeTasks, "\tExisting standby tasks: {}", activeTasks.keySet(), standbyTasks.keySet(), activeTaskIds(), standbyTaskIds()); - final Map> activeTasksToCreate = new HashMap<>(activeTasks); - final Map> standbyTasksToCreate = new HashMap<>(standbyTasks); - final Set tasksToRecycle = new HashSet<>(); - final Set dirtyTasks = new HashSet<>(); - builder.addSubscribedTopicsFromAssignment( activeTasks.values().stream().flatMap(Collection::stream).collect(Collectors.toList()), logPrefix ); - // first rectify all existing tasks final LinkedHashMap taskCloseExceptions = new LinkedHashMap<>(); + final Map> activeTasksToCreate = new HashMap<>(activeTasks); + final Map> standbyTasksToCreate = new HashMap<>(standbyTasks); + final LinkedList tasksToClose = new LinkedList<>(); + final Set tasksToRecycle = new HashSet<>(); + final Set dirtyTasks = new HashSet<>(); + + // first rectify all existing tasks for (final Task task : tasks.values()) { if (activeTasks.containsKey(task.id()) && task.isActive()) { updateInputPartitionsAndResume(task, activeTasks.get(task.id())); @@ -235,46 +236,49 @@ public void handleAssignment(final Map> activeTasks, } else if (standbyTasks.containsKey(task.id()) && !task.isActive()) { updateInputPartitionsAndResume(task, standbyTasks.get(task.id())); standbyTasksToCreate.remove(task.id()); - // check for tasks that were owned previously but have changed active/standby status } else if (activeTasks.containsKey(task.id()) || standbyTasks.containsKey(task.id())) { + // check for tasks that were owned previously but have changed active/standby status tasksToRecycle.add(task); } else { - try { - if (task.commitNeeded()) { - if (task.isActive()) { - log.error("Active task {} was revoked and should have already been committed", task.id()); - throw new IllegalStateException("Revoked active task was not committed during handleRevocation"); - } else { - task.suspend(); - task.prepareCommit(); - task.postCommit(); - } + tasksToClose.add(task); + } + } + + for (final Task task : tasksToClose) { + try { + task.suspend(); // Should be a no-op for active tasks, unless we hit an exception during handleRevocation + if (task.commitNeeded()) { + if (task.isActive()) { + log.error("Active task {} was revoked and should have already been committed", task.id()); + throw new IllegalStateException("Revoked active task was not committed during handleRevocation"); + } else { + task.prepareCommit(); + task.postCommit(); } - completeTaskCloseClean(task); - cleanUpTaskProducer(task, taskCloseExceptions); - tasks.remove(task.id()); - } catch (final RuntimeException e) { - final String uncleanMessage = String.format( - "Failed to close task %s cleanly. Attempting to close remaining tasks before re-throwing:", - task.id()); - log.error(uncleanMessage, e); - taskCloseExceptions.put(task.id(), e); - // We've already recorded the exception (which is the point of clean). - // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. - dirtyTasks.add(task); } + completeTaskCloseClean(task); + cleanUpTaskProducer(task, taskCloseExceptions); + tasks.remove(task.id()); + } catch (final RuntimeException e) { + final String uncleanMessage = String.format( + "Failed to close task %s cleanly. Attempting to close remaining tasks before re-throwing:", + task.id()); + log.error(uncleanMessage, e); + taskCloseExceptions.put(task.id(), e); + // We've already recorded the exception (which is the point of clean). + // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. + dirtyTasks.add(task); } } for (final Task oldTask : tasksToRecycle) { final Task newTask; try { + oldTask.suspend(); // Should be a no-op for active tasks, unless we hit an exception during handleRevocation if (oldTask.isActive()) { - // If active, the task should have already been suspended and committed during handleRevocation final Set partitions = standbyTasksToCreate.remove(oldTask.id()); newTask = standbyTaskCreator.createStandbyTaskFromActive((StreamTask) oldTask, partitions); } else { - oldTask.suspend(); final Set partitions = activeTasksToCreate.remove(oldTask.id()); newTask = activeTaskCreator.createActiveTaskFromStandby((StandbyTask) oldTask, partitions, mainConsumer); } @@ -425,6 +429,11 @@ boolean tryToCompleteRestoration() { } /** + * Handle the revoked partitions and prepare for closing the associated tasks in {@link #handleAssignment(Map, Map)} + * We should commit the revoked tasks now as we will not officially own them anymore when {@link #handleAssignment(Map, Map)} + * is called. Note that only active task partitions are passed in from the rebalance listener, so we only need to + * consider/commit active tasks here + * * @throws TaskMigratedException if the task producer got fenced (EOS only) */ void handleRevocation(final Collection revokedPartitions) { 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 8784cf1b28599..3f4b410358c07 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 @@ -57,6 +57,9 @@ import static org.apache.kafka.common.utils.Utils.mkEntry; import static org.apache.kafka.common.utils.Utils.mkMap; import static org.apache.kafka.common.utils.Utils.mkProperties; +import static org.apache.kafka.streams.processor.internals.Task.State.CREATED; +import static org.apache.kafka.streams.processor.internals.Task.State.RUNNING; +import static org.apache.kafka.streams.processor.internals.Task.State.SUSPENDED; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertEquals; @@ -139,7 +142,7 @@ public void cleanup() throws IOException { try { task.suspend(); } catch (final IllegalStateException maybeSwallow) { - if (!maybeSwallow.getMessage().startsWith("Invalid transition from CLOSED to SUSPENDED")) { + if (!maybeSwallow.getMessage().startsWith("Illegal state CLOSED while suspending standby task")) { throw maybeSwallow; } } @@ -171,16 +174,16 @@ public void shouldTransitToRunningAfterInitialization() { task = createStandbyTask(); - assertEquals(Task.State.CREATED, task.state()); + assertEquals(CREATED, task.state()); task.initializeIfNeeded(); - assertEquals(Task.State.RUNNING, task.state()); + assertEquals(RUNNING, task.state()); // initialize should be idempotent task.initializeIfNeeded(); - assertEquals(Task.State.RUNNING, task.state()); + assertEquals(RUNNING, task.state()); EasyMock.verify(stateManager); } @@ -263,7 +266,7 @@ public void shouldNotThrowFromStateManagerCloseInCloseDirty() { } @Test - public void shouldCommitOnCloseClean() { + public void shouldSuspendAndCommitBeforeCloseClean() { stateManager.close(); EasyMock.expectLastCall(); stateManager.checkpoint(EasyMock.eq(Collections.emptyMap())); @@ -288,6 +291,17 @@ public void shouldCommitOnCloseClean() { EasyMock.verify(stateManager); } + @Test + public void shouldRequireSuspendingCreatedTasksBeforeClose() { + EasyMock.replay(stateManager); + task = createStandbyTask(); + assertThat(task.state(), equalTo(CREATED)); + assertThrows(IllegalStateException.class, () -> task.closeClean()); + + task.suspend(); + task.closeClean(); + } + @Test public void shouldOnlyNeedCommitWhenChangelogOffsetChanged() { EasyMock.expect(stateManager.changelogPartitions()).andReturn(Collections.singleton(partition)).anyTimes(); @@ -355,7 +369,7 @@ public void shouldThrowOnCloseCleanCheckpointError() { task.prepareCommit(); assertThrows(RuntimeException.class, task::postCommit); - assertEquals(Task.State.RUNNING, task.state()); + assertEquals(RUNNING, task.state()); final double expectedCloseTaskMetric = 0.0; verifyCloseTaskMetric(expectedCloseTaskMetric, streamsMetrics, metricName); @@ -376,6 +390,7 @@ public void shouldCloseStateManagerOnTaskCreated() { final MetricName metricName = setupCloseTaskMetric(); task = createStandbyTask(); + task.suspend(); task.closeDirty(); @@ -405,6 +420,7 @@ public void shouldDeleteStateDirOnTaskCreatedAndEosAlphaUncleanClose() { ))); task = createStandbyTask(); + task.suspend(); task.closeDirty(); @@ -435,6 +451,7 @@ public void shouldDeleteStateDirOnTaskCreatedAndEosBetaUncleanClose() { task = createStandbyTask(); + task.suspend(); task.closeDirty(); final double expectedCloseTaskMetric = 1.0; @@ -447,20 +464,40 @@ public void shouldDeleteStateDirOnTaskCreatedAndEosBetaUncleanClose() { @Test public void shouldRecycleTask() { - stateManager.flush(); - EasyMock.expectLastCall(); stateManager.recycle(); - EasyMock.expectLastCall(); EasyMock.replay(stateManager); task = createStandbyTask(); + assertThrows(IllegalStateException.class, () -> task.closeAndRecycleState()); // CREATED + task.initializeIfNeeded(); + assertThrows(IllegalStateException.class, () -> task.closeAndRecycleState()); // RUNNING - task.closeAndRecycleState(); + task.suspend(); + task.closeAndRecycleState(); // SUSPENDED EasyMock.verify(stateManager); } + @Test + public void shouldAlwaysSuspendCreatedTasks() { + EasyMock.replay(stateManager); + task = createStandbyTask(); + assertThat(task.state(), equalTo(CREATED)); + task.suspend(); + assertThat(task.state(), equalTo(SUSPENDED)); + } + + @Test + public void shouldAlwaysSuspendRunningTasks() { + EasyMock.replay(stateManager); + task = createStandbyTask(); + task.initializeIfNeeded(); + assertThat(task.state(), equalTo(RUNNING)); + task.suspend(); + assertThat(task.state(), equalTo(SUSPENDED)); + } + private StandbyTask createStandbyTask() { final ThreadCache cache = new ThreadCache( 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 0426f689ae4c8..10eaee7889d42 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 @@ -87,6 +87,10 @@ import static org.apache.kafka.common.utils.Utils.mkProperties; import static org.apache.kafka.common.utils.Utils.mkSet; import static org.apache.kafka.streams.processor.internals.StreamTask.encodeTimestamp; +import static org.apache.kafka.streams.processor.internals.Task.State.CREATED; +import static org.apache.kafka.streams.processor.internals.Task.State.RESTORING; +import static org.apache.kafka.streams.processor.internals.Task.State.RUNNING; +import static org.apache.kafka.streams.processor.internals.Task.State.SUSPENDED; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.THREAD_ID_TAG; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.THREAD_ID_TAG_0100_TO_24; import static org.apache.kafka.test.StreamsTestUtils.getMetricByNameFilterByTags; @@ -329,14 +333,14 @@ public void shouldTransitToRestoringThenRunningAfterCreation() throws IOExceptio task.initializeIfNeeded(); - assertEquals(Task.State.RESTORING, task.state()); + assertEquals(RESTORING, task.state()); assertFalse(source1.initialized); assertFalse(source2.initialized); // initialize should be idempotent task.initializeIfNeeded(); - assertEquals(Task.State.RESTORING, task.state()); + assertEquals(RESTORING, task.state()); task.completeRestoration(); @@ -958,6 +962,7 @@ public void shouldCommitConsumerPositionIfRecordQueueIsEmpty() { @Test public void shouldFailOnCommitIfTaskIsClosed() { task = createStatelessTask(createConfig(false, "0"), StreamsConfig.METRICS_LATEST); + task.suspend(); task.transitionTo(Task.State.CLOSED); final IllegalStateException thrown = assertThrows( @@ -1299,7 +1304,7 @@ public void shouldReInitializeTopologyWhenResuming() throws IOException { task.resume(); - assertEquals(Task.State.RESTORING, task.state()); + assertEquals(RESTORING, task.state()); assertFalse(source1.initialized); assertFalse(source2.initialized); @@ -1506,6 +1511,7 @@ public void shouldThrowIfCommittingOnIllegalState() { task = createStatelessTask(createConfig(false, "100"), StreamsConfig.METRICS_LATEST); assertThrows(IllegalStateException.class, task::prepareCommit); + task.transitionTo(Task.State.SUSPENDED); task.transitionTo(Task.State.CLOSED); assertThrows(IllegalStateException.class, task::prepareCommit); } @@ -1515,6 +1521,7 @@ public void shouldThrowIfPostCommittingOnIllegalState() { task = createStatelessTask(createConfig(false, "100"), StreamsConfig.METRICS_LATEST); assertThrows(IllegalStateException.class, task::postCommit); + task.transitionTo(Task.State.SUSPENDED); task.transitionTo(Task.State.CLOSED); assertThrows(IllegalStateException.class, task::postCommit); } @@ -1695,7 +1702,7 @@ public void shouldThrowOnCloseCleanFlushError() { assertThrows(ProcessorStateException.class, task::prepareCommit); - assertEquals(Task.State.RESTORING, task.state()); + assertEquals(RESTORING, task.state()); final double expectedCloseTaskMetric = 0.0; verifyCloseTaskMetric(expectedCloseTaskMetric, streamsMetrics, metricName); @@ -1788,29 +1795,56 @@ public void shouldUpdatePartitions() { } @Test - public void shouldRecycleTask() { - EasyMock.expect(recordCollector.offsets()).andReturn(Collections.emptyMap()).anyTimes(); - recordCollector.flush(); - EasyMock.expectLastCall(); - stateManager.flush(); - EasyMock.expectLastCall(); - stateManager.checkpoint(Collections.emptyMap()); - EasyMock.expectLastCall(); + public void shouldOnlyRecycleSuspendedTasks() { stateManager.recycle(); - EasyMock.expectLastCall(); recordCollector.close(); - EasyMock.expectLastCall(); EasyMock.replay(stateManager, recordCollector); task = createStatefulTask(createConfig(false, "100"), true); + assertThrows(IllegalStateException.class, () -> task.closeAndRecycleState()); // CREATED + task.initializeIfNeeded(); + assertThrows(IllegalStateException.class, () -> task.closeAndRecycleState()); // RESTORING + task.completeRestoration(); + assertThrows(IllegalStateException.class, () -> task.closeAndRecycleState()); // RUNNING - task.closeAndRecycleState(); + task.suspend(); + task.closeAndRecycleState(); // SUSPENDED EasyMock.verify(stateManager, recordCollector); } + @Test + public void shouldAlwaysSuspendCreatedTasks() { + EasyMock.replay(stateManager); + task = createStatefulTask(createConfig(false, "100"), true); + assertThat(task.state(), equalTo(CREATED)); + task.suspend(); + assertThat(task.state(), equalTo(SUSPENDED)); + } + + @Test + public void shouldAlwaysSuspendRestoringTasks() { + EasyMock.replay(stateManager); + task = createStatefulTask(createConfig(false, "100"), true); + task.initializeIfNeeded(); + assertThat(task.state(), equalTo(RESTORING)); + task.suspend(); + assertThat(task.state(), equalTo(SUSPENDED)); + } + + @Test + public void shouldAlwaysSuspendRunningTasks() { + EasyMock.replay(stateManager); + task = createFaultyStatefulTask(createConfig(false, "100")); + task.initializeIfNeeded(); + task.completeRestoration(); + assertThat(task.state(), equalTo(RUNNING)); + assertThrows(RuntimeException.class , () -> task.suspend()); + assertThat(task.state(), equalTo(SUSPENDED)); + } + private StreamTask createOptimizedStatefulTask(final StreamsConfig config, final Consumer consumer) { final StateStore stateStore = new MockKeyValueStore(storeName, true); 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 76136b9b50cee..f7094630661d8 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 @@ -53,7 +53,6 @@ import org.junit.Before; import org.junit.Rule; import org.junit.Test; -import org.junit.function.ThrowingRunnable; import org.junit.rules.TemporaryFolder; import org.junit.runner.RunWith; @@ -84,7 +83,6 @@ import static org.apache.kafka.common.utils.Utils.mkSet; import static org.easymock.EasyMock.anyObject; import static org.easymock.EasyMock.anyString; -import static org.easymock.EasyMock.checkOrder; import static org.easymock.EasyMock.eq; import static org.easymock.EasyMock.expect; import static org.easymock.EasyMock.expectLastCall; @@ -421,10 +419,11 @@ public void shouldCloseActiveUnassignedSuspendedTasksWhenClosingRevokedTasks() { } @Test - public void shouldCloseDirtyActiveUnassignedSuspendedTasksWhenErrorCommittingRevokedTask() { + public void shouldCloseDirtyActiveUnassignedTasksWhenErrorSuspendingTask() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true) { @Override - public Map prepareCommit() { + public void suspend() { + transitionTo(State.SUSPENDED); throw new RuntimeException("KABOOM!"); } }; @@ -916,7 +915,7 @@ public void completeRestoration() { } @Test - public void shouldSuspendActiveTasks() { + public void shouldSuspendActiveTasksDuringRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets); @@ -937,10 +936,11 @@ public void shouldSuspendActiveTasks() { } @Test - public void shouldCommitAllActiveTasksTheNeedCommittingOnHandleAssignmentIfOneTaskClosed() { + public void shouldCommitAllActiveTasksThatNeedCommittingOnHandleRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); + task00.setCommitNeeded(); final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); final Map offsets01 = singletonMap(t1p1, new OffsetAndMetadata(1L, null)); @@ -986,8 +986,7 @@ public void shouldCommitAllActiveTasksTheNeedCommittingOnHandleAssignmentIfOneTa assertThat(task02.state(), is(Task.State.RUNNING)); assertThat(task10.state(), is(Task.State.RUNNING)); - assignmentActive.remove(taskId00); - taskManager.handleAssignment(assignmentActive, assignmentStandby); + taskManager.handleRevocation(taskId00Partitions); assertThat(task00.commitNeeded, is(false)); assertThat(task01.commitNeeded, is(false)); @@ -1055,58 +1054,11 @@ public void shouldNotCommitOnHandleAssignmentIfOnlyStandbyTaskClosed() { } @Test - public void shouldCleanupAnyTasksClosedAsDirtyAfterCommitException() { - final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); - final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); - task00.setCommittableOffsetsAndMetadata(offsets00); - - final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); - final Map offsets01 = singletonMap(t1p1, new OffsetAndMetadata(1L, null)); - task01.setCommittableOffsetsAndMetadata(offsets01); - task01.setCommitNeeded(); - - task01.setChangelogOffsets(singletonMap(t1p1, 0L)); - - final StateMachineTask task02 = new StateMachineTask(taskId02, taskId02Partitions, true); - final Map offsets02 = singletonMap(t1p2, new OffsetAndMetadata(2L, null)); - task02.setCommittableOffsetsAndMetadata(offsets02); - - final Map expectedCommittedOffsets = new HashMap<>(); - expectedCommittedOffsets.putAll(offsets00); - expectedCommittedOffsets.putAll(offsets01); - - final Map> assignmentActive = mkMap( - mkEntry(taskId00, taskId00Partitions), - mkEntry(taskId01, taskId01Partitions), - mkEntry(taskId02, taskId02Partitions) - ); - - expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))) - .andReturn(asList(task00, task01, task02)); - activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(EasyMock.anyObject(TaskId.class)); - expectLastCall().anyTimes(); - - consumer.commitSync(expectedCommittedOffsets); - expectLastCall().andThrow(new RuntimeException("Something went wrong!")); - - replay(activeTaskCreator, standbyTaskCreator, consumer, changeLogReader); - - taskManager.handleAssignment(assignmentActive, emptyMap()); - - assignmentActive.remove(taskId00); - assertThrows( - RuntimeException.class, - () -> taskManager.handleAssignment(assignmentActive, emptyMap()) - ); - - verify(changeLogReader); - } - - @Test - public void shouldCommitAllActiveTasksTheNeedCommittingOnRevocation() { + public void shouldCommitAllActiveTasksThatNeedCommittingOnRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); + task00.setCommitNeeded(); final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); final Map offsets01 = singletonMap(t1p1, new OffsetAndMetadata(1L, null)); @@ -1159,17 +1111,21 @@ public void shouldCommitAllActiveTasksTheNeedCommittingOnRevocation() { } @Test - public void shouldNotCommitCreatedTasksOnSuspend() { + public void shouldNotCommitCreatedTasksOnRevocationOrClosure() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); expect(activeTaskCreator.createTasks(anyObject(), eq(taskId00Assignment))).andReturn(singletonList(task00)); + activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(eq(taskId00)); replay(activeTaskCreator, consumer, changeLogReader); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(task00.state(), is(Task.State.CREATED)); taskManager.handleRevocation(taskId00Partitions); - assertThat(task00.state(), is(Task.State.CREATED)); + assertThat(task00.state(), is(Task.State.SUSPENDED)); + + taskManager.handleAssignment(emptyMap(), emptyMap()); + assertThat(task00.state(), is(Task.State.CLOSED)); } @Test @@ -1423,91 +1379,64 @@ public Collection changelogPartitions() { } @Test - public void shouldCloseActiveTasksDirtyAndPropagatePrepareCommitException() { + public void shouldOnlyCommitRevokedStandbyTaskAndPropagatePrepareCommitException() { setUpTaskManager(StreamThread.ProcessingMode.EXACTLY_ONCE_ALPHA); - final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, true); + final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, false); - final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true) { + final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, false) { @Override public Map prepareCommit() { throw new RuntimeException("task 0_1 prepare commit boom!"); } }; - - task01.setCommittableOffsetsAndMetadata(singletonMap(t1p1, new OffsetAndMetadata(0L, null))); task01.setCommitNeeded(); - final StateMachineTask task02 = new StateMachineTask(taskId02, taskId02Partitions, true); - final Map offsetsT02 = singletonMap(t1p2, new OffsetAndMetadata(1L, null)); - - task02.setCommittableOffsetsAndMetadata(offsetsT02); - task02.setCommitNeeded(); - taskManager.tasks().put(taskId00, task00); taskManager.tasks().put(taskId01, task01); - taskManager.tasks().put(taskId02, task02); - - checkOrder(activeTaskCreator, false); - - activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId01); - expectLastCall(); - - activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId02); - expectLastCall(); - - replay(activeTaskCreator); final RuntimeException thrown = assertThrows(RuntimeException.class, - () -> taskManager.handleAssignment(mkMap(mkEntry(taskId00, taskId00Partitions), - mkEntry(taskId01, taskId01Partitions)), Collections.emptyMap())); + () -> taskManager.handleAssignment( + Collections.emptyMap(), + singletonMap(taskId00, taskId00Partitions) + )); assertThat(thrown.getCause().getMessage(), is("task 0_1 prepare commit boom!")); assertThat(task00.state(), is(Task.State.CREATED)); assertThat(task01.state(), is(Task.State.CLOSED)); - assertThat(task02.state(), is(Task.State.CLOSED)); // All the tasks involving in the commit should already be removed. assertThat(taskManager.tasks(), is(Collections.singletonMap(taskId00, task00))); - - verify(activeTaskCreator); } @Test - public void shouldCloseActiveTasksDirtyAndPropagateCommitException() { + public void shouldCloseActiveTasksDirtyAndPropagateSuspendException() { setUpTaskManager(StreamThread.ProcessingMode.EXACTLY_ONCE_ALPHA); final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, true); - final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); - task01.setCommittableOffsetsAndMetadata(singletonMap(t1p1, new OffsetAndMetadata(0L, null))); - task01.setCommitNeeded(); + final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true) { + @Override + public void suspend() { + transitionTo(State.SUSPENDED); + throw new RuntimeException("task 0_1 suspend boom!"); + } + }; final StateMachineTask task02 = new StateMachineTask(taskId02, taskId02Partitions, true); - final Map offsetsT02 = singletonMap(t1p2, new OffsetAndMetadata(1L, null)); - - task02.setCommittableOffsetsAndMetadata(offsetsT02); - task02.setCommitNeeded(); taskManager.tasks().put(taskId00, task00); taskManager.tasks().put(taskId01, task01); taskManager.tasks().put(taskId02, task02); - expect(activeTaskCreator.streamsProducerForTask(taskId01)).andThrow(new RuntimeException("task 0_1 producer boom!")); - - checkOrder(activeTaskCreator, false); - - activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId01); - expectLastCall(); - activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId02); - expectLastCall(); + activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId01); replay(activeTaskCreator); final RuntimeException thrown = assertThrows(RuntimeException.class, () -> taskManager.handleAssignment(mkMap(mkEntry(taskId00, taskId00Partitions)), Collections.emptyMap())); - assertThat(thrown.getCause().getMessage(), is("task 0_1 producer boom!")); + assertThat(thrown.getCause().getMessage(), is("task 0_1 suspend boom!")); assertThat(task00.state(), is(Task.State.CREATED)); assertThat(task01.state(), is(Task.State.CLOSED)); @@ -2632,36 +2561,30 @@ public void shouldFailOnCommitFatal() { } @Test - public void shouldNotCloseTasksIfCommittingFailsDuringRevocation() { - shouldNotCloseTaskIfCommitFailsDuringAction(() -> taskManager.handleRevocation(singletonList(t1p0))); - } - - @Test - public void shouldNotCloseTasksIfCommittingFailsDuringShutdown() { - shouldNotCloseTaskIfCommitFailsDuringAction(() -> taskManager.shutdown(true)); - } - - private void shouldNotCloseTaskIfCommitFailsDuringAction(final ThrowingRunnable action) { - final Map offsets = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); + public void shouldNotSuspendOrCommitAllTasksIfSuspendingFailsDuringRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true) { @Override - public Map prepareCommit() { - return offsets; + public void suspend() { + super.suspend(); + throw new RuntimeException("KABOOM!"); } }; + final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); - expect(activeTaskCreator.createTasks(anyObject(), eq(taskId00Assignment))) - .andReturn(singletonList(task00)); - consumer.commitSync(offsets); - expectLastCall().andThrow(new RuntimeException("KABOOM!")); + final Map> assignment = new HashMap<>(taskId00Assignment); + assignment.putAll(taskId01Assignment); + expect(activeTaskCreator.createTasks(anyObject(), eq(assignment))) + .andReturn(asList(task00, task01)); replay(activeTaskCreator, consumer); - taskManager.handleAssignment(taskId00Assignment, Collections.emptyMap()); + taskManager.handleAssignment(assignment, Collections.emptyMap()); - final RuntimeException thrown = assertThrows(RuntimeException.class, action); + final RuntimeException thrown = assertThrows(RuntimeException.class, + () ->taskManager.handleRevocation(asList(t1p0, t1p1))); assertThat(thrown.getMessage(), is("KABOOM!")); - assertThat(task00.state(), is(Task.State.CREATED)); + assertThat(task00.state(), is(Task.State.SUSPENDED)); + assertThat(task01.state(), is(Task.State.CREATED)); } private static void expectRestoreToBeCompleted(final Consumer consumer, @@ -2782,7 +2705,9 @@ public void postCommit() { @Override public void suspend() { - if (state() == State.RUNNING) { + if (state() == State.CLOSED) { + throw new IllegalStateException("Illegal state " + state() + " while suspending active task " + id); + } else { transitionTo(State.SUSPENDED); } } From 7b08e28d940ae822b05da54460064c261b56e0ed Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Jun 2020 12:27:56 -0700 Subject: [PATCH 06/14] always close all tasks during shutdown --- .../streams/processor/internals/TaskManager.java | 15 ++++++++++----- .../processor/internals/StreamTaskTest.java | 2 +- .../processor/internals/TaskManagerTest.java | 5 +++-- 3 files changed, 14 insertions(+), 8 deletions(-) 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 1e843c41eccf6..c7e955815b89e 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 @@ -678,15 +678,20 @@ void shutdown(final boolean clean) { } } - if (clean && !consumedOffsetsAndMetadataPerTask.isEmpty()) { - commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); + try { + if (clean && !consumedOffsetsAndMetadataPerTask.isEmpty()) { + commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); + } + for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { + final Task task = tasks.get(taskId); + task.postCommit(); + } + } catch (final RuntimeException e) { + firstException.compareAndSet(null, e); } for (final Task task : tasksToClose) { try { - if (consumedOffsetsAndMetadataPerTask.containsKey(task.id())) { - task.postCommit(); - } completeTaskCloseClean(task); } catch (final RuntimeException e) { firstException.compareAndSet(null, e); 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 10eaee7889d42..9fef51e427961 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 @@ -1841,7 +1841,7 @@ public void shouldAlwaysSuspendRunningTasks() { task.initializeIfNeeded(); task.completeRestoration(); assertThat(task.state(), equalTo(RUNNING)); - assertThrows(RuntimeException.class , () -> task.suspend()); + assertThrows(RuntimeException.class, () -> task.suspend()); assertThat(task.state(), equalTo(SUSPENDED)); } 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 f7094630661d8..094e2c4db2cc8 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 @@ -2579,8 +2579,9 @@ public void suspend() { taskManager.handleAssignment(assignment, Collections.emptyMap()); - final RuntimeException thrown = assertThrows(RuntimeException.class, - () ->taskManager.handleRevocation(asList(t1p0, t1p1))); + final RuntimeException thrown = assertThrows( + RuntimeException.class, + () -> taskManager.handleRevocation(asList(t1p0, t1p1))); assertThat(thrown.getMessage(), is("KABOOM!")); assertThat(task00.state(), is(Task.State.SUSPENDED)); From dd03646b9c8c01602219d7afb25730df5663a7d0 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Jun 2020 15:30:12 -0700 Subject: [PATCH 07/14] only commit all for eos-beta --- .../processor/internals/TaskManager.java | 34 ++++++++++++++----- 1 file changed, 26 insertions(+), 8 deletions(-) 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 c7e955815b89e..0e3c338502470 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 @@ -434,31 +434,49 @@ boolean tryToCompleteRestoration() { * is called. Note that only active task partitions are passed in from the rebalance listener, so we only need to * consider/commit active tasks here * + * If eos-beta is used, we must commit ALL tasks. Otherwise, we can just commit those (active) tasks which are revoked + * * @throws TaskMigratedException if the task producer got fenced (EOS only) */ void handleRevocation(final Collection revokedPartitions) { final Set remainingRevokedPartitions = new HashSet<>(revokedPartitions); - final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); + final Set tasksToCommit = new HashSet<>(); + final Set additionalTasksForCommitting = new HashSet<>(); + for (final Task task : activeTaskIterable()) { if (remainingRevokedPartitions.containsAll(task.inputPartitions())) { task.suspend(); - } - if (task.commitNeeded()) { - final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.commitNeeded()) { + tasksToCommit.add(task); } + } else if (task.commitNeeded()) { + additionalTasksForCommitting.add(task); } remainingRevokedPartitions.removeAll(task.inputPartitions()); } + // If using eos-beta, if we must commit any task then we must commit all of them + // TODO: when KAFKA-9450 is done this will be less expensive, and we can simplify by always committing everything + if (processingMode == EXACTLY_ONCE_BETA && !tasksToCommit.isEmpty()) { + tasksToCommit.addAll(additionalTasksForCommitting); + } + + final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); + for (final Task task : tasksToCommit) { + final Map committableOffsets = task.prepareCommit(); + if (!committableOffsets.isEmpty()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + } else { + log.warn("Task {} claimed to need a commit but had no committable consumed offsets", task.id()); + } + } + if (!consumedOffsetsAndMetadataPerTask.isEmpty()) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); } - for (final TaskId committedTaskId : consumedOffsetsAndMetadataPerTask.keySet()) { - final Task task = tasks.get(committedTaskId); + for (final Task task : tasksToCommit) { task.postCommit(); } From 1f989fc78ffa6b89cf50775767e563a64120460a Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Jun 2020 19:50:56 -0700 Subject: [PATCH 08/14] addressing CR feedback --- .../processor/internals/StandbyTask.java | 9 +++++++-- .../processor/internals/StreamTask.java | 17 +++++++++++------ .../streams/processor/internals/Task.java | 4 ++-- .../processor/internals/TaskManager.java | 19 +++++++++++++++---- 4 files changed, 35 insertions(+), 14 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 20ff5f6715531..b407ad1f4f2ed 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 @@ -111,12 +111,17 @@ public void completeRestoration() { @Override public void suspend() { - log.trace("No-op suspend with state {}", state()); switch (state()) { case CREATED: case RUNNING: - case SUSPENDED: transitionTo(State.SUSPENDED); + log.info("Suspended {}", state()); + + break; + + case SUSPENDED: + log.info("Skip suspending since state is {}", state()); + break; case RESTORING: 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 07c3c872151bb..e46eec2c75683 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 @@ -251,7 +251,6 @@ public void suspend() { switch (state()) { case CREATED: case RESTORING: - case SUSPENDED: transitionTo(State.SUSPENDED); log.info("Suspended {}", state()); @@ -268,6 +267,12 @@ public void suspend() { break; + case SUSPENDED: + log.info("Skip suspending since state is {}", state()); + + break; + + case CLOSED: throw new IllegalStateException("Illegal state " + state() + " while suspending active task " + id); @@ -470,7 +475,6 @@ public void update(final Set topicPartitions, final Map * the following order must be followed: - * 1. checkpoint the state manager -- even if we crash before this step, EOS is still guaranteed + * 1. commit/checkpoint the state manager -- even if we crash before this step, EOS is still guaranteed * 2. then if we are closing on EOS and dirty, wipe out the state store directory * 3. finally release the state manager lock * */ private void close(final boolean clean) { - if (clean) { - executeAndMaybeSwallow(true, this::writeCheckpointIfNeed, "state manager checkpoint", log); + if (clean && commitNeeded && checkpoint != null) { + log.debug("Tried to close clean but there was an active scheduled checkpoint, this means we failed to " + + "commit and should close as dirty instead"); + throw new StreamsException("Tried to close dirty task as clean"); } - switch (state()) { case SUSPENDED: // first close state manager (which is idempotent) then close the record collector 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 a48d4ae9c77cd..0200870b7aa9a 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 @@ -56,7 +56,7 @@ public interface Task { * | | | | * | v | | * | +------+--------+ | | - * +---->| Suspended (3) | ----+ | //TODO Suspended(3) could be removed after we've stable on KIP-429 + * +---> | Suspended (3) | ----+ | //TODO Suspended(3) could be removed after we've stable on KIP-429 * +------+--------+ | * | | * v | @@ -69,7 +69,7 @@ enum State { CREATED(1, 3), // 0 RESTORING(2, 3), // 1 RUNNING(3), // 2 - SUSPENDED(1, 3, 4), // 3 + SUSPENDED(1, 4), // 3 CLOSED(0); // 4, we allow CLOSED to transit to CREATED to handle corrupted tasks private final Set validTransitions = new HashSet<>(); 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 0e3c338502470..8cec2cbad1283 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 @@ -224,7 +224,7 @@ public void handleAssignment(final Map> activeTasks, final Map> activeTasksToCreate = new HashMap<>(activeTasks); final Map> standbyTasksToCreate = new HashMap<>(standbyTasks); - final LinkedList tasksToClose = new LinkedList<>(); + final List tasksToClose = new LinkedList<>(); final Set tasksToRecycle = new HashSet<>(); final Set dirtyTasks = new HashSet<>(); @@ -444,11 +444,17 @@ void handleRevocation(final Collection revokedPartitions) { final Set tasksToCommit = new HashSet<>(); final Set additionalTasksForCommitting = new HashSet<>(); + final AtomicReference firstException = new AtomicReference<>(null); for (final Task task : activeTaskIterable()) { if (remainingRevokedPartitions.containsAll(task.inputPartitions())) { - task.suspend(); - if (task.commitNeeded()) { - tasksToCommit.add(task); + try { + task.suspend(); + if (task.commitNeeded()) { + tasksToCommit.add(task); + } + } catch (final RuntimeException e) { + log.error("Caught the following exception while trying to suspend revoked task " + task.id(), e); + firstException.compareAndSet(null, new StreamsException("Failed to suspend " + task.id(), e)); } } else if (task.commitNeeded()) { additionalTasksForCommitting.add(task); @@ -456,6 +462,11 @@ void handleRevocation(final Collection revokedPartitions) { remainingRevokedPartitions.removeAll(task.inputPartitions()); } + final RuntimeException exception = firstException.get(); + if (exception != null) { + throw exception; + } + // If using eos-beta, if we must commit any task then we must commit all of them // TODO: when KAFKA-9450 is done this will be less expensive, and we can simplify by always committing everything if (processingMode == EXACTLY_ONCE_BETA && !tasksToCommit.isEmpty()) { From cf21fb5d42e1a5b1521ae58b7f57839195a4791b Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 15 Jun 2020 16:19:22 -0700 Subject: [PATCH 09/14] get rid of checkpoint member --- .../processor/internals/StreamTask.java | 67 +++++-------------- .../processor/internals/TaskManager.java | 4 +- .../processor/internals/StreamTaskTest.java | 5 +- 3 files changed, 21 insertions(+), 55 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 e46eec2c75683..b7ccb62523509 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,6 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator, private boolean commitNeeded = false; private boolean commitRequested = false; - private Map checkpoint = null; - public StreamTask(final TaskId id, final Set partitions, final ProcessorTopology topology, @@ -343,7 +341,6 @@ public Map prepareCommit() { case RUNNING: case RESTORING: case SUSPENDED: - maybeScheduleCheckpoint(); stateMgr.flush(); recordCollector.flush(); @@ -410,6 +407,9 @@ private Map committableOffsetsAndMetadata() { return committableOffsets; } + /** + * This should only be called if the attempted commit succeeded for this task + */ @Override public void postCommit() { commitRequested = false; @@ -417,25 +417,18 @@ public void postCommit() { switch (state()) { case RESTORING: - writeCheckpointIfNeed(); + case SUSPENDED: + maybeWriteCheckpoint(); break; case RUNNING: - if (!eosEnabled) { // if RUNNING, checkpoint only for non-eos - writeCheckpointIfNeed(); + if (!eosEnabled) { + maybeWriteCheckpoint(); } break; - case SUSPENDED: - writeCheckpointIfNeed(); - // we cannot `clear()` the `PartitionGroup` in `suspend()` already, but only after committing, - // because otherwise we loose the partition-time information - partitionGroup.clear(); - - break; - case CREATED: case CLOSED: throw new IllegalStateException("Illegal state " + state() + " while post committing active task " + id); @@ -491,6 +484,8 @@ public void closeAndRecycleState() { throw new IllegalStateException("Unknown state " + state() + " while recycling active task " + id); } + // we cannot `clear()` the `PartitionGroup` in `suspend()` already, but only after committing, + // because otherwise we loose the partition-time information partitionGroup.clear(); closeTaskSensor.record(); @@ -499,53 +494,21 @@ public void closeAndRecycleState() { log.info("Closed clean and recycled state"); } - private void maybeScheduleCheckpoint() { - switch (state()) { - case RESTORING: - case SUSPENDED: - this.checkpoint = checkpointableOffsets(); - - break; - - case RUNNING: - if (!eosEnabled) { - this.checkpoint = checkpointableOffsets(); - } - - break; - - case CREATED: - case CLOSED: - throw new IllegalStateException("Illegal state " + state() + " while scheduling checkpoint for active task " + id); - - default: - throw new IllegalStateException("Unknown state " + state() + " while scheduling checkpoint for active task " + id); - } - } - - private void writeCheckpointIfNeed() { + private void maybeWriteCheckpoint() { if (commitNeeded) { log.error("Tried to write a checkpoint with pending uncommitted data, should complete the commit first."); throw new IllegalStateException("A checkpoint should only be written if no commit is needed."); } - if (checkpoint != null) { - stateMgr.checkpoint(checkpoint); - checkpoint = null; - } + stateMgr.checkpoint(checkpointableOffsets()); } /** - *
-     * the following order must be followed:
-     *  1. commit/checkpoint the state manager -- even if we crash before this step, EOS is still guaranteed
-     *  2. then if we are closing on EOS and dirty, wipe out the state store directory
-     *  3. finally release the state manager lock
-     * 
+ * You must commit a task and checkpoint the state manager before closing as this will release the state dir lock */ private void close(final boolean clean) { - if (clean && commitNeeded && checkpoint != null) { - log.debug("Tried to close clean but there was an active scheduled checkpoint, this means we failed to " - + "commit and should close as dirty instead"); + if (clean && commitNeeded ) { + log.debug("Tried to close clean but there was an active scheduled checkpoint, this means we failed to" + + " commit and should close as dirty instead"); throw new StreamsException("Tried to close dirty task as clean"); } switch (state()) { 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 8cec2cbad1283..538e82c308014 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 @@ -246,7 +246,7 @@ public void handleAssignment(final Map> activeTasks, for (final Task task : tasksToClose) { try { - task.suspend(); // Should be a no-op for active tasks, unless we hit an exception during handleRevocation + task.suspend(); // Should be a no-op for active tasks since they're suspended in handleRevocation if (task.commitNeeded()) { if (task.isActive()) { log.error("Active task {} was revoked and should have already been committed", task.id()); @@ -274,11 +274,11 @@ public void handleAssignment(final Map> activeTasks, for (final Task oldTask : tasksToRecycle) { final Task newTask; try { - oldTask.suspend(); // Should be a no-op for active tasks, unless we hit an exception during handleRevocation if (oldTask.isActive()) { final Set partitions = standbyTasksToCreate.remove(oldTask.id()); newTask = standbyTaskCreator.createStandbyTaskFromActive((StreamTask) oldTask, partitions); } else { + oldTask.suspend(); // Only need to suspend transitioning standbys, actives should be suspended already final Set partitions = activeTasksToCreate.remove(oldTask.id()); newTask = activeTaskCreator.createActiveTaskFromStandby((StandbyTask) oldTask, partitions, mainConsumer); } 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 9fef51e427961..7a2cf7ad49579 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 @@ -1615,6 +1615,7 @@ public void shouldNotCommitOnCloseRestoring() { task.completeRestoration(); task.suspend(); task.prepareCommit(); + task.postCommit(); task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -1640,6 +1641,7 @@ public void shouldCommitOnCloseClean() { task.completeRestoration(); task.suspend(); task.prepareCommit(); + task.postCommit(); task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -1669,6 +1671,7 @@ public void shouldSwallowExceptionOnCloseCleanError() { task.suspend(); task.prepareCommit(); + task.postCommit(); assertThrows(ProcessorStateException.class, () -> task.closeClean()); final double expectedCloseTaskMetric = 0.0; @@ -1676,7 +1679,7 @@ public void shouldSwallowExceptionOnCloseCleanError() { EasyMock.verify(stateManager); EasyMock.reset(stateManager); - EasyMock.expect(stateManager.changelogPartitions()).andReturn(Collections.singleton(changelogPartition)).anyTimes(); + EasyMock.expect(stateManager.changelogPartitions()).andStubReturn(Collections.singleton(changelogPartition)); stateManager.close(); EasyMock.expectLastCall(); EasyMock.replay(stateManager); From d2ba0df1befe929672850755d1610427c6234e57 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 15 Jun 2020 16:52:22 -0700 Subject: [PATCH 10/14] fixing up last few TM tests --- .../processor/internals/StreamTask.java | 4 +- .../processor/internals/TaskManagerTest.java | 130 +++++++++--------- 2 files changed, 65 insertions(+), 69 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 b7ccb62523509..ef1cc928a9702 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 @@ -506,9 +506,9 @@ private void maybeWriteCheckpoint() { * You must commit a task and checkpoint the state manager before closing as this will release the state dir lock */ private void close(final boolean clean) { - if (clean && commitNeeded ) { + if (clean && commitNeeded) { log.debug("Tried to close clean but there was an active scheduled checkpoint, this means we failed to" - + " commit and should close as dirty instead"); + + " commit and should close as dirty instead"); throw new StreamsException("Tried to close dirty task as clean"); } switch (state()) { 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 094e2c4db2cc8..389d187f93ad0 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 @@ -42,6 +42,7 @@ import org.apache.kafka.streams.errors.TaskMigratedException; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; +import org.apache.kafka.streams.processor.internals.StreamThread.ProcessingMode; import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender; import org.apache.kafka.streams.state.internals.OffsetCheckpoint; @@ -423,7 +424,7 @@ public void shouldCloseDirtyActiveUnassignedTasksWhenErrorSuspendingTask() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true) { @Override public void suspend() { - transitionTo(State.SUSPENDED); + super.suspend(); throw new RuntimeException("KABOOM!"); } }; @@ -509,21 +510,7 @@ public void shouldReInitializeThreadProducerOnHandleLostAllIfEosBetaEnabled() { expectLastCall(); replay(activeTaskCreator); - final StreamsMetricsImpl streamsMetrics = - new StreamsMetricsImpl(new Metrics(), "clientId", StreamsConfig.METRICS_LATEST); - taskManager = new TaskManager( - changeLogReader, - UUID.randomUUID(), - "taskManagerTest", - streamsMetrics, - activeTaskCreator, - standbyTaskCreator, - topologyBuilder, - adminClient, - stateDirectory, - StreamThread.ProcessingMode.EXACTLY_ONCE_BETA - ); - taskManager.setMainConsumer(consumer); + setUpTaskManager(ProcessingMode.EXACTLY_ONCE_BETA); taskManager.handleLostAll(); @@ -936,7 +923,10 @@ public void shouldSuspendActiveTasksDuringRevocation() { } @Test - public void shouldCommitAllActiveTasksThatNeedCommittingOnHandleRevocation() { + public void shouldCommitAllActiveTasksThatNeedCommittingOnHandleRevocationWithEosBeta() { + final StreamsProducer producer = mock(StreamsProducer.class); + setUpTaskManager(ProcessingMode.EXACTLY_ONCE_BETA); + final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); @@ -970,11 +960,14 @@ public void shouldCommitAllActiveTasksThatNeedCommittingOnHandleRevocation() { expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))) .andReturn(asList(task00, task01, task02)); + expect(activeTaskCreator.threadProducer()).andReturn(producer); activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId00); - expectLastCall(); expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) .andReturn(singletonList(task10)); - consumer.commitSync(expectedCommittedOffsets); + + final ConsumerGroupMetadata groupMetadata = new ConsumerGroupMetadata("appId"); + expect(consumer.groupMetadata()).andReturn(groupMetadata); + producer.commitTransaction(expectedCommittedOffsets, groupMetadata); expectLastCall(); replay(activeTaskCreator, standbyTaskCreator, consumer, changeLogReader); @@ -995,37 +988,65 @@ public void shouldCommitAllActiveTasksThatNeedCommittingOnHandleRevocation() { } @Test - public void shouldNotCommitOnHandleAssignmentIfNoTaskClosed() { + public void shouldCommitOnlyRevokedActiveTasksThatNeedCommittingOnHandleRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); task00.setCommitNeeded(); + final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); + final Map offsets01 = singletonMap(t1p1, new OffsetAndMetadata(1L, null)); + task01.setCommittableOffsetsAndMetadata(offsets01); + task01.setCommitNeeded(); + + final StateMachineTask task02 = new StateMachineTask(taskId02, taskId02Partitions, true); + final Map offsets02 = singletonMap(t1p2, new OffsetAndMetadata(2L, null)); + task02.setCommittableOffsetsAndMetadata(offsets02); + final StateMachineTask task10 = new StateMachineTask(taskId10, taskId10Partitions, false); - final Map> assignmentActive = singletonMap(taskId00, taskId00Partitions); - final Map> assignmentStandby = singletonMap(taskId10, taskId10Partitions); + final Map expectedCommittedOffsets = new HashMap<>(); + expectedCommittedOffsets.putAll(offsets00); + + final Map> assignmentActive = mkMap( + mkEntry(taskId00, taskId00Partitions), + mkEntry(taskId01, taskId01Partitions), + mkEntry(taskId02, taskId02Partitions) + ); + final Map> assignmentStandby = mkMap( + mkEntry(taskId10, taskId10Partitions) + ); expectRestoreToBeCompleted(consumer, changeLogReader); - expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))).andReturn(singleton(task00)); - expect(standbyTaskCreator.createTasks(eq(assignmentStandby))).andReturn(singletonList(task10)); + expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))) + .andReturn(asList(task00, task01, task02)); + activeTaskCreator.closeAndRemoveTaskProducerIfNeeded(taskId00); + expectLastCall(); + expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) + .andReturn(singletonList(task10)); + consumer.commitSync(expectedCommittedOffsets); + expectLastCall(); replay(activeTaskCreator, standbyTaskCreator, consumer, changeLogReader); taskManager.handleAssignment(assignmentActive, assignmentStandby); assertThat(taskManager.tryToCompleteRestoration(), is(true)); assertThat(task00.state(), is(Task.State.RUNNING)); + assertThat(task01.state(), is(Task.State.RUNNING)); + assertThat(task02.state(), is(Task.State.RUNNING)); assertThat(task10.state(), is(Task.State.RUNNING)); - taskManager.handleAssignment(assignmentActive, assignmentStandby); + taskManager.handleRevocation(taskId00Partitions); - assertThat(task00.commitNeeded, is(true)); + assertThat(task00.commitNeeded, is(false)); + assertThat(task01.commitPrepared, is(false)); + assertThat(task02.commitPrepared, is(false)); assertThat(task10.commitPrepared, is(false)); } @Test - public void shouldNotCommitOnHandleAssignmentIfOnlyStandbyTaskClosed() { + public void shouldNotCommitOnHandleAssignmentIfNoTaskClosed() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); @@ -1048,66 +1069,39 @@ public void shouldNotCommitOnHandleAssignmentIfOnlyStandbyTaskClosed() { assertThat(task00.state(), is(Task.State.RUNNING)); assertThat(task10.state(), is(Task.State.RUNNING)); - taskManager.handleAssignment(assignmentActive, Collections.emptyMap()); + taskManager.handleAssignment(assignmentActive, assignmentStandby); assertThat(task00.commitNeeded, is(true)); + assertThat(task10.commitPrepared, is(false)); } @Test - public void shouldCommitAllActiveTasksThatNeedCommittingOnRevocation() { + public void shouldNotCommitOnHandleAssignmentIfOnlyStandbyTaskClosed() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); task00.setCommitNeeded(); - final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true); - final Map offsets01 = singletonMap(t1p1, new OffsetAndMetadata(1L, null)); - task01.setCommittableOffsetsAndMetadata(offsets01); - task01.setCommitNeeded(); - - final StateMachineTask task02 = new StateMachineTask(taskId02, taskId02Partitions, true); - final Map offsets02 = singletonMap(t1p2, new OffsetAndMetadata(2L, null)); - task02.setCommittableOffsetsAndMetadata(offsets02); - final StateMachineTask task10 = new StateMachineTask(taskId10, taskId10Partitions, false); - final Map expectedCommittedOffsets = new HashMap<>(); - expectedCommittedOffsets.putAll(offsets00); - expectedCommittedOffsets.putAll(offsets01); - - final Map> assignmentActive = mkMap( - mkEntry(taskId00, taskId00Partitions), - mkEntry(taskId01, taskId01Partitions), - mkEntry(taskId02, taskId02Partitions) - ); + final Map> assignmentActive = singletonMap(taskId00, taskId00Partitions); + final Map> assignmentStandby = singletonMap(taskId10, taskId10Partitions); - final Map> assignmentStandby = mkMap( - mkEntry(taskId10, taskId10Partitions) - ); expectRestoreToBeCompleted(consumer, changeLogReader); - expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))) - .andReturn(asList(task00, task01, task02)); - expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) - .andReturn(singletonList(task10)); - consumer.commitSync(expectedCommittedOffsets); - expectLastCall(); + expect(activeTaskCreator.createTasks(anyObject(), eq(assignmentActive))).andReturn(singleton(task00)); + expect(standbyTaskCreator.createTasks(eq(assignmentStandby))).andReturn(singletonList(task10)); replay(activeTaskCreator, standbyTaskCreator, consumer, changeLogReader); taskManager.handleAssignment(assignmentActive, assignmentStandby); assertThat(taskManager.tryToCompleteRestoration(), is(true)); assertThat(task00.state(), is(Task.State.RUNNING)); - assertThat(task01.state(), is(Task.State.RUNNING)); - assertThat(task02.state(), is(Task.State.RUNNING)); assertThat(task10.state(), is(Task.State.RUNNING)); - taskManager.handleRevocation(taskId00Partitions); + taskManager.handleAssignment(assignmentActive, Collections.emptyMap()); - assertThat(task01.commitPrepared, is(true)); - assertThat(task01.commitNeeded, is(false)); - assertThat(task02.commitPrepared, is(false)); - assertThat(task10.commitPrepared, is(false)); + assertThat(task00.commitNeeded, is(true)); } @Test @@ -1418,7 +1412,7 @@ public void shouldCloseActiveTasksDirtyAndPropagateSuspendException() { final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, true) { @Override public void suspend() { - transitionTo(State.SUSPENDED); + super.suspend(); throw new RuntimeException("task 0_1 suspend boom!"); } }; @@ -2561,7 +2555,7 @@ public void shouldFailOnCommitFatal() { } @Test - public void shouldNotSuspendOrCommitAllTasksIfSuspendingFailsDuringRevocation() { + public void shouldSuspendAllTasksButSkipCommitIfSuspendingFailsDuringRevocation() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true) { @Override public void suspend() { @@ -2583,9 +2577,9 @@ public void suspend() { RuntimeException.class, () -> taskManager.handleRevocation(asList(t1p0, t1p1))); - assertThat(thrown.getMessage(), is("KABOOM!")); + assertThat(thrown.getCause().getMessage(), is("KABOOM!")); assertThat(task00.state(), is(Task.State.SUSPENDED)); - assertThat(task01.state(), is(Task.State.CREATED)); + assertThat(task01.state(), is(Task.State.SUSPENDED)); } private static void expectRestoreToBeCompleted(final Consumer consumer, @@ -2708,6 +2702,8 @@ public void postCommit() { public void suspend() { if (state() == State.CLOSED) { throw new IllegalStateException("Illegal state " + state() + " while suspending active task " + id); + } else if (state() == State.SUSPENDED) { + // do nothing } else { transitionTo(State.SUSPENDED); } From bd032c97d515c9dfff9ed940ae795105b2c466c7 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 15 Jun 2020 19:31:03 -0700 Subject: [PATCH 11/14] log before transition --- .../apache/kafka/streams/processor/internals/StandbyTask.java | 2 +- .../apache/kafka/streams/processor/internals/StreamTask.java | 2 +- 2 files changed, 2 insertions(+), 2 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 b407ad1f4f2ed..5df59f6e37c3a 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 @@ -114,8 +114,8 @@ public void suspend() { switch (state()) { case CREATED: case RUNNING: - transitionTo(State.SUSPENDED); log.info("Suspended {}", state()); + transitionTo(State.SUSPENDED); break; 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 ef1cc928a9702..5c97c9760ec02 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 @@ -249,8 +249,8 @@ public void suspend() { switch (state()) { case CREATED: case RESTORING: - transitionTo(State.SUSPENDED); log.info("Suspended {}", state()); + transitionTo(State.SUSPENDED); break; From 0349663e2836129908a73e9e587ece9afb3bc04f Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 15 Jun 2020 19:57:21 -0700 Subject: [PATCH 12/14] attempt to postCommit all tasks --- .../processor/internals/TaskManager.java | 36 +++++++++++++------ 1 file changed, 26 insertions(+), 10 deletions(-) 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 538e82c308014..b518e04ef799c 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 @@ -462,9 +462,15 @@ void handleRevocation(final Collection revokedPartitions) { remainingRevokedPartitions.removeAll(task.inputPartitions()); } - final RuntimeException exception = firstException.get(); - if (exception != null) { - throw exception; + if (!remainingRevokedPartitions.isEmpty()) { + log.warn("The following partitions {} are missing from the task partitions. It could potentially " + + "due to race condition of consumer detecting the heartbeat failure, or the tasks " + + "have been cleaned up by the handleAssignment callback.", remainingRevokedPartitions); + } + + final RuntimeException suspendException = firstException.get(); + if (suspendException != null) { + throw suspendException; } // If using eos-beta, if we must commit any task then we must commit all of them @@ -488,13 +494,17 @@ void handleRevocation(final Collection revokedPartitions) { } for (final Task task : tasksToCommit) { - task.postCommit(); + try { + task.postCommit(); + } catch (final RuntimeException e) { + log.error("Exception caught while post-committing task " + task.id(), e); + firstException.compareAndSet(null, e); + } } - if (!remainingRevokedPartitions.isEmpty()) { - log.warn("The following partitions {} are missing from the task partitions. It could potentially " + - "due to race condition of consumer detecting the heartbeat failure, or the tasks " + - "have been cleaned up by the handleAssignment callback.", remainingRevokedPartitions); + final RuntimeException commitException = firstException.get(); + if (commitException != null) { + throw commitException; } } @@ -712,10 +722,16 @@ void shutdown(final boolean clean) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); } for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { - final Task task = tasks.get(taskId); - task.postCommit(); + try { + final Task task = tasks.get(taskId); + task.postCommit(); + } catch (final RuntimeException e) { + log.error("Exception caught while post-committing task " + taskId, e); + firstException.compareAndSet(null, e); + } } } catch (final RuntimeException e) { + log.error("Exception caught while committing tasks during shutdown", e); firstException.compareAndSet(null, e); } From 3ff79ea351cd20c4021b053133061984c0bbd202 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Tue, 16 Jun 2020 12:48:20 -0700 Subject: [PATCH 13/14] always commit when commitNeeded --- .../processor/internals/StreamTask.java | 25 +++-- .../processor/internals/TaskManager.java | 97 ++++++++----------- 2 files changed, 60 insertions(+), 62 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 5c97c9760ec02..fa8b94ba854a0 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 @@ -417,18 +417,30 @@ public void postCommit() { switch (state()) { case RESTORING: - case SUSPENDED: - maybeWriteCheckpoint(); + writeCheckpoint(); break; case RUNNING: if (!eosEnabled) { - maybeWriteCheckpoint(); + writeCheckpoint(); } break; + case SUSPENDED: + /* + * We must clear the `PartitionGroup` only after committing, and not in `suspend()`, + * because otherwise we lose the partition-time information. + * We also must clear it when the task is revoked, and not in `close()`, as the consumer will clear + * its internal buffer when the corresponding partition is revoked but the task may be reassigned + */ + partitionGroup.clear(); + + writeCheckpoint(); + + break; + case CREATED: case CLOSED: throw new IllegalStateException("Illegal state " + state() + " while post committing active task " + id); @@ -484,9 +496,6 @@ public void closeAndRecycleState() { throw new IllegalStateException("Unknown state " + state() + " while recycling active task " + id); } - // we cannot `clear()` the `PartitionGroup` in `suspend()` already, but only after committing, - // because otherwise we loose the partition-time information - partitionGroup.clear(); closeTaskSensor.record(); transitionTo(State.CLOSED); @@ -494,7 +503,7 @@ public void closeAndRecycleState() { log.info("Closed clean and recycled state"); } - private void maybeWriteCheckpoint() { + private void writeCheckpoint() { if (commitNeeded) { log.error("Tried to write a checkpoint with pending uncommitted data, should complete the commit first."); throw new IllegalStateException("A checkpoint should only be written if no commit is needed."); @@ -507,7 +516,7 @@ private void maybeWriteCheckpoint() { */ private void close(final boolean clean) { if (clean && commitNeeded) { - log.debug("Tried to close clean but there was an active scheduled checkpoint, this means we failed to" + log.debug("Tried to close clean but there was pending uncommitted data, this means we failed to" + " commit and should close as dirty instead"); throw new StreamsException("Tried to close dirty task as clean"); } 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 b518e04ef799c..48c18f09cb00b 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 @@ -482,16 +482,10 @@ void handleRevocation(final Collection revokedPartitions) { final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); for (final Task task : tasksToCommit) { final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); - } else { - log.warn("Task {} claimed to need a commit but had no committable consumed offsets", task.id()); - } + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); } - if (!consumedOffsetsAndMetadataPerTask.isEmpty()) { - commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - } + commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); for (final Task task : tasksToCommit) { try { @@ -700,9 +694,7 @@ void shutdown(final boolean clean) { task.suspend(); if (task.commitNeeded()) { final Map committableOffsets = task.prepareCommit(); - if (!committableOffsets.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); - } + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); } tasksToClose.add(task); } catch (final TaskMigratedException e) { @@ -718,16 +710,16 @@ void shutdown(final boolean clean) { } try { - if (clean && !consumedOffsetsAndMetadataPerTask.isEmpty()) { + if (clean) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - } - for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { - try { - final Task task = tasks.get(taskId); - task.postCommit(); - } catch (final RuntimeException e) { - log.error("Exception caught while post-committing task " + taskId, e); - firstException.compareAndSet(null, e); + for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { + try { + final Task task = tasks.get(taskId); + task.postCommit(); + } catch (final RuntimeException e) { + log.error("Exception caught while post-committing task " + taskId, e); + firstException.compareAndSet(null, e); + } } } } catch (final RuntimeException e) { @@ -851,30 +843,25 @@ void addRecordsToTasks(final ConsumerRecords records) { * or if the task producer got fenced (EOS) * @return number of committed offsets, or -1 if we are in the middle of a rebalance and cannot commit */ - int commit(final Collection tasks) { + int commit(final Collection tasksToCommit) { if (rebalanceInProgress) { return -1; } else { int committed = 0; final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); - for (final Task task : tasks) { + for (final Task task : tasksToCommit) { if (task.commitNeeded()) { final Map offsetAndMetadata = task.prepareCommit(); - if (!offsetAndMetadata.isEmpty()) { - consumedOffsetsAndMetadataPerTask.put(task.id(), offsetAndMetadata); - } + consumedOffsetsAndMetadataPerTask.put(task.id(), offsetAndMetadata); } } - if (!consumedOffsetsAndMetadataPerTask.isEmpty()) { - commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - } + commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - for (final Task task : tasks) { - if (task.commitNeeded()) { - ++committed; - task.postCommit(); - } + for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { + final Task task = tasks.get(taskId); + ++committed; + task.postCommit(); } return committed; @@ -899,28 +886,30 @@ int maybeCommitActiveTasksPerUserRequested() { } private void commitOffsetsOrTransaction(final Map> offsetsPerTask) { - if (processingMode == EXACTLY_ONCE_ALPHA) { - for (final Map.Entry> taskToCommit : offsetsPerTask.entrySet()) { - activeTaskCreator.streamsProducerForTask(taskToCommit.getKey()) - .commitTransaction(taskToCommit.getValue(), mainConsumer.groupMetadata()); - } - } else { - final Map allOffsets = offsetsPerTask.values().stream() - .flatMap(e -> e.entrySet().stream()).collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); - - if (processingMode == EXACTLY_ONCE_BETA) { - activeTaskCreator.threadProducer().commitTransaction(allOffsets, mainConsumer.groupMetadata()); + if (!offsetsPerTask.isEmpty()) { + if (processingMode == EXACTLY_ONCE_ALPHA) { + for (final Map.Entry> taskToCommit : offsetsPerTask.entrySet()) { + activeTaskCreator.streamsProducerForTask(taskToCommit.getKey()) + .commitTransaction(taskToCommit.getValue(), mainConsumer.groupMetadata()); + } } else { - try { - mainConsumer.commitSync(allOffsets); - } catch (final CommitFailedException error) { - throw new TaskMigratedException("Consumer committing offsets failed, " + - "indicating the corresponding thread is no longer part of the group", error); - } catch (final TimeoutException error) { - // TODO KIP-447: we can consider treating it as non-fatal and retry on the thread level - throw new StreamsException("Timed out while committing offsets via consumer", error); - } catch (final KafkaException error) { - throw new StreamsException("Error encountered committing offsets via consumer", error); + final Map allOffsets = offsetsPerTask.values().stream() + .flatMap(e -> e.entrySet().stream()).collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + + if (processingMode == EXACTLY_ONCE_BETA) { + activeTaskCreator.threadProducer().commitTransaction(allOffsets, mainConsumer.groupMetadata()); + } else { + try { + mainConsumer.commitSync(allOffsets); + } catch (final CommitFailedException error) { + throw new TaskMigratedException("Consumer committing offsets failed, " + + "indicating the corresponding thread is no longer part of the group", error); + } catch (final TimeoutException error) { + // TODO KIP-447: we can consider treating it as non-fatal and retry on the thread level + throw new StreamsException("Timed out while committing offsets via consumer", error); + } catch (final KafkaException error) { + throw new StreamsException("Error encountered committing offsets via consumer", error); + } } } } From 2529fb91b156a520ceb4178665ab313051f68cfd Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Tue, 16 Jun 2020 13:13:05 -0700 Subject: [PATCH 14/14] only commit active tasks --- .../processor/internals/TaskManager.java | 24 ++++++++++++------- .../processor/internals/TaskManagerTest.java | 4 +++- 2 files changed, 18 insertions(+), 10 deletions(-) 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 48c18f09cb00b..92885fd1c35d3 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 @@ -686,6 +686,7 @@ void shutdown(final boolean clean) { final AtomicReference firstException = new AtomicReference<>(null); final Set tasksToClose = new HashSet<>(); + final Set tasksToCommit = new HashSet<>(); final Map> consumedOffsetsAndMetadataPerTask = new HashMap<>(); for (final Task task : tasks.values()) { @@ -693,8 +694,11 @@ void shutdown(final boolean clean) { try { task.suspend(); if (task.commitNeeded()) { + tasksToCommit.add(task); final Map committableOffsets = task.prepareCommit(); - consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + if (task.isActive()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), committableOffsets); + } } tasksToClose.add(task); } catch (final TaskMigratedException e) { @@ -712,12 +716,11 @@ void shutdown(final boolean clean) { try { if (clean) { commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { + for (final Task task : tasksToCommit) { try { - final Task task = tasks.get(taskId); task.postCommit(); } catch (final RuntimeException e) { - log.error("Exception caught while post-committing task " + taskId, e); + log.error("Exception caught while post-committing task " + task.id(), e); firstException.compareAndSet(null, e); } } @@ -852,16 +855,19 @@ int commit(final Collection tasksToCommit) { for (final Task task : tasksToCommit) { if (task.commitNeeded()) { final Map offsetAndMetadata = task.prepareCommit(); - consumedOffsetsAndMetadataPerTask.put(task.id(), offsetAndMetadata); + if (task.isActive()) { + consumedOffsetsAndMetadataPerTask.put(task.id(), offsetAndMetadata); + } } } commitOffsetsOrTransaction(consumedOffsetsAndMetadataPerTask); - for (final TaskId taskId : consumedOffsetsAndMetadataPerTask.keySet()) { - final Task task = tasks.get(taskId); - ++committed; - task.postCommit(); + for (final Task task : tasksToCommit) { + if (task.commitNeeded()) { + ++committed; + task.postCommit(); + } } return committed; 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 389d187f93ad0..a0f3be552dba7 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 @@ -635,6 +635,7 @@ public void shouldCommitNonCorruptedTasksOnTaskCorruptedException() { expectLastCall().anyTimes(); expectRestoreToBeCompleted(consumer, changeLogReader); + consumer.commitSync(eq(emptyMap())); replay(activeTaskCreator, topologyBuilder, consumer, changeLogReader); @@ -1664,7 +1665,8 @@ public void shouldCommitProvidedTasksIfNeeded() { .andReturn(Arrays.asList(task00, task01, task02)).anyTimes(); expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) .andReturn(Arrays.asList(task03, task04, task05)).anyTimes(); - expectLastCall(); + + consumer.commitSync(eq(emptyMap())); replay(activeTaskCreator, standbyTaskCreator, consumer, changeLogReader);