From 1c38b0a560fe45b1a0154176747ea79e35974e6f Mon Sep 17 00:00:00 2001 From: Takla Gerguis Date: Wed, 21 Jun 2023 15:45:39 +0100 Subject: [PATCH 1/2] KAFKA-14133: Migrate StateDirectory mock in TaskManagerTest to Mockito --- .../processor/internals/TaskManagerTest.java | 60 ++++++------------- 1 file changed, 19 insertions(+), 41 deletions(-) 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 22da72feecdc6..3b39d68b33c32 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 @@ -185,7 +185,7 @@ public class TaskManagerTest { @org.mockito.Mock private InternalTopologyBuilder topologyBuilder; - @Mock(type = MockType.DEFAULT) + @org.mockito.Mock private StateDirectory stateDirectory; @org.mockito.Mock private ChangelogReader changeLogReader; @@ -1489,13 +1489,9 @@ public void shouldAddSubscribedTopicsFromAssignmentToTopologyMetadata() { @Test public void shouldNotLockAnythingIfStateDirIsEmpty() { - expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(new ArrayList<>()).once(); - - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); - - verify(stateDirectory); assertTrue(taskManager.lockedTaskDirectories().isEmpty()); + Mockito.verify(stateDirectory).listNonEmptyTaskDirectories(); } @Test @@ -1508,10 +1504,7 @@ public void shouldTryToLockValidTaskDirsAtRebalanceStart() throws Exception { taskId10.toString(), "dummy" ); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); - - verify(stateDirectory); assertThat(taskManager.lockedTaskDirectories(), is(singleton(taskId01))); } @@ -1555,7 +1548,6 @@ public void shouldReleaseLockForUnassignedTasksAfterRebalance() throws Exception taskId01.toString(), // standby task taskId02.toString() // unassigned but able to lock ); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); assertThat(taskManager.lockedTaskDirectories(), is(mkSet(taskId00, taskId01, taskId02))); @@ -1567,7 +1559,6 @@ public void shouldReleaseLockForUnassignedTasksAfterRebalance() throws Exception taskManager.handleRebalanceComplete(); assertThat(taskManager.lockedTaskDirectories(), is(mkSet(taskId00, taskId01))); - verify(stateDirectory); } @Test @@ -1677,7 +1668,6 @@ private void computeOffsetSumAndVerify(final Map changelog final Map expectedOffsetSums) throws Exception { expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); final StateMachineTask restoringTask = handleAssignment( @@ -1700,7 +1690,6 @@ public void shouldComputeOffsetSumForStandbyTask() throws Exception { expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); final StateMachineTask restoringTask = handleAssignment( @@ -1724,8 +1713,6 @@ public void shouldComputeOffsetSumForUnassignedTaskWeCanLock() throws Exception expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); writeCheckpointFile(taskId00, changelogOffsets); - - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); assertThat(taskManager.getTaskOffsetSums(), is(expectedOffsetSums)); @@ -1742,7 +1729,6 @@ public void shouldComputeOffsetSumFromCheckpointFileForUninitializedTask() throw expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); writeCheckpointFile(taskId00, changelogOffsets); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); final StateMachineTask uninitializedTask = new StateMachineTask(taskId00, taskId00Partitions, true); @@ -1766,7 +1752,6 @@ public void shouldComputeOffsetSumFromCheckpointFileForClosedTask() throws Excep expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); writeCheckpointFile(taskId00, changelogOffsets); - replay(stateDirectory); final StateMachineTask closedTask = new StateMachineTask(taskId00, taskId00Partitions, true); @@ -1787,7 +1772,6 @@ public void shouldComputeOffsetSumFromCheckpointFileForClosedTask() throws Excep public void shouldNotReportOffsetSumsForTaskWeCantLock() throws Exception { expectLockFailedFor(taskId00); makeTaskFolders(taskId00.toString()); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); assertTrue(taskManager.lockedTaskDirectories().isEmpty()); @@ -1798,12 +1782,10 @@ public void shouldNotReportOffsetSumsForTaskWeCantLock() throws Exception { public void shouldNotReportOffsetSumsAndReleaseLockForUnassignedTaskWithoutCheckpoint() throws Exception { expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); - expect(stateDirectory.checkpointFileFor(taskId00)).andReturn(getCheckpointFile(taskId00)); - replay(stateDirectory); + when(stateDirectory.checkpointFileFor(taskId00)).thenReturn(getCheckpointFile(taskId00)); taskManager.handleRebalanceStart(singleton("topic")); assertTrue(taskManager.getTaskOffsetSums().isEmpty()); - verify(stateDirectory); } @Test @@ -1819,7 +1801,6 @@ public void shouldPinOffsetSumToLongMaxValueInCaseOfOverflow() throws Exception expectLockObtainedFor(taskId00); makeTaskFolders(taskId00.toString()); writeCheckpointFile(taskId00, changelogOffsets); - replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); assertThat(taskManager.getTaskOffsetSums(), is(expectedOffsetSums)); @@ -1907,9 +1888,7 @@ public void shouldCloseActiveTasksWhenHandlingLostTasks() throws Exception { expectLockObtainedFor(taskId00, taskId01); // The second attempt will return empty tasks. - makeTaskFolders(); expectLockObtainedFor(); - replay(stateDirectory); taskManager.handleRebalanceStart(emptySet()); assertThat(taskManager.lockedTaskDirectories(), Matchers.is(mkSet(taskId00, taskId01))); @@ -1932,6 +1911,7 @@ public void shouldCloseActiveTasksWhenHandlingLostTasks() throws Exception { // The locked task map will not be cleared. assertThat(taskManager.lockedTaskDirectories(), is(mkSet(taskId00, taskId01))); + makeTaskFolders(); taskManager.handleRebalanceStart(emptySet()); assertThat(taskManager.lockedTaskDirectories(), is(emptySet())); @@ -2161,7 +2141,7 @@ public void shouldNotCommitNonRunningNonCorruptedTasks() { } @Test - public void shouldCleanAndReviveCorruptedStandbyTasksBeforeCommittingNonCorruptedTasks() { + public void shouldCleanAndReviveCorruptedStandbyTaskzBeforeCommittingNonCorruptedTaskz() { final ProcessorStateManager stateManager = EasyMock.createStrictMock(ProcessorStateManager.class); stateManager.markChangelogAsCorrupted(taskId00Partitions); replay(stateManager); @@ -2204,7 +2184,6 @@ public Map prepareCommit() { @Test public void shouldNotAttemptToCommitInHandleCorruptedDuringARebalance() { final ProcessorStateManager stateManager = EasyMock.createNiceMock(ProcessorStateManager.class); - expect(stateDirectory.listNonEmptyTaskDirectories()).andStubReturn(new ArrayList<>()); final StateMachineTask corruptedActive = new StateMachineTask(taskId00, taskId00Partitions, true, stateManager); @@ -2224,7 +2203,7 @@ public void shouldNotAttemptToCommitInHandleCorruptedDuringARebalance() { expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions)); - replay(consumer, stateDirectory, stateManager); + replay(consumer, stateManager); uncorruptedActive.setCommittableOffsetsAndMetadata(offsets); @@ -2848,7 +2827,7 @@ public void shouldCommitAllNeededTasksOnHandleRevocation() { } @Test - public void shouldNotCommitOnHandleAssignmentIfNoTaskClosed() { + public void shouldNotCommitOnHandleAssignmentIfNoTaskClozed() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); @@ -2878,7 +2857,7 @@ public void shouldNotCommitOnHandleAssignmentIfNoTaskClosed() { } @Test - public void shouldNotCommitOnHandleAssignmentIfOnlyStandbyTaskClosed() { + public void shouldNotCommitOnHandleAzzignmentIfOnlyStandbyTaskClozed() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final Map offsets00 = singletonMap(t1p0, new OffsetAndMetadata(0L, null)); task00.setCommittableOffsetsAndMetadata(offsets00); @@ -3428,8 +3407,8 @@ public void shouldHandleRebalanceEvents() { expect(consumer.assignment()).andReturn(assignment); consumer.pause(assignment); expectLastCall(); - expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(new ArrayList<>()); - replay(consumer, stateDirectory); + when(stateDirectory.listNonEmptyTaskDirectories()).thenReturn(new ArrayList<>()); + replay(consumer); assertThat(taskManager.rebalanceInProgress(), is(false)); taskManager.handleRebalanceStart(emptySet()); assertThat(taskManager.rebalanceInProgress(), is(true)); @@ -3544,14 +3523,13 @@ public void shouldNotCommitActiveAndStandbyTasksWhileRebalanceInProgress() throw final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, false); makeTaskFolders(taskId00.toString(), task01.toString()); - expectLockObtainedFor(taskId00, taskId01); expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); when(standbyTaskCreator.createTasks(taskId01Assignment)) .thenReturn(singletonList(task01)); - replay(stateDirectory, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4305,22 +4283,22 @@ private Map handleAssignment(final Map offsets) throws Exception { final File checkpointFile = getCheckpointFile(task); Files.createFile(checkpointFile.toPath()); new OffsetCheckpoint(checkpointFile).write(offsets); - expect(stateDirectory.checkpointFileFor(task)).andReturn(checkpointFile); + when(stateDirectory.checkpointFileFor(task)).thenReturn(checkpointFile); } private File getCheckpointFile(final TaskId task) { From 1ff753697b6166337ab50508ac6336f041b11573 Mon Sep 17 00:00:00 2001 From: Takla Gerguis Date: Thu, 20 Jul 2023 13:56:14 +0100 Subject: [PATCH 2/2] Address comments --- .../kafka/streams/processor/internals/TaskManagerTest.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) 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 3b39d68b33c32..b42f540ea7c52 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 @@ -1541,7 +1541,6 @@ public void shouldNotPauseReadyTasksWithStateUpdaterOnRebalanceComplete() { @Test public void shouldReleaseLockForUnassignedTasksAfterRebalance() throws Exception { expectLockObtainedFor(taskId00, taskId01, taskId02); - expectUnlockFor(taskId02); makeTaskFolders( taskId00.toString(), // active task @@ -1559,6 +1558,7 @@ public void shouldReleaseLockForUnassignedTasksAfterRebalance() throws Exception taskManager.handleRebalanceComplete(); assertThat(taskManager.lockedTaskDirectories(), is(mkSet(taskId00, taskId01))); + expectUnlockFor(taskId02); } @Test @@ -3407,7 +3407,6 @@ public void shouldHandleRebalanceEvents() { expect(consumer.assignment()).andReturn(assignment); consumer.pause(assignment); expectLastCall(); - when(stateDirectory.listNonEmptyTaskDirectories()).thenReturn(new ArrayList<>()); replay(consumer); assertThat(taskManager.rebalanceInProgress(), is(false)); taskManager.handleRebalanceStart(emptySet()); @@ -3522,7 +3521,8 @@ public void shouldNotCommitActiveAndStandbyTasksWhileRebalanceInProgress() throw final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, true); final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, false); - makeTaskFolders(taskId00.toString(), task01.toString()); + makeTaskFolders(taskId00.toString(), taskId01.toString()); + expectLockObtainedFor(taskId00, taskId01); expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); @@ -4297,7 +4297,6 @@ private void expectLockFailedFor(final TaskId... tasks) { private void expectUnlockFor(final TaskId... tasks) { for (final TaskId task : tasks) { - stateDirectory.unlock(task); Mockito.verify(stateDirectory, Mockito.atLeastOnce()).unlock(task); } }