From 3d90058d967518dd3915a13be8a5f80df3def361 Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Tue, 22 Feb 2022 02:34:31 -0800 Subject: [PATCH 1/8] introduce TaskExecutionMetadata --- .../internals/TaskExecutionMetadata.java | 144 ++++++++++++++++++ .../processor/internals/TaskExecutor.java | 32 ++-- .../processor/internals/TaskManager.java | 8 +- .../streams/processor/internals/Tasks.java | 11 ++ .../processor/internals/TopologyMetadata.java | 15 +- .../processor/internals/TasksTest.java | 106 +++++++++++++ 6 files changed, 303 insertions(+), 13 deletions(-) create mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java create mode 100644 streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java new file mode 100644 index 0000000000000..b3a3749a28439 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java @@ -0,0 +1,144 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.processor.internals; + +import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.streams.errors.UnknownTopologyException; +import org.apache.kafka.streams.processor.TaskId; + +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentSkipListSet; +import org.slf4j.Logger; + +import static org.apache.kafka.streams.processor.internals.TopologyMetadata.UNNAMED_TOPOLOGY; + +/** + * Multi-threaded class that tracks the status of active tasks being processed. A single instance of this class is + * shared between all StreamThreads. + */ +public class TaskExecutionMetadata { + private Logger log; + + private final boolean hasNamedTopologies; + private final ConcurrentHashMap topologyNameToMetadata = new ConcurrentHashMap<>(); + + public TaskExecutionMetadata(final Set allTopologyNames) { + this.hasNamedTopologies = !(allTopologyNames.size() == 1 && allTopologyNames.contains(UNNAMED_TOPOLOGY)); + allTopologyNames.forEach(name -> topologyNameToMetadata.put(name, new NamedTopologyMetadata(name))); + } + + public void setLog(final LogContext logContext) { + log = logContext.logger(getClass()); + } + + public boolean canProcessTopology(final String topologyName) { + if (!hasNamedTopologies) { + return true; + } else { + return getMetadata(topologyName).canProcess(); + } + } + + public boolean canProcessTask(final Task task) { + final String topologyName = task.id().topologyName(); + if (!hasNamedTopologies) { + // TODO implement error handling/backoff for non-named topologies (needs KIP) + return true; + } else { + return getMetadata(topologyName).canProcessTask(task); + } + } + + public void registerTaskError(final Task task, final Throwable t) { + if (hasNamedTopologies) { + getMetadata(task.id().topologyName()).registerTaskError(task, t); + } + } + + public void registerTopology(final String topologyName) { + if (topologyNameToMetadata.containsKey(topologyName)) { + log.error("Topology {} is already registered in topology map.\n" + + "topologyNameToMetadata: {}", topologyName, topologyNameToMetadata); + throw new IllegalStateException("Tried to register new topology with execution metadata but " + + topologyName + " was already registered"); + } + topologyNameToMetadata.put(topologyName, new NamedTopologyMetadata(topologyName)); + log.debug("Registered topology {} with execution metadata", topologyName); + } + + public void unregisterTopology(final String topologyName) { + if (!topologyNameToMetadata.containsKey(topologyName)) { + log.error("Topology {} is not already registered in topology map.\n" + + "topologyNameToMetadata: {}", topologyName, topologyNameToMetadata); + throw new IllegalStateException("Tried to unregister a topology with execution metadata but " + + topologyName + " was not currently registered"); + } + topologyNameToMetadata.put(topologyName, new NamedTopologyMetadata(topologyName)); + log.debug("Unregistered topology {} with execution metadata", topologyName); + } + + /** + * Look up the metadata for this named topology + * + * @throws IllegalStateException if topology name is invalid + * @throws org.apache.kafka.streams.errors.UnknownTopologyException if the topology name is not found + */ + private NamedTopologyMetadata getMetadata(final String topologyName) { + if (topologyName == null || topologyName.equals(UNNAMED_TOPOLOGY)) { + log.error("Tried to look up metadata for named topology but topologyName was '{}'", topologyName); + throw new IllegalStateException("Invalid topology name for "); + } + + final NamedTopologyMetadata topologyMetadata = topologyNameToMetadata.get(topologyName); + if (topologyMetadata == null) { + log.error("Tried to look up metadata for named topology but could not find topologyName = '{}'", topologyName); + throw new UnknownTopologyException("Failed to check execution status", topologyName); + } else { + return topologyMetadata; + } + } + + static class NamedTopologyMetadata { + private final Logger log; + private final Set tasksToBackoff = new ConcurrentSkipListSet<>(); + + public NamedTopologyMetadata(final String topologyName) { + final LogContext logContext = new LogContext(String.format("topology-name [%s] ", topologyName)); + this.log = logContext.logger(NamedTopologyMetadata.class); + } + + public boolean canProcess() { + // TODO: during long task backoffs, pause the full topology to avoid it getting out of sync + return true; + } + + public boolean canProcessTask(final Task task) { + // TODO: implement true backoff, for now we just skip one iteration of processing a task upon error + final boolean canProcess = !tasksToBackoff.remove(task.id()); + if (!canProcess) { + log.info("Skipping processing iteration for task {}", task.id()); + } + return canProcess; + } + + public synchronized void registerTaskError(final Task task, final Throwable t) { + log.warn("Registered error {} for task {}", t.getMessage(), task.id()); + tasksToBackoff.add(task.id()); + } + } +} diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java index 4edc35b08b99b..efb533b66425d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java @@ -51,12 +51,15 @@ public class TaskExecutor { private final boolean hasNamedTopologies; private final ProcessingMode processingMode; private final Tasks tasks; + private final TaskExecutionMetadata taskExecutionMetadata; public TaskExecutor(final Tasks tasks, + final TaskExecutionMetadata taskExecutionMetadata, final ProcessingMode processingMode, final boolean hasNamedTopologies, final LogContext logContext) { this.tasks = tasks; + this.taskExecutionMetadata = taskExecutionMetadata; this.processingMode = processingMode; this.hasNamedTopologies = hasNamedTopologies; this.log = logContext.logger(getClass()); @@ -69,15 +72,25 @@ public TaskExecutor(final Tasks tasks, int process(final int maxNumRecords, final Time time) { int totalProcessed = 0; Task lastProcessed = null; - try { - for (final Task task : tasks.activeTasks()) { - lastProcessed = task; - totalProcessed += processTask(task, maxNumRecords, time); + + for (final Map.Entry> topologyEntry : tasks.activeTasksByTopology().entrySet()) { + final String topologyName = topologyEntry.getKey(); + final Set topologyActiveTasks = topologyEntry.getValue(); + if (taskExecutionMetadata.canProcessTopology(topologyName)) { + for (final Task task : topologyActiveTasks) { + try { + if (taskExecutionMetadata.canProcessTask(task)) { + lastProcessed = task; + totalProcessed += processTask(task, maxNumRecords, time); + } + } catch (final Throwable t) { + taskExecutionMetadata.registerTaskError(task, t); + tasks.removeTaskFromCuccessfullyProcessedBeforeClosing(lastProcessed); + commitSuccessfullyProcessedTasks(); + throw t; + } + } } - } catch (final Exception e) { - tasks.removeTaskFromCuccessfullyProcessedBeforeClosing(lastProcessed); - commitSuccessfullyProcessedTasks(); - throw e; } return totalProcessed; @@ -98,6 +111,7 @@ private long processTask(final Task task, final int maxNumRecords, final Time ti tasks.addToSuccessfullyProcessed(task); } } catch (final TimeoutException timeoutException) { + // TODO consolidate TimeoutException retries with general error handling task.maybeInitTaskTimeoutOrThrow(now, timeoutException); log.debug( String.format( @@ -132,7 +146,7 @@ private long processTask(final Task task, final int maxNumRecords, final Time ti * @return number of committed offsets, or -1 if we are in the middle of a rebalance and cannot commit */ int commitTasksAndMaybeUpdateCommittableOffsets(final Collection tasksToCommit, - final Map> consumedOffsetsAndMetadata) { + final Map> consumedOffsetsAndMetadata) { int committed = 0; for (final Task task : tasksToCommit) { // we need to call commitNeeded first since we need to update committable offsets 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 0ff6dc21ab2da..9cb2b9b2862ce 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 @@ -116,7 +116,13 @@ public class TaskManager { this.log = logContext.logger(getClass()); this.tasks = new Tasks(logContext, topologyMetadata, streamsMetrics, activeTaskCreator, standbyTaskCreator); - this.taskExecutor = new TaskExecutor(tasks, processingMode, topologyMetadata.hasNamedTopologies(), logContext); + this.taskExecutor = new TaskExecutor( + tasks, + topologyMetadata.taskExecutionMetadata(), + processingMode, + topologyMetadata.hasNamedTopologies(), + logContext + ); } void setMainConsumer(final Consumer mainConsumer) { 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 2740791f8f628..75c1f2687c0bd 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 @@ -35,6 +35,8 @@ import java.util.TreeMap; import java.util.stream.Collectors; +import static org.apache.kafka.streams.processor.internals.TopologyMetadata.getTopologyNameOrElseUnnamed; + class Tasks { private final Logger log; private final TopologyMetadata topologyMetadata; @@ -270,6 +272,15 @@ Collection activeTasks() { return readOnlyActiveTasks; } + Map> activeTasksByTopology() { + final Map> activeTasksByTopology = new HashMap<>(); + for (final Map.Entry taskEntry : readOnlyActiveTasksPerId.entrySet()) { + final String topologyName = getTopologyNameOrElseUnnamed(taskEntry.getKey().topologyName()); + activeTasksByTopology.computeIfAbsent(topologyName, k -> new HashSet<>()).add(taskEntry.getValue()); + } + return activeTasksByTopology; + } + Collection allTasks() { return readOnlyTasks; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java index 602abd34f364d..a569744b103f2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java @@ -69,6 +69,7 @@ public class TopologyMetadata { private final StreamsConfig config; private final ProcessingMode processingMode; private final TopologyVersion version; + private final TaskExecutionMetadata taskExecutionMetadata; private final ConcurrentNavigableMap builders; // Keep sorted by topology name for readability @@ -96,8 +97,8 @@ public TopologyVersionWaiters(final long topologyVersion, final KafkaFutureImpl< public TopologyMetadata(final InternalTopologyBuilder builder, final StreamsConfig config) { - version = new TopologyVersion(); - processingMode = StreamsConfigUtils.processingMode(config); + this.version = new TopologyVersion(); + this.processingMode = StreamsConfigUtils.processingMode(config); this.config = config; builders = new ConcurrentSkipListMap<>(); @@ -106,6 +107,7 @@ public TopologyMetadata(final InternalTopologyBuilder builder, } else { builders.put(UNNAMED_TOPOLOGY, builder); } + this.taskExecutionMetadata = new TaskExecutionMetadata(builders.keySet()); } public TopologyMetadata(final ConcurrentNavigableMap builders, @@ -115,16 +117,17 @@ public TopologyMetadata(final ConcurrentNavigableMap future, fina lock(); buildAndVerifyTopology(newTopologyBuilder); log.info("New NamedTopology passed validation and will be added {}, old topology version is {}", newTopologyBuilder.topologyName(), version.topologyVersion.get()); + taskExecutionMetadata.registerTopology(newTopologyBuilder.topologyName()); version.topologyVersion.incrementAndGet(); version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), future)); builders.put(newTopologyBuilder.topologyName(), newTopologyBuilder); @@ -238,6 +246,7 @@ public KafkaFuture unregisterTopology(final KafkaFutureImpl removeTo try { lock(); log.info("Beginning removal of NamedTopology {}, old topology version is {}", topologyName, version.topologyVersion.get()); + taskExecutionMetadata.unregisterTopology(topologyName); version.topologyVersion.incrementAndGet(); version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), removeTopologyFuture)); final InternalTopologyBuilder removedBuilder = builders.remove(topologyName); 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 new file mode 100644 index 0000000000000..997696fe06f5b --- /dev/null +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java @@ -0,0 +1,106 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.processor.internals; + +import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.processor.TaskId; +import org.apache.kafka.test.StreamsTestUtils; +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentSkipListMap; + +import static org.apache.kafka.streams.processor.internals.TopologyMetadata.UNNAMED_TOPOLOGY; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.core.IsEqual.equalTo; +import static org.junit.jupiter.api.Assertions.assertIterableEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class TasksTest { + + private static final String TOPOLOGY_NAME_0 = "TOPOLOGY_NAME_0"; + private static final String TOPOLOGY_NAME_1 = "TOPOLOGY_NAME_1"; + + private static final StreamTask TASK_0_0 = createActiveTask(new TaskId(0, 0)); + private static final StreamTask TASK_0_1 = createActiveTask(new TaskId(0, 1)); + + private static final StreamTask TASK_0_0_0 = createActiveTask(new TaskId(0, 0, TOPOLOGY_NAME_0)); + private static final StreamTask TASK_0_1_0 = createActiveTask(new TaskId(0, 1, TOPOLOGY_NAME_0)); + private static final StreamTask TASK_0_2_0 = createActiveTask(new TaskId(0, 2, TOPOLOGY_NAME_0)); + private static final StreamTask TASK_0_0_1 = createActiveTask(new TaskId(0, 0, TOPOLOGY_NAME_1)); + private static final StreamTask TASK_0_1_1 = createActiveTask(new TaskId(0, 1, TOPOLOGY_NAME_1)); + + @Test + void shouldReturnAllTasksInUnnamedTopology() { + final StreamsConfig streamsConfig = new StreamsConfig(StreamsTestUtils.getStreamsConfig()); + final Tasks tasks = new Tasks( + new LogContext("[test]"), + new TopologyMetadata(new ConcurrentSkipListMap<>(), streamsConfig), + null, + null, + null); + tasks.addTask(TASK_0_0); + tasks.addTask(TASK_0_1); + + final Map> activeTasksByTopology = tasks.activeTasksByTopology(); + + assertThat(activeTasksByTopology.size(), equalTo(1)); + assertThat(activeTasksByTopology.containsKey(UNNAMED_TOPOLOGY), is(true)); + assertIterableEquals(Arrays.asList(TASK_0_0, TASK_0_1), activeTasksByTopology.get(UNNAMED_TOPOLOGY)); + } + + @Test + void shouldReturnStreamTasksInNamedTopologies() { + final StreamsConfig streamsConfig = new StreamsConfig(StreamsTestUtils.getStreamsConfig()); + final Tasks tasks = new Tasks( + new LogContext("[test]"), + new TopologyMetadata(new ConcurrentSkipListMap<>(), streamsConfig), + null, + null, + null); + tasks.addTask(TASK_0_0_0); + tasks.addTask(TASK_0_0_1); + tasks.addTask(TASK_0_1_0); + tasks.addTask(TASK_0_1_1); + tasks.addTask(TASK_0_2_0); + + final Map> activeTasksByTopology = tasks.activeTasksByTopology(); + + assertThat(activeTasksByTopology.size(), equalTo(2)); + assertIterableEquals( + Arrays.asList(TASK_0_0_0, TASK_0_1_0, TASK_0_2_0), + activeTasksByTopology.get(TOPOLOGY_NAME_0) + ); + assertIterableEquals( + Arrays.asList(TASK_0_0_1, TASK_0_1_1), + activeTasksByTopology.get(TOPOLOGY_NAME_1) + ); + } + + private static StreamTask createActiveTask(final TaskId taskId) { + final StreamTask task = mock(StreamTask.class); + when(task.id()).thenReturn(taskId); + when(task.isActive()).thenReturn(true); + return task; + } +} \ No newline at end of file From ba81469c3f6af804252ec7518cb23ac8dbc84a0b Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Wed, 23 Feb 2022 01:52:28 -0800 Subject: [PATCH 2/8] fix unregister topology --- .../streams/processor/internals/TaskExecutionMetadata.java | 2 +- .../kafka/streams/integration/ErrorHandlingIntegrationTest.java | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) create mode 100644 streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java index b3a3749a28439..ac855826571b3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java @@ -88,7 +88,7 @@ public void unregisterTopology(final String topologyName) { throw new IllegalStateException("Tried to unregister a topology with execution metadata but " + topologyName + " was not currently registered"); } - topologyNameToMetadata.put(topologyName, new NamedTopologyMetadata(topologyName)); + topologyNameToMetadata.remove(topologyName); log.debug("Unregistered topology {} with execution metadata", topologyName); } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java new file mode 100644 index 0000000000000..dd1ce8359b249 --- /dev/null +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java @@ -0,0 +1,2 @@ +package org.apache.kafka.streams.integration;public class ErrorHandlingIntegrationTest { +} From b9b270bec4d6fc01347c776ad696c84e3840c733 Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Wed, 23 Feb 2022 23:25:42 -0800 Subject: [PATCH 3/8] backoff by 15s --- .../errors/UnknownTopologyException.java | 4 +- .../processor/internals/StreamTask.java | 1 + .../internals/TaskExecutionMetadata.java | 101 +++----- .../processor/internals/TaskExecutor.java | 22 +- .../processor/internals/TaskManager.java | 2 + .../processor/internals/TopologyMetadata.java | 3 - .../EmitOnChangeIntegrationTest.java | 80 +----- .../ErrorHandlingIntegrationTest.java | 230 +++++++++++++++++- 8 files changed, 279 insertions(+), 164 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/errors/UnknownTopologyException.java b/streams/src/main/java/org/apache/kafka/streams/errors/UnknownTopologyException.java index 2124405878def..d7644841a23b0 100644 --- a/streams/src/main/java/org/apache/kafka/streams/errors/UnknownTopologyException.java +++ b/streams/src/main/java/org/apache/kafka/streams/errors/UnknownTopologyException.java @@ -24,11 +24,11 @@ public class UnknownTopologyException extends StreamsException { private static final long serialVersionUID = 1L; public UnknownTopologyException(final String message, final String namedTopology) { - super(message + "due to being unable to locate a Topology named " + namedTopology); + super(message + " due to being unable to locate a Topology named " + namedTopology); } public UnknownTopologyException(final String message, final Throwable throwable, final String namedTopology) { - super(message + "due to being unable to locate a Topology named " + namedTopology, throwable); + super(message + " due to being unable to locate a Topology named " + namedTopology, throwable); } } 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 f86e89f73ff63..6f8227200e340 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 @@ -708,6 +708,7 @@ record = partitionGroup.nextRecord(recordInfo, wallClockTime); final TopicPartition partition = recordInfo.partition(); if (!(record instanceof CorruptedRecord)) { + log.info("SOPHIE: actually processing task {}", id()); doProcess(wallClockTime); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java index ac855826571b3..c6fa6f4811aa3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java @@ -17,12 +17,11 @@ package org.apache.kafka.streams.processor.internals; import org.apache.kafka.common.utils.LogContext; -import org.apache.kafka.streams.errors.UnknownTopologyException; import org.apache.kafka.streams.processor.TaskId; +import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentSkipListSet; import org.slf4j.Logger; import static org.apache.kafka.streams.processor.internals.TopologyMetadata.UNNAMED_TOPOLOGY; @@ -32,90 +31,45 @@ * shared between all StreamThreads. */ public class TaskExecutionMetadata { - private Logger log; - private final boolean hasNamedTopologies; - private final ConcurrentHashMap topologyNameToMetadata = new ConcurrentHashMap<>(); + // map of topologies experiencing errors/currently under backoff + private final ConcurrentHashMap topologyNameToErrorMetadata = new ConcurrentHashMap<>(); public TaskExecutionMetadata(final Set allTopologyNames) { this.hasNamedTopologies = !(allTopologyNames.size() == 1 && allTopologyNames.contains(UNNAMED_TOPOLOGY)); - allTopologyNames.forEach(name -> topologyNameToMetadata.put(name, new NamedTopologyMetadata(name))); - } - - public void setLog(final LogContext logContext) { - log = logContext.logger(getClass()); } public boolean canProcessTopology(final String topologyName) { if (!hasNamedTopologies) { return true; } else { - return getMetadata(topologyName).canProcess(); + final NamedTopologyMetadata metadata = topologyNameToErrorMetadata.get(topologyName); + return metadata == null || metadata.canProcess(); } } - public boolean canProcessTask(final Task task) { + public boolean canProcessTask(final Task task, final long now) { final String topologyName = task.id().topologyName(); if (!hasNamedTopologies) { // TODO implement error handling/backoff for non-named topologies (needs KIP) return true; } else { - return getMetadata(topologyName).canProcessTask(task); + final NamedTopologyMetadata metadata = topologyNameToErrorMetadata.get(topologyName); + return metadata == null || metadata.canProcessTask(task, now); } } - public void registerTaskError(final Task task, final Throwable t) { + public void registerTaskError(final Task task, final Throwable t, final long now) { if (hasNamedTopologies) { - getMetadata(task.id().topologyName()).registerTaskError(task, t); - } - } - - public void registerTopology(final String topologyName) { - if (topologyNameToMetadata.containsKey(topologyName)) { - log.error("Topology {} is already registered in topology map.\n" + - "topologyNameToMetadata: {}", topologyName, topologyNameToMetadata); - throw new IllegalStateException("Tried to register new topology with execution metadata but " - + topologyName + " was already registered"); - } - topologyNameToMetadata.put(topologyName, new NamedTopologyMetadata(topologyName)); - log.debug("Registered topology {} with execution metadata", topologyName); - } - - public void unregisterTopology(final String topologyName) { - if (!topologyNameToMetadata.containsKey(topologyName)) { - log.error("Topology {} is not already registered in topology map.\n" + - "topologyNameToMetadata: {}", topologyName, topologyNameToMetadata); - throw new IllegalStateException("Tried to unregister a topology with execution metadata but " - + topologyName + " was not currently registered"); - } - topologyNameToMetadata.remove(topologyName); - log.debug("Unregistered topology {} with execution metadata", topologyName); - } - - /** - * Look up the metadata for this named topology - * - * @throws IllegalStateException if topology name is invalid - * @throws org.apache.kafka.streams.errors.UnknownTopologyException if the topology name is not found - */ - private NamedTopologyMetadata getMetadata(final String topologyName) { - if (topologyName == null || topologyName.equals(UNNAMED_TOPOLOGY)) { - log.error("Tried to look up metadata for named topology but topologyName was '{}'", topologyName); - throw new IllegalStateException("Invalid topology name for "); - } - - final NamedTopologyMetadata topologyMetadata = topologyNameToMetadata.get(topologyName); - if (topologyMetadata == null) { - log.error("Tried to look up metadata for named topology but could not find topologyName = '{}'", topologyName); - throw new UnknownTopologyException("Failed to check execution status", topologyName); - } else { - return topologyMetadata; + final String topologyName = task.id().topologyName(); + topologyNameToErrorMetadata.computeIfAbsent(topologyName, n -> new NamedTopologyMetadata(topologyName)) + .registerTaskError(task, t, now); } } - static class NamedTopologyMetadata { + class NamedTopologyMetadata { private final Logger log; - private final Set tasksToBackoff = new ConcurrentSkipListSet<>(); + private final Map tasksToErrorTime = new ConcurrentHashMap<>(); public NamedTopologyMetadata(final String topologyName) { final LogContext logContext = new LogContext(String.format("topology-name [%s] ", topologyName)); @@ -127,18 +81,27 @@ public boolean canProcess() { return true; } - public boolean canProcessTask(final Task task) { - // TODO: implement true backoff, for now we just skip one iteration of processing a task upon error - final boolean canProcess = !tasksToBackoff.remove(task.id()); - if (!canProcess) { - log.info("Skipping processing iteration for task {}", task.id()); + public boolean canProcessTask(final Task task, final long now) { + // TODO: implement exponential backoff, for now we just wait 15s + final Long errorTime = tasksToErrorTime.get(task.id()); + if (errorTime == null) { + return true; + } else if (now - errorTime > 15000L) { + log.info("End backoff for task {} at t={}", task.id(), now); + tasksToErrorTime.remove(task.id()); + if (tasksToErrorTime.isEmpty()) { + topologyNameToErrorMetadata.remove(task.id().topologyName()); + } + return true; + } else { + log.debug("Skipping processing for unhealthy task {} at t={}", task.id(), now); + return false; } - return canProcess; } - public synchronized void registerTaskError(final Task task, final Throwable t) { - log.warn("Registered error {} for task {}", t.getMessage(), task.id()); - tasksToBackoff.add(task.id()); + public synchronized void registerTaskError(final Task task, final Throwable t, final long now) { + log.info("Begin backoff for unhealthy task {} at t={}", task.id(), now); + tasksToErrorTime.put(task.id(), now); } } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java index efb533b66425d..c0da661fc9699 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java @@ -78,13 +78,14 @@ int process(final int maxNumRecords, final Time time) { final Set topologyActiveTasks = topologyEntry.getValue(); if (taskExecutionMetadata.canProcessTopology(topologyName)) { for (final Task task : topologyActiveTasks) { + final long now = time.milliseconds(); try { - if (taskExecutionMetadata.canProcessTask(task)) { + if (taskExecutionMetadata.canProcessTask(task, now)) { lastProcessed = task; - totalProcessed += processTask(task, maxNumRecords, time); + totalProcessed += processTask(task, maxNumRecords, now, time); } } catch (final Throwable t) { - taskExecutionMetadata.registerTaskError(task, t); + taskExecutionMetadata.registerTaskError(task, t, now); tasks.removeTaskFromCuccessfullyProcessedBeforeClosing(lastProcessed); commitSuccessfullyProcessedTasks(); throw t; @@ -96,9 +97,9 @@ int process(final int maxNumRecords, final Time time) { return totalProcessed; } - private long processTask(final Task task, final int maxNumRecords, final Time time) { + private long processTask(final Task task, final int maxNumRecords, final long begin, final Time time) { int processed = 0; - long now = time.milliseconds(); + long now = begin; final long then = now; try { @@ -107,13 +108,14 @@ private long processTask(final Task task, final int maxNumRecords, final Time ti processed++; } // TODO: enable regardless of whether using named topologies - if (hasNamedTopologies && processingMode != EXACTLY_ONCE_V2) { + if (processed > 0 && hasNamedTopologies && processingMode != EXACTLY_ONCE_V2) { + log.trace("Successfully processed task {}", task.id()); tasks.addToSuccessfullyProcessed(task); } } catch (final TimeoutException timeoutException) { // TODO consolidate TimeoutException retries with general error handling task.maybeInitTaskTimeoutOrThrow(now, timeoutException); - log.debug( + log.error( String.format( "Could not complete processing records for %s due to the following exception; will move to next task and retry later", task.id()), @@ -124,11 +126,11 @@ private long processTask(final Task task, final int maxNumRecords, final Time ti "Will trigger a new rebalance and close all tasks as zombies together.", task.id()); throw e; } catch (final StreamsException e) { - log.error("Failed to process stream task {} due to the following error:", task.id(), e); + log.error(String.format("Failed to process stream task %s due to the following error:", task.id()), e); e.setTaskId(task.id()); throw e; } catch (final RuntimeException e) { - log.error("Failed to process stream task {} due to the following error:", task.id(), e); + log.error(String.format("Failed to process stream task %s due to the following error:", task.id()), e); throw new StreamsException(e, task.id()); } finally { now = time.milliseconds(); @@ -267,7 +269,7 @@ private void commitSuccessfullyProcessedTasks() { if (!tasks.successfullyProcessed().isEmpty()) { log.info("Streams encountered an error when processing tasks." + " Will commit all previously successfully processed tasks {}", - tasks.successfullyProcessed().toString()); + tasks.successfullyProcessed().stream().map(Task::id)); commitTasksAndMaybeUpdateCommittableOffsets(tasks.successfullyProcessed(), new HashMap<>()); } tasks.clearSuccessfullyProcessed(); 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 9cb2b9b2862ce..5864558017eb5 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 @@ -1062,6 +1062,8 @@ void addRecordsToTasks(final ConsumerRecords records) { throw new NullPointerException("Task was unexpectedly missing for partition " + partition); } + log.info("SOPHIE: adding records to task {}: {}", activeTask.id(), records.records(partition)); + activeTask.addRecords(partition, records.records(partition)); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java index a569744b103f2..dbcacf7534144 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java @@ -127,7 +127,6 @@ public TopologyMetadata(final ConcurrentNavigableMap future, fina lock(); buildAndVerifyTopology(newTopologyBuilder); log.info("New NamedTopology passed validation and will be added {}, old topology version is {}", newTopologyBuilder.topologyName(), version.topologyVersion.get()); - taskExecutionMetadata.registerTopology(newTopologyBuilder.topologyName()); version.topologyVersion.incrementAndGet(); version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), future)); builders.put(newTopologyBuilder.topologyName(), newTopologyBuilder); @@ -246,7 +244,6 @@ public KafkaFuture unregisterTopology(final KafkaFutureImpl removeTo try { lock(); log.info("Beginning removal of NamedTopology {}, old topology version is {}", topologyName, version.topologyVersion.get()); - taskExecutionMetadata.unregisterTopology(topologyName); version.topologyVersion.incrementAndGet(); version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), removeTopologyFuture)); final InternalTopologyBuilder removedBuilder = builders.remove(topologyName); diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java index 2d04070b402be..f5c92996bf9b1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java @@ -30,8 +30,6 @@ import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.processor.internals.namedtopology.KafkaStreamsNamedTopologyWrapper; -import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyBuilder; import org.apache.kafka.test.IntegrationTest; import org.apache.kafka.test.StreamsTestUtils; import org.apache.kafka.test.TestUtils; @@ -47,14 +45,11 @@ import java.util.Arrays; import java.util.Properties; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; 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.mkObjectProperties; import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.safeUniqueTestName; -import static org.hamcrest.CoreMatchers.equalTo; -import static org.hamcrest.MatcherAssert.assertThat; @Category(IntegrationTest.class) public class EmitOnChangeIntegrationTest { @@ -103,7 +98,7 @@ public void shouldEmitSameRecordAfterFailover() throws Exception { mkEntry(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 300000L), mkEntry(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class), mkEntry(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class), - mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000) + mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000) ) ); @@ -177,77 +172,4 @@ public void shouldEmitSameRecordAfterFailover() throws Exception { ); } } - - @Test - public void shouldEmitRecordsAfterFailures() throws Exception { - final Properties properties = mkObjectProperties( - mkMap( - mkEntry(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, CLUSTER.bootstrapServers()), - mkEntry(StreamsConfig.APPLICATION_ID_CONFIG, appId), - mkEntry(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory().getPath()), - mkEntry(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1), - mkEntry(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0), - mkEntry(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 300000L), - mkEntry(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class), - mkEntry(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class), - mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000) - ) - ); - - try (final KafkaStreamsNamedTopologyWrapper kafkaStreams = new KafkaStreamsNamedTopologyWrapper(properties)) { - kafkaStreams.setUncaughtExceptionHandler(exception -> StreamThreadExceptionResponse.REPLACE_THREAD); - - final NamedTopologyBuilder builder = kafkaStreams.newNamedTopologyBuilder("topology_A"); - final AtomicInteger noOutputExpected = new AtomicInteger(0); - final AtomicInteger twoOutputExpected = new AtomicInteger(0); - builder.stream(inputTopic2).peek((k, v) -> twoOutputExpected.incrementAndGet()).to(outputTopic2); - builder.stream(inputTopic) - .peek((k, v) -> { - throw new RuntimeException("Kaboom"); - }) - .peek((k, v) -> noOutputExpected.incrementAndGet()) - .to(outputTopic); - - kafkaStreams.addNamedTopology(builder.build()); - - StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams); - IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( - inputTopic, - Arrays.asList( - new KeyValue<>(1, "A") - ), - TestUtils.producerConfig( - CLUSTER.bootstrapServers(), - IntegerSerializer.class, - StringSerializer.class, - new Properties()), - 0L); - IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( - inputTopic2, - Arrays.asList( - new KeyValue<>(1, "A"), - new KeyValue<>(1, "B") - ), - TestUtils.producerConfig( - CLUSTER.bootstrapServers(), - IntegerSerializer.class, - StringSerializer.class, - new Properties()), - 0L); - IntegrationTestUtils.waitUntilFinalKeyValueRecordsReceived( - TestUtils.consumerConfig( - CLUSTER.bootstrapServers(), - IntegerDeserializer.class, - StringDeserializer.class - ), - outputTopic2, - Arrays.asList( - new KeyValue<>(1, "A"), - new KeyValue<>(1, "B") - ) - ); - assertThat(noOutputExpected.get(), equalTo(0)); - assertThat(twoOutputExpected.get(), equalTo(2)); - } - } } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java index dd1ce8359b249..2415e3b0f82a6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java @@ -1,2 +1,230 @@ -package org.apache.kafka.streams.integration;public class ErrorHandlingIntegrationTest { +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.integration; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.serialization.IntegerDeserializer; +import org.apache.kafka.common.serialization.IntegerSerializer; +import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse; +import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; +import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; +import org.apache.kafka.streams.processor.internals.namedtopology.KafkaStreamsNamedTopologyWrapper; +import org.apache.kafka.streams.processor.internals.namedtopology.NamedTopologyBuilder; +import org.apache.kafka.test.IntegrationTest; +import org.apache.kafka.test.StreamsTestUtils; +import org.apache.kafka.test.TestUtils; + +import java.io.IOException; +import java.util.Arrays; +import java.util.Properties; +import java.util.concurrent.atomic.AtomicInteger; +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 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.mkObjectProperties; +import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.safeUniqueTestName; + +import static org.hamcrest.CoreMatchers.equalTo; +import static org.hamcrest.MatcherAssert.assertThat; + +@Category(IntegrationTest.class) +public class ErrorHandlingIntegrationTest { + + private static final EmbeddedKafkaCluster CLUSTER = new EmbeddedKafkaCluster(1); + + @BeforeClass + public static void startCluster() throws IOException { + CLUSTER.start(); + } + + @AfterClass + public static void closeCluster() { + CLUSTER.stop(); + } + + @Rule + public TestName testName = new TestName(); + + private final String testId = safeUniqueTestName(getClass(), testName); + private final String appId = "appId_" + testId; + private final Properties properties = props(); + + // Task 0 + private final String inputTopic = "input" + testId; + private final String outputTopic = "output" + testId; + // Task 1 + private final String errorInputTopic = "error-input" + testId; + private final String errorOutputTopic = "error-output" + testId; + + @Before + public void setup() { + IntegrationTestUtils.cleanStateBeforeTest(CLUSTER, errorInputTopic, errorOutputTopic, inputTopic, outputTopic); + } + + private Properties props() { + return mkObjectProperties( + mkMap( + mkEntry(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, CLUSTER.bootstrapServers()), + mkEntry(StreamsConfig.APPLICATION_ID_CONFIG, appId), + mkEntry(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory(appId).getPath()), + mkEntry(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0), + mkEntry(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 15000L), + mkEntry(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class), + mkEntry(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class), + mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000), + mkEntry(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 10000) + ) + ); + } +/* + @Test + public void shouldSkipTaskIteration() throws Exception { + final AtomicInteger numErrors = new AtomicInteger(0); + final AtomicInteger task1Processed = new AtomicInteger(0); + final AtomicInteger task2Processed = new AtomicInteger(0); + + try (final KafkaStreamsNamedTopologyWrapper kafkaStreams = new KafkaStreamsNamedTopologyWrapper(properties)) { + kafkaStreams.setUncaughtExceptionHandler(exception -> { + numErrors.incrementAndGet(); + return StreamThreadExceptionResponse.REPLACE_THREAD; + }); + + final NamedTopologyBuilder builder = kafkaStreams.newNamedTopologyBuilder("topology_A"); + + builder.stream(inputTopic2).peek((k, v) -> task2Processed.incrementAndGet()).to(outputTopic2); + builder.stream(inputTopic) + .peek((k, v) -> { + throw new RuntimeException("Kaboom"); + }) + .peek((k, v) -> task1Processed.incrementAndGet()) + .to(outputTopic); + + kafkaStreams.addNamedTopology(builder.build()); + + StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams); + IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( + inputTopic, + Arrays.asList( + new KeyValue<>(1, "A") + ), + TestUtils.producerConfig( + CLUSTER.bootstrapServers(), + IntegerSerializer.class, + StringSerializer.class, + new Properties()), + 0L); + IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( + inputTopic2, + Arrays.asList( + new KeyValue<>(1, "A"), + new KeyValue<>(1, "B") + ), + TestUtils.producerConfig( + CLUSTER.bootstrapServers(), + IntegerSerializer.class, + StringSerializer.class, + new Properties()), + 0L); + IntegrationTestUtils.waitUntilFinalKeyValueRecordsReceived( + TestUtils.consumerConfig( + CLUSTER.bootstrapServers(), + IntegerDeserializer.class, + StringDeserializer.class + ), + outputTopic2, + Arrays.asList( + new KeyValue<>(1, "A"), + new KeyValue<>(1, "B") + ) + ); + assertThat(task1Processed.get(), equalTo(0)); + assertThat(task2Processed.get(), equalTo(2)); + } + } +*/ + @Test + public void shouldBackOffTaskAndEmitDataWithinSameTopology() throws Exception { + final AtomicInteger noOutputExpected = new AtomicInteger(0); + final AtomicInteger outputExpected = new AtomicInteger(0); + + try (final KafkaStreamsNamedTopologyWrapper kafkaStreams = new KafkaStreamsNamedTopologyWrapper(properties)) { + kafkaStreams.setUncaughtExceptionHandler(exception -> StreamThreadExceptionResponse.REPLACE_THREAD); + + final NamedTopologyBuilder builder = kafkaStreams.newNamedTopologyBuilder("topology_A"); + builder.stream(inputTopic).peek((k, v) -> outputExpected.incrementAndGet()).to(outputTopic); + builder.stream(errorInputTopic) + .peek((k, v) -> { + throw new RuntimeException("Kaboom"); + }) + .peek((k, v) -> noOutputExpected.incrementAndGet()) + .to(errorOutputTopic); + + kafkaStreams.addNamedTopology(builder.build()); + + StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams); + IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( + errorInputTopic, + Arrays.asList( + new KeyValue<>(1, "A") + ), + TestUtils.producerConfig( + CLUSTER.bootstrapServers(), + IntegerSerializer.class, + StringSerializer.class, + new Properties()), + 0L); + IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( + inputTopic, + Arrays.asList( + new KeyValue<>(1, "A"), + new KeyValue<>(1, "B") + ), + TestUtils.producerConfig( + CLUSTER.bootstrapServers(), + IntegerSerializer.class, + StringSerializer.class, + new Properties()), + 0L); + IntegrationTestUtils.waitUntilFinalKeyValueRecordsReceived( + TestUtils.consumerConfig( + CLUSTER.bootstrapServers(), + IntegerDeserializer.class, + StringDeserializer.class + ), + outputTopic, + Arrays.asList( + new KeyValue<>(1, "A"), + new KeyValue<>(1, "B") + ) + ); + assertThat(noOutputExpected.get(), equalTo(0)); + assertThat(outputExpected.get(), equalTo(2)); + } + } } From f204b84604ba68383514bf4f79c99d4dc1887555 Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Wed, 23 Feb 2022 23:31:49 -0800 Subject: [PATCH 4/8] don't group tasks by topology during processing --- .../internals/TaskExecutionMetadata.java | 11 +------- .../processor/internals/TaskExecutor.java | 28 ++++++++----------- 2 files changed, 12 insertions(+), 27 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java index c6fa6f4811aa3..63b53aebdcf84 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutionMetadata.java @@ -39,15 +39,6 @@ public TaskExecutionMetadata(final Set allTopologyNames) { this.hasNamedTopologies = !(allTopologyNames.size() == 1 && allTopologyNames.contains(UNNAMED_TOPOLOGY)); } - public boolean canProcessTopology(final String topologyName) { - if (!hasNamedTopologies) { - return true; - } else { - final NamedTopologyMetadata metadata = topologyNameToErrorMetadata.get(topologyName); - return metadata == null || metadata.canProcess(); - } - } - public boolean canProcessTask(final Task task, final long now) { final String topologyName = task.id().topologyName(); if (!hasNamedTopologies) { @@ -55,7 +46,7 @@ public boolean canProcessTask(final Task task, final long now) { return true; } else { final NamedTopologyMetadata metadata = topologyNameToErrorMetadata.get(topologyName); - return metadata == null || metadata.canProcessTask(task, now); + return metadata == null || (metadata.canProcess() && metadata.canProcessTask(task, now)); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java index c0da661fc9699..cad03fbd1b3ec 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskExecutor.java @@ -73,24 +73,18 @@ int process(final int maxNumRecords, final Time time) { int totalProcessed = 0; Task lastProcessed = null; - for (final Map.Entry> topologyEntry : tasks.activeTasksByTopology().entrySet()) { - final String topologyName = topologyEntry.getKey(); - final Set topologyActiveTasks = topologyEntry.getValue(); - if (taskExecutionMetadata.canProcessTopology(topologyName)) { - for (final Task task : topologyActiveTasks) { - final long now = time.milliseconds(); - try { - if (taskExecutionMetadata.canProcessTask(task, now)) { - lastProcessed = task; - totalProcessed += processTask(task, maxNumRecords, now, time); - } - } catch (final Throwable t) { - taskExecutionMetadata.registerTaskError(task, t, now); - tasks.removeTaskFromCuccessfullyProcessedBeforeClosing(lastProcessed); - commitSuccessfullyProcessedTasks(); - throw t; - } + for (final Task task : tasks.activeTasks()) { + final long now = time.milliseconds(); + try { + if (taskExecutionMetadata.canProcessTask(task, now)) { + lastProcessed = task; + totalProcessed += processTask(task, maxNumRecords, now, time); } + } catch (final Throwable t) { + taskExecutionMetadata.registerTaskError(task, t, now); + tasks.removeTaskFromCuccessfullyProcessedBeforeClosing(lastProcessed); + commitSuccessfullyProcessedTasks(); + throw t; } } From 2f007c25ce1ffe16188a563959482213fdd289b7 Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Thu, 24 Feb 2022 01:31:38 -0800 Subject: [PATCH 5/8] improve logging --- .../apache/kafka/streams/processor/internals/StreamTask.java | 1 - .../apache/kafka/streams/processor/internals/StreamThread.java | 3 ++- .../apache/kafka/streams/processor/internals/TaskManager.java | 2 -- 3 files changed, 2 insertions(+), 4 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 6f8227200e340..f86e89f73ff63 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 @@ -708,7 +708,6 @@ record = partitionGroup.nextRecord(recordInfo, wallClockTime); final TopicPartition partition = recordInfo.partition(); if (!(record instanceof CorruptedRecord)) { - log.info("SOPHIE: actually processing task {}", id()); doProcess(wallClockTime); } 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 2a53e32d98f28..3af07bbbd646e 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 @@ -932,7 +932,8 @@ private long pollPhase() { .ifPresent(t -> taskManager.updateTaskEndMetadata(topicPartition, t.offset())); } - log.debug("Main Consumer poll completed in {} ms and fetched {} records", pollLatency, numRecords); + log.debug("Main Consumer poll completed in {} ms and fetched {} records from partitions {}", + pollLatency, numRecords, records.partitions()); pollSensor.record(pollLatency, now); 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 5864558017eb5..9cb2b9b2862ce 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 @@ -1062,8 +1062,6 @@ void addRecordsToTasks(final ConsumerRecords records) { throw new NullPointerException("Task was unexpectedly missing for partition " + partition); } - log.info("SOPHIE: adding records to task {}: {}", activeTask.id(), records.records(partition)); - activeTask.addRecords(partition, records.records(partition)); } } From 20c0b4d0135d7a294fa5d008aec82724829403e3 Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Thu, 24 Feb 2022 01:34:02 -0800 Subject: [PATCH 6/8] clean up test --- .../ErrorHandlingIntegrationTest.java | 65 ------------------- 1 file changed, 65 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java index 2415e3b0f82a6..872d66f9cefd9 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java @@ -102,72 +102,7 @@ private Properties props() { ) ); } -/* - @Test - public void shouldSkipTaskIteration() throws Exception { - final AtomicInteger numErrors = new AtomicInteger(0); - final AtomicInteger task1Processed = new AtomicInteger(0); - final AtomicInteger task2Processed = new AtomicInteger(0); - - try (final KafkaStreamsNamedTopologyWrapper kafkaStreams = new KafkaStreamsNamedTopologyWrapper(properties)) { - kafkaStreams.setUncaughtExceptionHandler(exception -> { - numErrors.incrementAndGet(); - return StreamThreadExceptionResponse.REPLACE_THREAD; - }); - - final NamedTopologyBuilder builder = kafkaStreams.newNamedTopologyBuilder("topology_A"); - - builder.stream(inputTopic2).peek((k, v) -> task2Processed.incrementAndGet()).to(outputTopic2); - builder.stream(inputTopic) - .peek((k, v) -> { - throw new RuntimeException("Kaboom"); - }) - .peek((k, v) -> task1Processed.incrementAndGet()) - .to(outputTopic); - - kafkaStreams.addNamedTopology(builder.build()); - StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams); - IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( - inputTopic, - Arrays.asList( - new KeyValue<>(1, "A") - ), - TestUtils.producerConfig( - CLUSTER.bootstrapServers(), - IntegerSerializer.class, - StringSerializer.class, - new Properties()), - 0L); - IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp( - inputTopic2, - Arrays.asList( - new KeyValue<>(1, "A"), - new KeyValue<>(1, "B") - ), - TestUtils.producerConfig( - CLUSTER.bootstrapServers(), - IntegerSerializer.class, - StringSerializer.class, - new Properties()), - 0L); - IntegrationTestUtils.waitUntilFinalKeyValueRecordsReceived( - TestUtils.consumerConfig( - CLUSTER.bootstrapServers(), - IntegerDeserializer.class, - StringDeserializer.class - ), - outputTopic2, - Arrays.asList( - new KeyValue<>(1, "A"), - new KeyValue<>(1, "B") - ) - ); - assertThat(task1Processed.get(), equalTo(0)); - assertThat(task2Processed.get(), equalTo(2)); - } - } -*/ @Test public void shouldBackOffTaskAndEmitDataWithinSameTopology() throws Exception { final AtomicInteger noOutputExpected = new AtomicInteger(0); From 36560d24d7e4b9e0a5f6e1a392b069b0bc3f093f Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Thu, 24 Feb 2022 01:47:25 -0800 Subject: [PATCH 7/8] undo test config change --- .../streams/integration/EmitOnChangeIntegrationTest.java | 2 +- .../streams/integration/ErrorHandlingIntegrationTest.java | 4 +--- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java index f5c92996bf9b1..25f0f3055042d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EmitOnChangeIntegrationTest.java @@ -98,7 +98,7 @@ public void shouldEmitSameRecordAfterFailover() throws Exception { mkEntry(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 300000L), mkEntry(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class), mkEntry(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class), - mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000) + mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000) ) ); diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java index 872d66f9cefd9..40d2cf98027ca 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ErrorHandlingIntegrationTest.java @@ -97,9 +97,7 @@ private Properties props() { mkEntry(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 15000L), mkEntry(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class), mkEntry(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class), - mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000), - mkEntry(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 10000) - ) + mkEntry(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000)) ); } From c9dc1dd9ffb0cba1775e08b4467caabbb355c0ca Mon Sep 17 00:00:00 2001 From: "A. Sophie Blee-Goldman" Date: Thu, 24 Feb 2022 01:59:38 -0800 Subject: [PATCH 8/8] remove unused method --- .../streams/processor/internals/Tasks.java | 11 -- .../processor/internals/TasksTest.java | 106 ------------------ 2 files changed, 117 deletions(-) delete mode 100644 streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java 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 75c1f2687c0bd..2740791f8f628 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 @@ -35,8 +35,6 @@ import java.util.TreeMap; import java.util.stream.Collectors; -import static org.apache.kafka.streams.processor.internals.TopologyMetadata.getTopologyNameOrElseUnnamed; - class Tasks { private final Logger log; private final TopologyMetadata topologyMetadata; @@ -272,15 +270,6 @@ Collection activeTasks() { return readOnlyActiveTasks; } - Map> activeTasksByTopology() { - final Map> activeTasksByTopology = new HashMap<>(); - for (final Map.Entry taskEntry : readOnlyActiveTasksPerId.entrySet()) { - final String topologyName = getTopologyNameOrElseUnnamed(taskEntry.getKey().topologyName()); - activeTasksByTopology.computeIfAbsent(topologyName, k -> new HashSet<>()).add(taskEntry.getValue()); - } - return activeTasksByTopology; - } - Collection allTasks() { return readOnlyTasks; } 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 deleted file mode 100644 index 997696fe06f5b..0000000000000 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TasksTest.java +++ /dev/null @@ -1,106 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.kafka.streams.processor.internals; - -import org.apache.kafka.common.utils.LogContext; -import org.apache.kafka.streams.StreamsConfig; -import org.apache.kafka.streams.processor.TaskId; -import org.apache.kafka.test.StreamsTestUtils; -import org.junit.jupiter.api.Test; - -import java.util.Arrays; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.ConcurrentSkipListMap; - -import static org.apache.kafka.streams.processor.internals.TopologyMetadata.UNNAMED_TOPOLOGY; - -import static org.hamcrest.MatcherAssert.assertThat; -import static org.hamcrest.Matchers.is; -import static org.hamcrest.core.IsEqual.equalTo; -import static org.junit.jupiter.api.Assertions.assertIterableEquals; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; - -class TasksTest { - - private static final String TOPOLOGY_NAME_0 = "TOPOLOGY_NAME_0"; - private static final String TOPOLOGY_NAME_1 = "TOPOLOGY_NAME_1"; - - private static final StreamTask TASK_0_0 = createActiveTask(new TaskId(0, 0)); - private static final StreamTask TASK_0_1 = createActiveTask(new TaskId(0, 1)); - - private static final StreamTask TASK_0_0_0 = createActiveTask(new TaskId(0, 0, TOPOLOGY_NAME_0)); - private static final StreamTask TASK_0_1_0 = createActiveTask(new TaskId(0, 1, TOPOLOGY_NAME_0)); - private static final StreamTask TASK_0_2_0 = createActiveTask(new TaskId(0, 2, TOPOLOGY_NAME_0)); - private static final StreamTask TASK_0_0_1 = createActiveTask(new TaskId(0, 0, TOPOLOGY_NAME_1)); - private static final StreamTask TASK_0_1_1 = createActiveTask(new TaskId(0, 1, TOPOLOGY_NAME_1)); - - @Test - void shouldReturnAllTasksInUnnamedTopology() { - final StreamsConfig streamsConfig = new StreamsConfig(StreamsTestUtils.getStreamsConfig()); - final Tasks tasks = new Tasks( - new LogContext("[test]"), - new TopologyMetadata(new ConcurrentSkipListMap<>(), streamsConfig), - null, - null, - null); - tasks.addTask(TASK_0_0); - tasks.addTask(TASK_0_1); - - final Map> activeTasksByTopology = tasks.activeTasksByTopology(); - - assertThat(activeTasksByTopology.size(), equalTo(1)); - assertThat(activeTasksByTopology.containsKey(UNNAMED_TOPOLOGY), is(true)); - assertIterableEquals(Arrays.asList(TASK_0_0, TASK_0_1), activeTasksByTopology.get(UNNAMED_TOPOLOGY)); - } - - @Test - void shouldReturnStreamTasksInNamedTopologies() { - final StreamsConfig streamsConfig = new StreamsConfig(StreamsTestUtils.getStreamsConfig()); - final Tasks tasks = new Tasks( - new LogContext("[test]"), - new TopologyMetadata(new ConcurrentSkipListMap<>(), streamsConfig), - null, - null, - null); - tasks.addTask(TASK_0_0_0); - tasks.addTask(TASK_0_0_1); - tasks.addTask(TASK_0_1_0); - tasks.addTask(TASK_0_1_1); - tasks.addTask(TASK_0_2_0); - - final Map> activeTasksByTopology = tasks.activeTasksByTopology(); - - assertThat(activeTasksByTopology.size(), equalTo(2)); - assertIterableEquals( - Arrays.asList(TASK_0_0_0, TASK_0_1_0, TASK_0_2_0), - activeTasksByTopology.get(TOPOLOGY_NAME_0) - ); - assertIterableEquals( - Arrays.asList(TASK_0_0_1, TASK_0_1_1), - activeTasksByTopology.get(TOPOLOGY_NAME_1) - ); - } - - private static StreamTask createActiveTask(final TaskId taskId) { - final StreamTask task = mock(StreamTask.class); - when(task.id()).thenReturn(taskId); - when(task.isActive()).thenReturn(true); - return task; - } -} \ No newline at end of file