From 19b8c702b360a189527646ebd8843aa080426536 Mon Sep 17 00:00:00 2001 From: Bruno Cadonna Date: Thu, 3 Aug 2023 18:02:40 +0200 Subject: [PATCH 1/2] KAFKA-10199: Change to RUNNING if no pending task to recycle exist A stream thread should only change to RUNNING if there are no active tasks in restoration in the state updater and if there are no pending tasks to recycle. There are situations in which a stream thread might only have standby tasks that are recycled to active task after a rebalance. In such situations, the stream thread might be faster in checking active tasks in restoration then the state updater removing the standby task to recycle from the state updater. If that happens the stream thread changes to RUNNING although it should wait until the standby tasks are recycled to active tasks and restored. --- .../processor/internals/TaskManager.java | 2 +- .../streams/processor/internals/Tasks.java | 5 ++++ .../processor/internals/TasksRegistry.java | 2 ++ .../processor/internals/TaskManagerTest.java | 30 +++++++++++++++++++ .../processor/internals/TasksTest.java | 25 ++++++++++++++++ 5 files changed, 63 insertions(+), 1 deletion(-) 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 9152d20721faa..203ef842b656d 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 @@ -763,7 +763,7 @@ public boolean checkStateUpdater(final long now, if (stateUpdater.restoresActiveTasks()) { handleRestoredTasksFromStateUpdater(now, offsetResetter); } - return !stateUpdater.restoresActiveTasks(); + return !stateUpdater.restoresActiveTasks() && !tasks.pendingTasksToRecycleExist(); } private void recycleTaskFromStateUpdater(final Task task, diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java index 7b3f7860fb72d..bcbfa7e42d4f3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java @@ -114,6 +114,11 @@ public void addPendingTaskToRecycle(final TaskId taskId, final Set action.getAction() == Action.RECYCLE); + } + @Override public Set removePendingTaskToUpdateInputPartitions(final TaskId taskId) { if (containsTaskIdWithAction(taskId, Action.UPDATE_INPUT_PARTITIONS)) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java index c93c4f145c6ab..8c2e854f6cfae 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java @@ -37,6 +37,8 @@ public interface TasksRegistry { Set removePendingTaskToRecycle(final TaskId taskId); + boolean pendingTasksToRecycleExist(); + void addPendingTaskToRecycle(final TaskId taskId, final Set inputPartitions); Set removePendingTaskToUpdateInputPartitions(final TaskId taskId); 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 3c626a8adba27..7969eecfeae2e 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 @@ -894,6 +894,7 @@ public void shouldSuspendRevokedTaskRemovedFromStateUpdater() { Mockito.verify(statefulTask).suspend(); Mockito.verify(tasks).addTask(statefulTask); } + @Test public void shouldHandleMultipleRemovedTasksFromStateUpdater() { final StreamTask taskToRecycle0 = statefulTask(taskId00, taskId00ChangelogPartitions) @@ -947,6 +948,35 @@ public void shouldHandleMultipleRemovedTasksFromStateUpdater() { Mockito.verify(stateUpdater).add(taskToUpdateInputPartitions); } + @Test + public void shouldReturnFalseFromCheckStateUpdaterIfActiveTasksAreRestoring() { + when(stateUpdater.restoresActiveTasks()).thenReturn(true); + final TasksRegistry tasks = mock(TasksRegistry.class); + final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + + assertFalse(taskManager.checkStateUpdater(time.milliseconds(), noOpResetter)); + } + + @Test + public void shouldReturnFalseFromCheckStateUpdaterIfActiveTasksAreNotRestoringButPendingTasksToRecycle() { + when(stateUpdater.restoresActiveTasks()).thenReturn(false); + final TasksRegistry tasks = mock(TasksRegistry.class); + when(tasks.pendingTasksToRecycleExist()).thenReturn(true); + final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + + assertFalse(taskManager.checkStateUpdater(time.milliseconds(), noOpResetter)); + } + + @Test + public void shouldReturnTrueFromCheckStateUpdaterIfActiveTasksAreNotRestoringAndNoPendingTasksToRecycle() { + when(stateUpdater.restoresActiveTasks()).thenReturn(false); + final TasksRegistry tasks = mock(TasksRegistry.class); + when(tasks.pendingTasksToRecycleExist()).thenReturn(false); + final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + + assertTrue(taskManager.checkStateUpdater(time.milliseconds(), noOpResetter)); + } + @Test public void shouldAddActiveTaskWithRevokedInputPartitionsInStateUpdaterToPendingTasksToSuspend() { final StreamTask task = statefulTask(taskId00, taskId00ChangelogPartitions) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java index d303eb4f60383..a35e5617c1997 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java @@ -45,7 +45,10 @@ public class TasksTest { private final static TopicPartition TOPIC_PARTITION_B_1 = new TopicPartition("topicB", 1); private final static TaskId TASK_0_0 = new TaskId(0, 0); private final static TaskId TASK_0_1 = new TaskId(0, 1); + private final static TaskId TASK_0_2 = new TaskId(0, 2); private final static TaskId TASK_1_0 = new TaskId(1, 0); + private final static TaskId TASK_1_1 = new TaskId(1, 1); + private final static TaskId TASK_1_2 = new TaskId(1, 2); private final Tasks tasks = new Tasks(new LogContext()); @@ -122,6 +125,28 @@ public void shouldAddAndRemovePendingTaskToRecycle() { assertNull(tasks.removePendingTaskToRecycle(TASK_0_0)); } + @Test + public void shouldVerifyIfPendingTaskToRecycleExist() { + assertFalse(tasks.pendingTasksToRecycleExist()); + tasks.addPendingTaskToRecycle(TASK_0_0, mkSet(TOPIC_PARTITION_A_0)); + assertTrue(tasks.pendingTasksToRecycleExist()); + + tasks.addPendingTaskToRecycle(TASK_1_0, mkSet(TOPIC_PARTITION_A_1)); + assertTrue(tasks.pendingTasksToRecycleExist()); + + tasks.addPendingTaskToCloseClean(TASK_0_1); + tasks.addPendingTaskToCloseDirty(TASK_0_2); + tasks.addPendingTaskToUpdateInputPartitions(TASK_1_1, mkSet(TOPIC_PARTITION_B_0)); + tasks.addPendingActiveTaskToSuspend(TASK_1_2); + assertTrue(tasks.pendingTasksToRecycleExist()); + + tasks.removePendingTaskToRecycle(TASK_0_0); + assertTrue(tasks.pendingTasksToRecycleExist()); + + tasks.removePendingTaskToRecycle(TASK_1_0); + assertFalse(tasks.pendingTasksToRecycleExist()); + } + @Test public void shouldAddAndRemovePendingTaskToUpdateInputPartitions() { final Set expectedInputPartitions = mkSet(TOPIC_PARTITION_A_0); From a92179ee9268c8f38c80b3ffd5c961b74e646f2f Mon Sep 17 00:00:00 2001 From: Bruno Cadonna Date: Thu, 3 Aug 2023 21:16:02 +0200 Subject: [PATCH 2/2] Rename pendingTasksToRecycleExist to hasPendingTaskToRecycle --- .../streams/processor/internals/TaskManager.java | 2 +- .../kafka/streams/processor/internals/Tasks.java | 2 +- .../streams/processor/internals/TasksRegistry.java | 2 +- .../streams/processor/internals/TaskManagerTest.java | 6 +++--- .../kafka/streams/processor/internals/TasksTest.java | 12 ++++++------ 5 files changed, 12 insertions(+), 12 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 203ef842b656d..13d01dd94ae18 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 @@ -763,7 +763,7 @@ public boolean checkStateUpdater(final long now, if (stateUpdater.restoresActiveTasks()) { handleRestoredTasksFromStateUpdater(now, offsetResetter); } - return !stateUpdater.restoresActiveTasks() && !tasks.pendingTasksToRecycleExist(); + return !stateUpdater.restoresActiveTasks() && !tasks.hasPendingTasksToRecycle(); } private void recycleTaskFromStateUpdater(final Task task, diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java index bcbfa7e42d4f3..d809bc8338e35 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Tasks.java @@ -115,7 +115,7 @@ public void addPendingTaskToRecycle(final TaskId taskId, final Set action.getAction() == Action.RECYCLE); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java index 8c2e854f6cfae..a04cd0fa6f8ec 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TasksRegistry.java @@ -37,7 +37,7 @@ public interface TasksRegistry { Set removePendingTaskToRecycle(final TaskId taskId); - boolean pendingTasksToRecycleExist(); + boolean hasPendingTasksToRecycle(); void addPendingTaskToRecycle(final TaskId taskId, final Set inputPartitions); 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 7969eecfeae2e..63d2a4b631f1b 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 @@ -885,7 +885,7 @@ public void shouldSuspendRevokedTaskRemovedFromStateUpdater() { when(tasks.removePendingActiveTaskToSuspend(statefulTask.id())).thenReturn(true); when(stateUpdater.hasRemovedTasks()).thenReturn(true); when(stateUpdater.drainRemovedTasks()).thenReturn(mkSet(statefulTask)); - final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); replay(consumer); taskManager.checkStateUpdater(time.milliseconds(), noOpResetter); @@ -961,7 +961,7 @@ public void shouldReturnFalseFromCheckStateUpdaterIfActiveTasksAreRestoring() { public void shouldReturnFalseFromCheckStateUpdaterIfActiveTasksAreNotRestoringButPendingTasksToRecycle() { when(stateUpdater.restoresActiveTasks()).thenReturn(false); final TasksRegistry tasks = mock(TasksRegistry.class); - when(tasks.pendingTasksToRecycleExist()).thenReturn(true); + when(tasks.hasPendingTasksToRecycle()).thenReturn(true); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); assertFalse(taskManager.checkStateUpdater(time.milliseconds(), noOpResetter)); @@ -971,7 +971,7 @@ public void shouldReturnFalseFromCheckStateUpdaterIfActiveTasksAreNotRestoringBu public void shouldReturnTrueFromCheckStateUpdaterIfActiveTasksAreNotRestoringAndNoPendingTasksToRecycle() { when(stateUpdater.restoresActiveTasks()).thenReturn(false); final TasksRegistry tasks = mock(TasksRegistry.class); - when(tasks.pendingTasksToRecycleExist()).thenReturn(false); + when(tasks.hasPendingTasksToRecycle()).thenReturn(false); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); assertTrue(taskManager.checkStateUpdater(time.milliseconds(), noOpResetter)); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java index a35e5617c1997..8ff3bea570bd1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java @@ -127,24 +127,24 @@ public void shouldAddAndRemovePendingTaskToRecycle() { @Test public void shouldVerifyIfPendingTaskToRecycleExist() { - assertFalse(tasks.pendingTasksToRecycleExist()); + assertFalse(tasks.hasPendingTasksToRecycle()); tasks.addPendingTaskToRecycle(TASK_0_0, mkSet(TOPIC_PARTITION_A_0)); - assertTrue(tasks.pendingTasksToRecycleExist()); + assertTrue(tasks.hasPendingTasksToRecycle()); tasks.addPendingTaskToRecycle(TASK_1_0, mkSet(TOPIC_PARTITION_A_1)); - assertTrue(tasks.pendingTasksToRecycleExist()); + assertTrue(tasks.hasPendingTasksToRecycle()); tasks.addPendingTaskToCloseClean(TASK_0_1); tasks.addPendingTaskToCloseDirty(TASK_0_2); tasks.addPendingTaskToUpdateInputPartitions(TASK_1_1, mkSet(TOPIC_PARTITION_B_0)); tasks.addPendingActiveTaskToSuspend(TASK_1_2); - assertTrue(tasks.pendingTasksToRecycleExist()); + assertTrue(tasks.hasPendingTasksToRecycle()); tasks.removePendingTaskToRecycle(TASK_0_0); - assertTrue(tasks.pendingTasksToRecycleExist()); + assertTrue(tasks.hasPendingTasksToRecycle()); tasks.removePendingTaskToRecycle(TASK_1_0); - assertFalse(tasks.pendingTasksToRecycleExist()); + assertFalse(tasks.hasPendingTasksToRecycle()); } @Test