From 35fe831ea9031d5033e191476f1d57bccccde0ef Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Tue, 4 Jul 2023 22:46:28 +0200 Subject: [PATCH 1/3] Kafka Streams Threading: Exception handling Catch any exceptions that escape the processing logic inside TaskExecutors and record them in the TaskManager. Make sure the TaskExecutor survives, but the task is unassigned. Add a method to TaskManager to drain the exceptions. The aim here is that the polling thread will drain the exceptions to be able to execute the uncaught exception handler, abort transactions, etc. --- .../internals/tasks/DefaultTaskExecutor.java | 14 ++++++- .../internals/tasks/DefaultTaskManager.java | 33 +++++++++++++++ .../internals/tasks/TaskManager.java | 23 +++++++++++ .../tasks/DefaultTaskExecutorTest.java | 16 ++++++++ .../tasks/DefaultTaskManagerTest.java | 40 +++++++++++++++++++ 5 files changed, 125 insertions(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java index c03384f4daafb..b59d695d9b0e0 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java @@ -54,7 +54,10 @@ public void run() { while (isRunning.get()) { runOnce(time.milliseconds()); } - // TODO: add exception handling + } catch (final StreamsException e) { + handleException(e); + } catch (final Exception e) { + handleException(new StreamsException(e)); } finally { if (currentTask != null) { unassignCurrentTask(); @@ -65,6 +68,15 @@ public void run() { } } + private void handleException(final StreamsException e) { + if (currentTask != null) { + taskManager.setUncaughtException(e, currentTask.id()); + } else { + // If we do not currently have a task assigned and still get an error, this is fatal for the executor thread + throw e; + } + } + private void runOnce(final long nowMs) { final KafkaFutureImpl pauseFuture; if ((pauseFuture = pauseRequested.getAndSet(null)) != null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java index 3f97de85cebcd..9b1feefef0d93 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.ReadOnlyTask; import org.apache.kafka.streams.processor.internals.StreamTask; @@ -55,6 +56,7 @@ public class DefaultTaskManager implements TaskManager { private final Lock tasksLock = new ReentrantLock(); private final List lockedTasks = new ArrayList<>(); + private final Map uncaughtExceptions = new HashMap<>(); private final Map assignedTasks = new HashMap<>(); private final List taskExecutors; @@ -225,6 +227,37 @@ public Set getTasks() { return returnWithTasksLocked(() -> tasks.activeTasks().stream().map(ReadOnlyTask::new).collect(Collectors.toSet())); } + @Override + public void setUncaughtException(final StreamsException exception, final TaskId taskId) { + executeWithTasksLocked(() -> { + + if (!assignedTasks.containsKey(taskId)) { + throw new IllegalArgumentException("An uncaught exception can only be set as long as the task is still assigned"); + } + + if (uncaughtExceptions.containsKey(taskId)) { + throw new IllegalArgumentException("The uncaught exception must be cleared before restarting processing"); + } + + uncaughtExceptions.put(taskId, exception); + }); + + log.info("Set an uncaught exception for task {}", taskId); + } + + public Map drainUncaughtExceptions() { + final Map returnValue = returnWithTasksLocked(() -> { + final Map result = new HashMap<>(uncaughtExceptions); + uncaughtExceptions.clear(); + return result; + }); + + log.info("Drained {} uncaught exceptions", returnValue.size()); + + return returnValue; + } + + private void executeWithTasksLocked(final Runnable action) { tasksLock.lock(); try { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/TaskManager.java index cfca9f75fe820..707719a00a60a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/TaskManager.java @@ -16,7 +16,9 @@ */ package org.apache.kafka.streams.processor.internals.tasks; +import java.util.Map; import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.ReadOnlyTask; import org.apache.kafka.streams.processor.internals.StreamTask; @@ -98,4 +100,25 @@ public interface TaskManager { * @return set of all managed active tasks */ Set getTasks(); + + /** + * Called whenever an existing task has thrown an uncaught exception. + * + * Setting an uncaught exception for a task prevents it from being reassigned until the + * corresponding exception has been handled in the polling thread. + * + */ + void setUncaughtException(StreamsException exception, TaskId taskId); + + /** + * Returns and clears all uncaught exceptions that were fell through to the processing + * threads and need to be handled in the polling thread. + * + * Called by the polling thread to handle processing exceptions, e.g. to abort + * transactions or shut down the application. + * + * @return A map from task ID to the exception that occurred. + */ + Map drainUncaughtExceptions(); + } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java index c2faeef880bae..ee8a0053c296a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Time; +import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.internals.StreamTask; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -104,4 +105,19 @@ public void shouldUnassignTaskWhenRequired() throws Exception { assertTrue(future.isDone(), "Unassign is not completed"); assertEquals(task, future.get(), "Unexpected task was unassigned"); } + + @Test + public void shouldSetUncaughtStreamsException() { + final StreamsException exception = mock(StreamsException.class); + when(task.process(anyLong())).thenThrow(exception); + + taskExecutor.start(); + + verify(taskManager, timeout(VERIFICATION_TIMEOUT)).assignNextTask(taskExecutor); + + verify(taskManager, timeout(VERIFICATION_TIMEOUT)).setUncaughtException(exception, task.id()); + verify(taskManager, timeout(VERIFICATION_TIMEOUT)).unassignTask(task, taskExecutor); + assertNull(taskExecutor.currentTask()); + } + } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java index e17a724f3652c..2a12f6445718c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.StreamTask; import org.apache.kafka.streams.processor.internals.TasksRegistry; @@ -50,6 +51,7 @@ public class DefaultTaskManagerTest { private final StreamTask task = mock(StreamTask.class); private final TasksRegistry tasks = mock(TasksRegistry.class); private final TaskExecutor taskExecutor = mock(TaskExecutor.class); + private final StreamsException exception = mock(StreamsException.class); private final StreamsConfig config = new StreamsConfig(configProps()); private final TaskManager taskManager = new DefaultTaskManager(time, "TaskManager", tasks, config, @@ -166,6 +168,44 @@ public void shouldNotAssignAnyLockedTask() { assertNull(taskManager.assignNextTask(taskExecutor)); } + @Test + public void shouldNotSetUncaughtExceptionsForUnassignedTasks() { + taskManager.add(Collections.singleton(task)); + + assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + } + + @Test + public void shouldNotSetUncaughtExceptionsTwice() { + taskManager.add(Collections.singleton(task)); + when(tasks.activeTasks()).thenReturn(Collections.singleton(task)); + taskManager.assignNextTask(taskExecutor); + taskManager.setUncaughtException(exception, task.id()); + + assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + } + + @Test + public void shouldReturnExceptionsOnDrainExceptions() { + taskManager.add(Collections.singleton(task)); + when(tasks.activeTasks()).thenReturn(Collections.singleton(task)); + taskManager.assignNextTask(taskExecutor); + taskManager.setUncaughtException(exception, task.id()); + + assertEquals(taskManager.drainUncaughtExceptions(), Collections.singletonMap(task.id(), exception)); + } + + @Test + public void shouldClearExceptionsOnDrainExceptions() { + taskManager.add(Collections.singleton(task)); + when(tasks.activeTasks()).thenReturn(Collections.singleton(task)); + taskManager.assignNextTask(taskExecutor); + taskManager.setUncaughtException(exception, task.id()); + taskManager.drainUncaughtExceptions(); + + assertEquals(taskManager.drainUncaughtExceptions(), Collections.emptyMap()); + } + @Test public void shouldUnassignLockingTask() { final KafkaFutureImpl future = new KafkaFutureImpl<>(); From f160b528a6f1a8a1437a9b1575089e693c4d16cf Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Mon, 10 Jul 2023 10:41:00 +0200 Subject: [PATCH 2/3] Update streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java Co-authored-by: Bruno Cadonna --- .../processor/internals/tasks/DefaultTaskExecutorTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java index ee8a0053c296a..db33a0d7c677e 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutorTest.java @@ -114,7 +114,6 @@ public void shouldSetUncaughtStreamsException() { taskExecutor.start(); verify(taskManager, timeout(VERIFICATION_TIMEOUT)).assignNextTask(taskExecutor); - verify(taskManager, timeout(VERIFICATION_TIMEOUT)).setUncaughtException(exception, task.id()); verify(taskManager, timeout(VERIFICATION_TIMEOUT)).unassignTask(task, taskExecutor); assertNull(taskExecutor.currentTask()); From c8536b50937594efcfd094e9815a6310d048f6c4 Mon Sep 17 00:00:00 2001 From: Lucas Brutschy Date: Mon, 10 Jul 2023 12:05:10 +0200 Subject: [PATCH 3/3] Bruno's comments --- .../internals/tasks/DefaultTaskManager.java | 5 ++++- .../tasks/DefaultTaskManagerTest.java | 18 +++++------------- 2 files changed, 9 insertions(+), 14 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java index 9b1feefef0d93..41526356d61c2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java @@ -242,7 +242,10 @@ public void setUncaughtException(final StreamsException exception, final TaskId uncaughtExceptions.put(taskId, exception); }); - log.info("Set an uncaught exception for task {}", taskId); + log.info("Set an uncaught exception of type {} for task {}, with error message: {}", + exception.getClass().getName(), + taskId, + exception.getMessage()); } public Map drainUncaughtExceptions() { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java index 2a12f6445718c..16e944d970401 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManagerTest.java @@ -172,7 +172,8 @@ public void shouldNotAssignAnyLockedTask() { public void shouldNotSetUncaughtExceptionsForUnassignedTasks() { taskManager.add(Collections.singleton(task)); - assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + final Exception e = assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + assertEquals("An uncaught exception can only be set as long as the task is still assigned", e.getMessage()); } @Test @@ -182,27 +183,18 @@ public void shouldNotSetUncaughtExceptionsTwice() { taskManager.assignNextTask(taskExecutor); taskManager.setUncaughtException(exception, task.id()); - assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + final Exception e = assertThrows(IllegalArgumentException.class, () -> taskManager.setUncaughtException(exception, task.id())); + assertEquals("The uncaught exception must be cleared before restarting processing", e.getMessage()); } @Test - public void shouldReturnExceptionsOnDrainExceptions() { + public void shouldReturnAndClearExceptionsOnDrainExceptions() { taskManager.add(Collections.singleton(task)); when(tasks.activeTasks()).thenReturn(Collections.singleton(task)); taskManager.assignNextTask(taskExecutor); taskManager.setUncaughtException(exception, task.id()); assertEquals(taskManager.drainUncaughtExceptions(), Collections.singletonMap(task.id(), exception)); - } - - @Test - public void shouldClearExceptionsOnDrainExceptions() { - taskManager.add(Collections.singleton(task)); - when(tasks.activeTasks()).thenReturn(Collections.singleton(task)); - taskManager.assignNextTask(taskExecutor); - taskManager.setUncaughtException(exception, task.id()); - taskManager.drainUncaughtExceptions(); - assertEquals(taskManager.drainUncaughtExceptions(), Collections.emptyMap()); }