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 6e7a3aa9ac1e9..3a57cc6258fd8 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 @@ -644,7 +644,7 @@ void runOnce() { if (records != null && !records.isEmpty()) { pollSensor.record(pollLatency, now); pollRecordsSensor.record(records.count(), now); - addRecordsToTasks(records); + taskManager.addRecordsToTasks(records); } // Shutdown hook could potentially be triggered and transit the thread state to PENDING_SHUTDOWN during #pollRequests(). @@ -819,25 +819,6 @@ private void addToResetList(final TopicPartition partition, final Set records) { - for (final TopicPartition partition : records.partitions()) { - final Task task = taskManager.taskForInputPartition(partition); - - if (task == null) { - log.error("Unable to locate active task for received-record partition {}. Current tasks: {}", - partition, taskManager.toString(">")); - throw new NullPointerException("Task was unexpectedly missing for partition " + partition); - } - - task.addRecords(partition, records.records(partition)); - } - } - /** * 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 d1be8a3f3055e..5c55093e1a711 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 @@ -21,6 +21,7 @@ import org.apache.kafka.clients.admin.RecordsToDelete; import org.apache.kafka.clients.consumer.CommitFailedException; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Metric; @@ -86,7 +87,7 @@ public class TaskManager { private boolean rebalanceInProgress = false; // if we are in the middle of a rebalance, it is not safe to commit // includes assigned & initialized tasks and unassigned tasks we locked temporarily during rebalance - private Set lockedTaskDirectories = new HashSet<>(); + private final Set lockedTaskDirectories = new HashSet<>(); TaskManager(final ChangelogReader changelogReader, final UUID processId, @@ -743,10 +744,6 @@ Set standbyTaskIds() { .collect(Collectors.toSet()); } - Task taskForInputPartition(final TopicPartition partition) { - return partitionToTask.get(partition); - } - Map tasks() { // 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. @@ -778,6 +775,25 @@ int commitAll() { return commit(tasks.values()); } + /** + * Take records and add them to each respective task + * + * @param records Records, can be null + */ + void addRecordsToTasks(final ConsumerRecords records) { + for (final TopicPartition partition : records.partitions()) { + final Task task = partitionToTask.get(partition); + + if (task == null) { + log.error("Unable to locate active task for received-record partition {}. Current tasks: {}", + partition, toString(">")); + throw new NullPointerException("Task was unexpectedly missing for partition " + partition); + } + + task.addRecords(partition, records.records(partition)); + } + } + /** * @throws TaskMigratedException if committing offsets failed (non-EOS) * or if the task producer got fenced (EOS) 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 4d38b2f3d9f4d..7f31d7e8fd261 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 @@ -694,7 +694,6 @@ public void shouldUpdateInputPartitionsAfterRebalance() { assertThat(taskManager.tryToCompleteRestoration(), is(true)); assertThat(task00.state(), is(Task.State.RUNNING)); assertEquals(newPartitionsSet, task00.inputPartitions()); - assertEquals(task00, taskManager.taskForInputPartition(t1p1)); verify(activeTaskCreator, consumer, changeLogReader); }