From 50b03d2bfe84a7d4ca0aa71a3ad464e6eeb8ff5b Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Mon, 21 Nov 2022 15:40:50 +0100 Subject: [PATCH 1/9] Minor: enable Pause/Resume integration test This wasn't enabled and didn't pass for state updater. --- .../PauseResumeIntegrationTest.java | 100 ++++++++++-------- 1 file changed, 56 insertions(+), 44 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/PauseResumeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/PauseResumeIntegrationTest.java index e09ba997b2ec7..5e1331bf67525 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/PauseResumeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/PauseResumeIntegrationTest.java @@ -35,17 +35,16 @@ import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyBuilder; import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.Stores; -import org.apache.kafka.test.IntegrationTest; import org.apache.kafka.test.TestUtils; import org.hamcrest.CoreMatchers; -import org.junit.After; -import org.junit.AfterClass; -import org.junit.Before; -import org.junit.BeforeClass; -import org.junit.Rule; -import org.junit.Test; -import org.junit.experimental.categories.Category; -import org.junit.rules.TestName; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.TestInfo; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.MethodSource; import java.time.Duration; import java.util.ArrayList; @@ -53,6 +52,7 @@ import java.util.Collection; import java.util.List; import java.util.Properties; +import java.util.stream.Stream; import static java.util.Arrays.asList; import static java.util.Collections.singletonList; @@ -69,7 +69,7 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; -@Category({IntegrationTest.class}) +@Tag("integration") public class PauseResumeIntegrationTest { private static final Duration STARTUP_TIMEOUT = Duration.ofSeconds(45); public static final EmbeddedKafkaCluster CLUSTER = new EmbeddedKafkaCluster(1); @@ -101,10 +101,14 @@ public class PauseResumeIntegrationTest { private KafkaStreams kafkaStreams, kafkaStreams2; private KafkaStreamsNamedTopologyWrapper streamsNamedTopologyWrapper; - @Rule - public final TestName testName = new TestName(); - @BeforeClass + private static Stream parameters() { + return Stream.of( + Boolean.TRUE, + Boolean.FALSE); + } + + @BeforeAll public static void startCluster() throws Exception { CLUSTER.start(); producerConfig = TestUtils.producerConfig(CLUSTER.bootstrapServers(), @@ -113,18 +117,18 @@ public static void startCluster() throws Exception { StringDeserializer.class, LongDeserializer.class); } - @AfterClass + @AfterAll public static void closeCluster() { CLUSTER.stop(); } - @Before - public void createTopics() throws InterruptedException { + @BeforeEach + public void createTopics(final TestInfo testInfo) throws InterruptedException { cleanStateBeforeTest(CLUSTER, 1, INPUT_STREAM_1, INPUT_STREAM_2, OUTPUT_STREAM_1, OUTPUT_STREAM_2); - appId = safeUniqueTestName(PauseResumeIntegrationTest.class, testName); + appId = safeUniqueTestName(PauseResumeIntegrationTest.class, testInfo); } - private Properties props() { + private Properties props(final boolean stateUpdaterEnabled) { final Properties properties = new Properties(); properties.put(StreamsConfig.APPLICATION_ID_CONFIG, appId); properties.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, CLUSTER.bootstrapServers()); @@ -136,10 +140,11 @@ private Properties props() { properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); properties.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 100); properties.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 1000); + properties.put(StreamsConfig.InternalConfig.STATE_UPDATER_ENABLED, stateUpdaterEnabled); return properties; } - @After + @AfterEach public void shutdown() throws InterruptedException { for (final KafkaStreams streams : Arrays.asList(kafkaStreams, kafkaStreams2, streamsNamedTopologyWrapper)) { if (streams != null) { @@ -152,9 +157,10 @@ private static void produceToInputTopics(final String topic, final Collection Date: Tue, 20 Dec 2022 17:25:52 +0100 Subject: [PATCH 2/9] KAFKA-14299: Make sure no progress is made on paused topologies The state updater restored one round of polls from the restore consumer before realizing that a newly added task was already in paused state when being added. This is tested by the existing `PauseResumeIntegrationTest`. --- .../streams/processor/internals/DefaultStateUpdater.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java index d96593c5011d7..5c1cfac821493 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java @@ -277,6 +277,11 @@ private void addTask(final Task task) { changelogReader.transitToUpdateStandby(); } } + + // Move task to paused tasks immediately to not make any progress + if (topologyMetadata.isPaused(task.id().topologyName())) { + pauseTask(task); + } } } From 1c56e836f3b365c6265715ada5dec65ad69861b3 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Tue, 3 Jan 2023 15:03:38 +0100 Subject: [PATCH 3/9] Bruno's comments --- .../internals/DefaultStateUpdater.java | 21 ++++++++++--------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java index 5c1cfac821493..9fd02cc4daf2c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java @@ -258,30 +258,31 @@ private List getTasksAndActions() { } private void addTask(final Task task) { + final TaskId taskId = task.id(); if (isStateless(task)) { addToRestoredTasks((StreamTask) task); - log.info("Stateless active task " + task.id() + " was added to the restored tasks of the state updater"); + log.info("Stateless active task " + taskId + " was added to the restored tasks of the state updater"); + } else if (topologyMetadata.isPaused(taskId.topologyName())) { + pausedTasks.put(taskId, task); + changelogReader.register(task.changelogPartitions(), task.stateManager()); + log.debug((task.isActive() ? "Active" : "Standby") + + " task " + taskId + " was directly added to the paused tasks."); } else { - final Task existingTask = updatingTasks.putIfAbsent(task.id(), task); + final Task existingTask = updatingTasks.putIfAbsent(taskId, task); if (existingTask != null) { - throw new IllegalStateException((existingTask.isActive() ? "Active" : "Standby") + " task " + task.id() + " already exist, " + + throw new IllegalStateException((existingTask.isActive() ? "Active" : "Standby") + " task " + taskId + " already exist, " + "should not try to add another " + (task.isActive() ? "active" : "standby") + " task with the same id. " + BUG_ERROR_MESSAGE); } changelogReader.register(task.changelogPartitions(), task.stateManager()); if (task.isActive()) { - log.info("Stateful active task " + task.id() + " was added to the state updater"); + log.info("Stateful active task " + taskId + " was added to the state updater"); changelogReader.enforceRestoreActive(); } else { - log.info("Standby task " + task.id() + " was added to the state updater"); + log.info("Standby task " + taskId + " was added to the state updater"); if (updatingTasks.size() == 1) { changelogReader.transitToUpdateStandby(); } } - - // Move task to paused tasks immediately to not make any progress - if (topologyMetadata.isPaused(task.id().topologyName())) { - pauseTask(task); - } } } From 8d306e2fca2b3a530f16120b61df3d5a3a44e80e Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Tue, 3 Jan 2023 15:05:11 +0100 Subject: [PATCH 4/9] All tasks in TaskManager should return tasks in restoration --- .../kafka/streams/processor/internals/ReadOnlyTask.java | 4 ++-- .../kafka/streams/processor/internals/TaskManager.java | 8 +++++++- .../streams/processor/internals/ReadOnlyTaskTest.java | 2 ++ 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java index e26b6ca29fd54..1275e4201112d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java @@ -190,7 +190,7 @@ public void clearTaskTimeout() { @Override public boolean commitNeeded() { - throw new UnsupportedOperationException("This task is read-only"); + return task.commitNeeded(); } @Override @@ -200,7 +200,7 @@ public StateStore getStore(final String name) { @Override public Map changelogOffsets() { - throw new UnsupportedOperationException("This task is read-only"); + return task.changelogOffsets(); } @Override 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 583662f092476..b71aca3fdd5b6 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 @@ -1534,7 +1534,13 @@ Set standbyTaskIds() { Map allTasks() { // not bothering with an unmodifiable map, since the tasks themselves are mutable, but // if any outside code modifies the map or the tasks, it would be a severe transgression. - return tasks.allTasksPerId(); + if (stateUpdater != null) { + final Map ret = stateUpdater.getTasks().stream().collect(Collectors.toMap(Task::id, x -> x)); + ret.putAll(tasks.allTasksPerId()); + return ret; + } else { + return tasks.allTasksPerId(); + } } Map notPausedTasks() { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java index cd5da8739818f..c9e98f248d93c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java @@ -41,6 +41,8 @@ class ReadOnlyTaskTest { add("changelogPartitions"); add("commitRequested"); add("isActive"); + add("commitNeeded"); + add("changelogOffsets"); add("state"); add("id"); } From 827f7eca99d1843455dc4f90379cae601d624417 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Wed, 4 Jan 2023 16:49:01 +0100 Subject: [PATCH 5/9] Wake up state updater when tasks are being resumed --- .../apache/kafka/streams/KafkaStreams.java | 1 + .../internals/DefaultStateUpdater.java | 31 ++++++++++++++----- .../processor/internals/StateUpdater.java | 5 +++ .../processor/internals/StreamThread.java | 5 +++ .../processor/internals/TaskManager.java | 6 ++++ .../KafkaStreamsNamedTopologyWrapper.java | 2 ++ .../internals/DefaultStateUpdaterTest.java | 1 + 7 files changed, 44 insertions(+), 7 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index bc3a1243ce2e1..df428a8503d31 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -1747,6 +1747,7 @@ public void resume() { } else { topologyMetadata.resumeTopology(UNNAMED_TOPOLOGY); } + threads.forEach(StreamThread::signalResume); } /** diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java index 9fd02cc4daf2c..63960a77574e6 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java @@ -116,6 +116,7 @@ public void run() { private void runOnce() throws InterruptedException { performActionsOnTasks(); + resumeTasks(); restoreTasks(); checkAllUpdatingTaskStates(time.milliseconds()); waitIfAllChangelogsCompletelyRead(); @@ -140,6 +141,16 @@ private void performActionsOnTasks() { } } + private void resumeTasks() { + if (needsResumeCheck.compareAndSet(true, false)) { + for (final Task task : pausedTasks.values()) { + if (!topologyMetadata.isPaused(task.id().topologyName())) { + resumeTask(task); + } + } + } + } + private void restoreTasks() { try { changelogReader.restore(updatingTasks); @@ -229,7 +240,7 @@ private void waitIfAllChangelogsCompletelyRead() throws InterruptedException { if (isRunning.get() && changelogReader.allChangelogsCompleted()) { tasksAndActionsLock.lock(); try { - while (tasksAndActions.isEmpty()) { + while (tasksAndActions.isEmpty() && !needsResumeCheck.get()) { tasksAndActionsCondition.await(); } } finally { @@ -394,12 +405,6 @@ private void checkAllUpdatingTaskStates(final long now) { } } - for (final Task task : pausedTasks.values()) { - if (!topologyMetadata.isPaused(task.id().topologyName())) { - resumeTask(task); - } - } - lastCommitMs = now; } } @@ -417,6 +422,7 @@ private void checkAllUpdatingTaskStates(final long now) { private final Condition restoredActiveTasksCondition = restoredActiveTasksLock.newCondition(); private final BlockingQueue exceptionsAndFailedTasks = new LinkedBlockingQueue<>(); private final BlockingQueue removedTasks = new LinkedBlockingQueue<>(); + private final AtomicBoolean needsResumeCheck = new AtomicBoolean(false); private final long commitIntervalMs; private long lastCommitMs; @@ -512,6 +518,17 @@ public void remove(final TaskId taskId) { } } + @Override + public void signalResume() { + tasksAndActionsLock.lock(); + try { + needsResumeCheck.set(true); + tasksAndActionsCondition.signalAll(); + } finally { + tasksAndActionsLock.unlock(); + } + } + @Override public Set drainRestoredActiveTasks(final Duration timeout) { final long timeoutMs = timeout.toMillis(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateUpdater.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateUpdater.java index 10ac51874d264..3153472bf929f 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateUpdater.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateUpdater.java @@ -94,6 +94,11 @@ public int hashCode() { */ void remove(final TaskId taskId); + /** + * Wakes up the state updater if it is currently dormant, to check if a paused task should be resumed. + */ + void signalResume(); + /** * Drains the restored active tasks from the state updater. * diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 1ea75d17dc09f..1d4a84de79f53 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -1090,6 +1090,11 @@ public boolean isThreadAlive() { return isAlive(); } + // Call method when a topology is resumed + public void signalResume() { + taskManager.signalResume(); + } + /** * Try to commit all active tasks owned by this thread. * 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 b71aca3fdd5b6..cb70a35364bee 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 @@ -1115,6 +1115,12 @@ private void removeLostActiveTasksFromStateUpdater() { } } + public void signalResume() { + if (stateUpdater != null) { + stateUpdater.signalResume(); + } + } + /** * Compute the offset total summed across all stores in a task. Includes offset sum for any tasks we own the * lock for, which includes assigned and unassigned tasks we locked in {@link #tryToLockAllNonEmptyTaskDirectories()}. diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java index 3d22c583373c2..10bcff08d93d3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java @@ -40,6 +40,7 @@ import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.DefaultKafkaClientSupplier; import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder; +import org.apache.kafka.streams.processor.internals.StreamThread; import org.apache.kafka.streams.processor.internals.Task; import org.apache.kafka.streams.processor.internals.TopologyMetadata; import org.slf4j.Logger; @@ -276,6 +277,7 @@ public boolean isNamedTopologyPaused(final String topologyName) { */ public void resumeNamedTopology(final String topologyName) { topologyMetadata.resumeTopology(topologyName); + threads.forEach(StreamThread::signalResume); } /** diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdaterTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdaterTest.java index 9f1e642a17e75..26e8754e286b0 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdaterTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdaterTest.java @@ -1047,6 +1047,7 @@ private void shouldResumeStatefulTask(final Task task) throws Exception { verifyUpdatingTasks(); when(topologyMetadata.isPaused(null)).thenReturn(false); + stateUpdater.signalResume(); verifyPausedTasks(); verifyUpdatingTasks(task); From 124fedf92287be89a65d209274e395062ef676df Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Mon, 6 Feb 2023 13:10:15 +0100 Subject: [PATCH 6/9] Reviewer comments addressed --- .../internals/DefaultStateUpdater.java | 30 +++++++++++++------ .../processor/internals/TaskManagerTest.java | 16 ++++++++++ 2 files changed, 37 insertions(+), 9 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java index 63960a77574e6..ae6618c304f70 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java @@ -142,7 +142,7 @@ private void performActionsOnTasks() { } private void resumeTasks() { - if (needsResumeCheck.compareAndSet(true, false)) { + if (isTopologyResumed.compareAndSet(true, false)) { for (final Task task : pausedTasks.values()) { if (!topologyMetadata.isPaused(task.id().topologyName())) { resumeTask(task); @@ -240,7 +240,7 @@ private void waitIfAllChangelogsCompletelyRead() throws InterruptedException { if (isRunning.get() && changelogReader.allChangelogsCompleted()) { tasksAndActionsLock.lock(); try { - while (tasksAndActions.isEmpty() && !needsResumeCheck.get()) { + while (tasksAndActions.isEmpty() && !isTopologyResumed.get()) { tasksAndActionsCondition.await(); } } finally { @@ -270,6 +270,22 @@ private List getTasksAndActions() { private void addTask(final Task task) { final TaskId taskId = task.id(); + + Task existingTask = pausedTasks.get(taskId); + if (existingTask != null) { + throw new IllegalStateException( + (existingTask.isActive() ? "Active" : "Standby") + " task " + taskId + " already exist in paused tasks, " + + "should not try to add another " + (task.isActive() ? "active" : "standby") + " task with the same id. " + + BUG_ERROR_MESSAGE); + } + existingTask = updatingTasks.get(taskId); + if (existingTask != null) { + throw new IllegalStateException( + (existingTask.isActive() ? "Active" : "Standby") + " task " + taskId + " already exist in updating tasks, " + + "should not try to add another " + (task.isActive() ? "active" : "standby") + " task with the same id. " + + BUG_ERROR_MESSAGE); + } + if (isStateless(task)) { addToRestoredTasks((StreamTask) task); log.info("Stateless active task " + taskId + " was added to the restored tasks of the state updater"); @@ -279,11 +295,7 @@ private void addTask(final Task task) { log.debug((task.isActive() ? "Active" : "Standby") + " task " + taskId + " was directly added to the paused tasks."); } else { - final Task existingTask = updatingTasks.putIfAbsent(taskId, task); - if (existingTask != null) { - throw new IllegalStateException((existingTask.isActive() ? "Active" : "Standby") + " task " + taskId + " already exist, " + - "should not try to add another " + (task.isActive() ? "active" : "standby") + " task with the same id. " + BUG_ERROR_MESSAGE); - } + updatingTasks.put(taskId, task); changelogReader.register(task.changelogPartitions(), task.stateManager()); if (task.isActive()) { log.info("Stateful active task " + taskId + " was added to the state updater"); @@ -422,7 +434,7 @@ private void checkAllUpdatingTaskStates(final long now) { private final Condition restoredActiveTasksCondition = restoredActiveTasksLock.newCondition(); private final BlockingQueue exceptionsAndFailedTasks = new LinkedBlockingQueue<>(); private final BlockingQueue removedTasks = new LinkedBlockingQueue<>(); - private final AtomicBoolean needsResumeCheck = new AtomicBoolean(false); + private final AtomicBoolean isTopologyResumed = new AtomicBoolean(false); private final long commitIntervalMs; private long lastCommitMs; @@ -522,7 +534,7 @@ public void remove(final TaskId taskId) { public void signalResume() { tasksAndActionsLock.lock(); try { - needsResumeCheck.set(true); + isTopologyResumed.set(true); tasksAndActionsCondition.signalAll(); } finally { tasksAndActionsLock.unlock(); 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 00d44b045b13d..2841fae8c7a07 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 @@ -497,6 +497,22 @@ public void shouldAssignMultipleTasksInStateUpdater() { Mockito.verify(tasks).addPendingTaskToRecycle(standbyTaskToRecycle.id(), standbyTaskToRecycle.inputPartitions()); } + @Test + public void shouldReturnStateUpdaterTasksInAllTasks() { + final StreamTask activeTask = statefulTask(taskId03, taskId03ChangelogPartitions) + .inState(State.RUNNING) + .withInputPartitions(taskId03Partitions).build(); + final StandbyTask standbyTask = standbyTask(taskId02, taskId02ChangelogPartitions) + .inState(State.RUNNING) + .withInputPartitions(taskId02Partitions).build(); + final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); + final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + + when(stateUpdater.getTasks()).thenReturn(mkSet(standbyTask)); + when(tasks.allTasksPerId()).thenReturn(mkMap(mkEntry(taskId03, activeTask))); + assertEquals(taskManager.allTasks(), mkMap(mkEntry(taskId03, activeTask), mkEntry(taskId02, standbyTask))); + } + @Test public void shouldCreateActiveTaskDuringAssignment() { final StreamTask activeTaskToBeCreated = statefulTask(taskId03, taskId03ChangelogPartitions) From 9d83fe792e8fc9f0ba3b87d79e48e00e7f73eef7 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Thu, 16 Feb 2023 13:54:12 +0100 Subject: [PATCH 7/9] Fix: only process owned tasks in `maybeCommit` --- .../processor/internals/ReadOnlyTask.java | 2 +- .../processor/internals/StreamThread.java | 2 +- .../streams/processor/internals/TaskManager.java | 11 +++++++++++ .../processor/internals/ReadOnlyTaskTest.java | 1 - .../processor/internals/TaskManagerTest.java | 16 ++++++++++++++++ 5 files changed, 29 insertions(+), 3 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java index 1275e4201112d..ee3989cf62e48 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ReadOnlyTask.java @@ -190,7 +190,7 @@ public void clearTaskTimeout() { @Override public boolean commitNeeded() { - return task.commitNeeded(); + throw new UnsupportedOperationException("This task is read-only"); } @Override diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index dbb0c5d750a39..aad2c79b383ae 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -1115,7 +1115,7 @@ int maybeCommit() { } committed = taskManager.commit( - taskManager.allTasks() + taskManager.allOwnedTasks() .values() .stream() .filter(t -> t.state() == Task.State.RUNNING || t.state() == Task.State.RESTORING) 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 cb70a35364bee..732d0fb52516f 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 @@ -1549,6 +1549,17 @@ Map allTasks() { } } + /** + * Returns tasks owned by the stream thread. With state updater disabled, these are all tasks. With + * state updater enabled, this does not return any tasks currently owned by the state updater. + * @return + */ + Map allOwnedTasks() { + // not bothering with an unmodifiable map, since the tasks themselves are mutable, but + // if any outside code modifies the map or the tasks, it would be a severe transgression. + return tasks.allTasksPerId(); + } + Map notPausedTasks() { return Collections.unmodifiableMap(tasks.allTasks() .stream() diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java index c9e98f248d93c..9b780dee6bbc3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ReadOnlyTaskTest.java @@ -41,7 +41,6 @@ class ReadOnlyTaskTest { add("changelogPartitions"); add("commitRequested"); add("isActive"); - add("commitNeeded"); add("changelogOffsets"); add("state"); add("id"); 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 2841fae8c7a07..f43de372388d0 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 @@ -513,6 +513,22 @@ public void shouldReturnStateUpdaterTasksInAllTasks() { assertEquals(taskManager.allTasks(), mkMap(mkEntry(taskId03, activeTask), mkEntry(taskId02, standbyTask))); } + @Test + public void shouldNotReturnStateUpdaterTasksInOwnedTasks() { + final StreamTask activeTask = statefulTask(taskId03, taskId03ChangelogPartitions) + .inState(State.RUNNING) + .withInputPartitions(taskId03Partitions).build(); + final StandbyTask standbyTask = standbyTask(taskId02, taskId02ChangelogPartitions) + .inState(State.RUNNING) + .withInputPartitions(taskId02Partitions).build(); + final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); + final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); + + when(stateUpdater.getTasks()).thenReturn(mkSet(standbyTask)); + when(tasks.allTasksPerId()).thenReturn(mkMap(mkEntry(taskId03, activeTask))); + assertEquals(taskManager.allOwnedTasks(), mkMap(mkEntry(taskId03, activeTask))); + } + @Test public void shouldCreateActiveTaskDuringAssignment() { final StreamTask activeTaskToBeCreated = statefulTask(taskId03, taskId03ChangelogPartitions) From 7673087dfc1eddf9df66361045b7781a2187d094 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Thu, 16 Feb 2023 14:03:14 +0100 Subject: [PATCH 8/9] Fix tests --- .../streams/processor/internals/StreamThreadTest.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index dc9259a1f91ac..7114fc13d8818 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -2734,7 +2734,7 @@ public void shouldNotCommitNonRunningNonRestoringTasks() { expect(task3.state()).andReturn(Task.State.CREATED).anyTimes(); expect(task3.id()).andReturn(taskId3).anyTimes(); - expect(taskManager.allTasks()).andReturn(mkMap( + expect(taskManager.allOwnedTasks()).andReturn(mkMap( mkEntry(taskId1, task1), mkEntry(taskId2, task2), mkEntry(taskId3, task3) @@ -3084,7 +3084,7 @@ private TaskManager mockTaskManager(final Task runningTask) { expect(runningTask.state()).andStubReturn(Task.State.RUNNING); expect(runningTask.id()).andStubReturn(taskId); - expect(taskManager.allTasks()).andStubReturn(Collections.singletonMap(taskId, runningTask)); + expect(taskManager.allOwnedTasks()).andStubReturn(Collections.singletonMap(taskId, runningTask)); expect(taskManager.commit(Collections.singleton(runningTask))).andStubReturn(1); return taskManager; } @@ -3159,7 +3159,7 @@ private void addRecord(final MockConsumer mockConsumer, } StreamTask activeTask(final TaskManager taskManager, final TopicPartition partition) { - final Stream standbys = taskManager.allTasks().values().stream().filter(Task::isActive); + final Stream standbys = taskManager.allOwnedTasks().values().stream().filter(Task::isActive); for (final Task task : (Iterable) standbys::iterator) { if (task.inputPartitions().contains(partition)) { return (StreamTask) task; @@ -3168,7 +3168,7 @@ StreamTask activeTask(final TaskManager taskManager, final TopicPartition partit return null; } StandbyTask standbyTask(final TaskManager taskManager, final TopicPartition partition) { - final Stream standbys = taskManager.allTasks().values().stream().filter(t -> !t.isActive()); + final Stream standbys = taskManager.allOwnedTasks().values().stream().filter(t -> !t.isActive()); for (final Task task : (Iterable) standbys::iterator) { if (task.inputPartitions().contains(partition)) { return (StandbyTask) task; From 019de467523f5b0bd28271374201e8d4112b7216 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Mon, 20 Feb 2023 13:46:45 +0100 Subject: [PATCH 9/9] Fix: add state updater tasks to `readOnlyAllTasks` as well. --- .../streams/processor/internals/TaskManager.java | 12 +++++++++--- .../processor/internals/StreamThreadTest.java | 4 ++-- 2 files changed, 11 insertions(+), 5 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 47698ff405d3b..cf83cb27ebaff 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 @@ -1561,9 +1561,15 @@ Map allOwnedTasks() { } Set readOnlyAllTasks() { - // need to make sure the returned set is unmodifiable as it could be accessed - // by other thread than the StreamThread owning this task manager; - return Collections.unmodifiableSet(tasks.allTasks()); + // not bothering with an unmodifiable map, since the tasks themselves are mutable, but + // if any outside code modifies the map or the tasks, it would be a severe transgression. + if (stateUpdater != null) { + final HashSet ret = new HashSet<>(stateUpdater.getTasks()); + ret.addAll(tasks.allTasks()); + return Collections.unmodifiableSet(ret); + } else { + return Collections.unmodifiableSet(tasks.allTasks()); + } } Map notPausedTasks() { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index 7114fc13d8818..a8768ddd3eb0d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -3159,7 +3159,7 @@ private void addRecord(final MockConsumer mockConsumer, } StreamTask activeTask(final TaskManager taskManager, final TopicPartition partition) { - final Stream standbys = taskManager.allOwnedTasks().values().stream().filter(Task::isActive); + final Stream standbys = taskManager.allTasks().values().stream().filter(Task::isActive); for (final Task task : (Iterable) standbys::iterator) { if (task.inputPartitions().contains(partition)) { return (StreamTask) task; @@ -3168,7 +3168,7 @@ StreamTask activeTask(final TaskManager taskManager, final TopicPartition partit return null; } StandbyTask standbyTask(final TaskManager taskManager, final TopicPartition partition) { - final Stream standbys = taskManager.allOwnedTasks().values().stream().filter(t -> !t.isActive()); + final Stream standbys = taskManager.allTasks().values().stream().filter(t -> !t.isActive()); for (final Task task : (Iterable) standbys::iterator) { if (task.inputPartitions().contains(partition)) { return (StandbyTask) task;