From 763aa52a4658848e69345e3e5708d36ccb2439c4 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Wed, 8 Jul 2020 17:54:45 -0700 Subject: [PATCH 1/6] improve log4j for per-consumer assignment --- .../internals/StreamsPartitionAssignor.java | 35 +++++-- .../internals/assignment/ClientState.java | 91 +++++++++++-------- 2 files changed, 81 insertions(+), 45 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 9352a3b3384f1..4396c902ce1d9 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -120,8 +120,8 @@ public int hashCode() { private static class ClientMetadata { private final HostInfo hostInfo; - private final SortedSet consumers; private final ClientState state; + private final SortedSet consumers; ClientMetadata(final String endPoint) { @@ -725,8 +725,10 @@ private boolean assignTasksToClients(final Set allSourceTopics, statefulTasks, assignmentConfigs); - log.info("Assigned tasks to clients as {}{}.", - Utils.NL, clientStates.entrySet().stream().map(Map.Entry::toString).collect(Collectors.joining(Utils.NL))); + log.info("Assigned tasks to clients as: {}{}.", Utils.NL, + clientStates.entrySet().stream() + .map(entry -> entry.getKey() + "=" + entry.getValue().currentAssignment()) + .collect(Collectors.joining(Utils.NL))); return probingRebalanceNeeded; } @@ -931,7 +933,9 @@ private Map computeNewAssignment(final Map computeNewAssignment(final Map assignment, if (!activeTasksRemovedPendingRevokation.isEmpty()) { // TODO: once KAFKA-10078 is resolved we can leave it to the client to trigger this rebalance - log.info("Requesting followup rebalance be scheduled immediately due to tasks changing ownership."); + log.info("Requesting {} followup rebalance be scheduled immediately due to tasks changing ownership.", consumer); info.setNextRebalanceTime(0L); followupRebalanceRequiredForRevokedTasks = true; // Don't bother to schedule a probing rebalance if an immediate one is already scheduled shouldEncodeProbingRebalance = false; } else if (shouldEncodeProbingRebalance) { final long nextRebalanceTimeMs = time.milliseconds() + probingRebalanceIntervalMs(); - log.info("Requesting followup rebalance be scheduled for {} ms to probe for caught-up replica tasks.", nextRebalanceTimeMs); + log.info("Requesting {} followup rebalance be scheduled for {} ms to probe for caught-up replica tasks.", + consumer, nextRebalanceTimeMs); info.setNextRebalanceTime(nextRebalanceTimeMs); shouldEncodeProbingRebalance = false; } @@ -1061,8 +1076,11 @@ private Set populateActiveTaskAndPartitionsLists(final List assignedPartitions = new ArrayList<>(); final Set removedActiveTasks = new TreeSet<>(); - // Build up list of all assigned partition-task pairs for (final TaskId taskId : activeTasksForConsumer) { + // Populate the consumer for assigned tasks without considering revocation, + // this is for debugging purposes only + clientState.assignActiveToConsumer(taskId, consumer); + final List assignedPartitionsForTask = new ArrayList<>(); for (final TopicPartition partition : partitionsForTask.get(taskId)) { final String oldOwner = clientState.previousOwnerForPartition(partition); @@ -1113,6 +1131,7 @@ private Map> buildStandbyTaskMap(final String consum final ClientState clientState) { final Map> standbyTaskMap = new HashMap<>(); for (final TaskId task : standbys) { + clientState.assignStandbyToConsumer(task, consumer); standbyTaskMap.put(task, partitionsForTask.get(task)); } for (final TaskId task : tasksRevoked) { @@ -1121,6 +1140,8 @@ private Map> buildStandbyTaskMap(final String consum task, consumer ); + + clientState.assignStandbyToConsumer(task, consumer); standbyTaskMap.put(task, partitionsForTask.get(task)); // This has no effect on the assignment, as we'll never consult the ClientState again, but diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java index c98ed9b8af542..318f1c0a9a1f2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java @@ -16,19 +16,20 @@ */ package org.apache.kafka.streams.processor.internals.assignment; -import java.util.stream.Collectors; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.Task; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; import java.util.Collection; import java.util.Comparator; import java.util.HashSet; +import java.util.List; import java.util.Map; import java.util.Set; -import java.util.SortedMap; +import java.util.stream.Collectors; import java.util.TreeMap; import java.util.TreeSet; import java.util.UUID; @@ -50,10 +51,14 @@ public class ClientState { private final Set prevStandbyTasks; private final Map> consumerToPreviousStatefulTaskIds; + private final Map> consumerToPreviousActiveTaskIds; + private final Map> consumerToAssignedActiveTaskIds; + private final Map> consumerToAssignedStandbyTaskIds; private final Map ownedPartitions; private final Map taskOffsetSums; // contains only stateful tasks we previously owned private final Map taskLagTotals; // contains lag for all stateful tasks in the app topology + private int capacity; public ClientState() { @@ -66,32 +71,16 @@ public ClientState() { prevActiveTasks = new TreeSet<>(); prevStandbyTasks = new TreeSet<>(); consumerToPreviousStatefulTaskIds = new TreeMap<>(); + consumerToPreviousActiveTaskIds = new TreeMap<>(); + consumerToAssignedActiveTaskIds = new TreeMap<>(); + consumerToAssignedStandbyTaskIds = new TreeMap<>(); ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = new TreeMap<>(); taskLagTotals = new TreeMap<>(); this.capacity = capacity; } - private ClientState(final Set activeTasks, - final Set standbyTasks, - final Set prevActiveTasks, - final Set prevStandbyTasks, - final Map> consumerToPreviousStatefulTaskIds, - final SortedMap ownedPartitions, - final Map taskOffsetSums, - final Map taskLagTotals, - final int capacity) { - this.activeTasks = activeTasks; - this.standbyTasks = standbyTasks; - this.prevActiveTasks = prevActiveTasks; - this.prevStandbyTasks = prevStandbyTasks; - this.consumerToPreviousStatefulTaskIds = consumerToPreviousStatefulTaskIds; - this.ownedPartitions = ownedPartitions; - this.taskOffsetSums = taskOffsetSums; - this.taskLagTotals = taskLagTotals; - this.capacity = capacity; - } - + // For testing only public ClientState(final Set previousActiveTasks, final Set previousStandbyTasks, final Map taskLagTotals, @@ -101,27 +90,15 @@ public ClientState(final Set previousActiveTasks, prevActiveTasks = unmodifiableSet(new TreeSet<>(previousActiveTasks)); prevStandbyTasks = unmodifiableSet(new TreeSet<>(previousStandbyTasks)); consumerToPreviousStatefulTaskIds = new TreeMap<>(); + consumerToPreviousActiveTaskIds = new TreeMap<>(); + consumerToAssignedActiveTaskIds = new TreeMap<>(); + consumerToAssignedStandbyTaskIds = new TreeMap<>(); ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = emptyMap(); this.taskLagTotals = unmodifiableMap(taskLagTotals); this.capacity = capacity; } - public ClientState copy() { - final TreeMap newOwnedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); - newOwnedPartitions.putAll(ownedPartitions); - return new ClientState( - new TreeSet<>(activeTasks), - new TreeSet<>(standbyTasks), - new TreeSet<>(prevActiveTasks), - new TreeSet<>(prevStandbyTasks), - new TreeMap<>(consumerToPreviousStatefulTaskIds), - newOwnedPartitions, - new TreeMap<>(taskOffsetSums), - new TreeMap<>(taskLagTotals), - capacity); - } - int capacity() { return capacity; } @@ -150,6 +127,39 @@ public void assignActiveTasks(final Collection tasks) { activeTasks.addAll(tasks); } + public void assignActiveToConsumer(final TaskId task, final String consumer) { + consumerToAssignedActiveTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + } + + public void assignStandbyToConsumer(final TaskId task, final String consumer) { + consumerToAssignedStandbyTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + } + + public Map> prevOwnedActiveByConsumer() { + return consumerToPreviousActiveTaskIds; + } + + public Map> prevOwnedStandbyByConsumer() { + // standbys are just those stateful tasks minus active tasks + final Map> consumerToPreviousStandbyTaskIds = new TreeMap<>(); + + for (final Map.Entry> entry: consumerToPreviousStatefulTaskIds.entrySet()) { + final List standbyTaskIds = new ArrayList<>(entry.getValue()); + standbyTaskIds.removeAll(consumerToPreviousActiveTaskIds.get(entry.getKey())); + consumerToPreviousStandbyTaskIds.put(entry.getKey(), standbyTaskIds); + } + + return consumerToPreviousStandbyTaskIds; + } + + public Map> assignedActiveByConsumer() { + return consumerToAssignedActiveTaskIds; + } + + public Map> assignedStandbyByConsumer() { + return consumerToAssignedStandbyTaskIds; + } + public void assignActive(final TaskId task) { assertNotAssigned(task); activeTasks.add(task); @@ -345,13 +355,17 @@ boolean hasMoreAvailableCapacityThan(final ClientState other) { } } + public String currentAssignment() { + return "[activeTasks: (" + activeTasks + + ") standbyTasks: (" + standbyTasks + ")]"; + } + @Override public String toString() { return "[activeTasks: (" + activeTasks + ") standbyTasks: (" + standbyTasks + ") prevActiveTasks: (" + prevActiveTasks + ") prevStandbyTasks: (" + prevStandbyTasks + - ") prevOwnedPartitionsByConsumerId: (" + ownedPartitions.keySet() + ") changelogOffsetTotalsByTask: (" + taskOffsetSums.entrySet() + ") taskLagTotals: (" + taskLagTotals.entrySet() + ") capacity: " + capacity + @@ -373,6 +387,7 @@ private void initializePrevActiveTasksFromOwnedPartitions(final Map()).add(task); } else { LOG.error("No task found for topic partition {}", tp); } From 3ae2df4f211dc247e3f71587ab9ec8dec3599e61 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Wed, 8 Jul 2020 18:26:55 -0700 Subject: [PATCH 2/6] add log4j for revoking active too --- .../internals/StreamsPartitionAssignor.java | 6 +++++- .../internals/assignment/ClientState.java | 14 +++++++++++++- 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 4396c902ce1d9..2afec9bdab357 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -961,10 +961,12 @@ private Map computeNewAssignment(final Map populateActiveTaskAndPartitionsLists(final List> consumerToPreviousActiveTaskIds; private final Map> consumerToAssignedActiveTaskIds; private final Map> consumerToAssignedStandbyTaskIds; + private final Map> consumerToRevokingActiveTaskIds; private final Map ownedPartitions; private final Map taskOffsetSums; // contains only stateful tasks we previously owned private final Map taskLagTotals; // contains lag for all stateful tasks in the app topology @@ -74,6 +75,7 @@ public ClientState() { consumerToPreviousActiveTaskIds = new TreeMap<>(); consumerToAssignedActiveTaskIds = new TreeMap<>(); consumerToAssignedStandbyTaskIds = new TreeMap<>(); + consumerToRevokingActiveTaskIds = new TreeMap<>(); ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = new TreeMap<>(); taskLagTotals = new TreeMap<>(); @@ -93,6 +95,7 @@ public ClientState(final Set previousActiveTasks, consumerToPreviousActiveTaskIds = new TreeMap<>(); consumerToAssignedActiveTaskIds = new TreeMap<>(); consumerToAssignedStandbyTaskIds = new TreeMap<>(); + consumerToRevokingActiveTaskIds = new TreeMap<>(); ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = emptyMap(); this.taskLagTotals = unmodifiableMap(taskLagTotals); @@ -135,6 +138,10 @@ public void assignStandbyToConsumer(final TaskId task, final String consumer) { consumerToAssignedStandbyTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); } + public void revokeActiveFromConsumer(final TaskId task, final String consumer) { + consumerToRevokingActiveTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + } + public Map> prevOwnedActiveByConsumer() { return consumerToPreviousActiveTaskIds; } @@ -145,7 +152,8 @@ public Map> prevOwnedStandbyByConsumer() { for (final Map.Entry> entry: consumerToPreviousStatefulTaskIds.entrySet()) { final List standbyTaskIds = new ArrayList<>(entry.getValue()); - standbyTaskIds.removeAll(consumerToPreviousActiveTaskIds.get(entry.getKey())); + if (consumerToPreviousActiveTaskIds.containsKey(entry.getKey())) + standbyTaskIds.removeAll(consumerToPreviousActiveTaskIds.get(entry.getKey())); consumerToPreviousStandbyTaskIds.put(entry.getKey(), standbyTaskIds); } @@ -156,6 +164,10 @@ public Map> assignedActiveByConsumer() { return consumerToAssignedActiveTaskIds; } + public Map> revokingActiveByConsumer() { + return consumerToRevokingActiveTaskIds; + } + public Map> assignedStandbyByConsumer() { return consumerToAssignedStandbyTaskIds; } From be0d18c701ffaca095e4404565aaa309b3200525 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Fri, 10 Jul 2020 10:46:22 -0700 Subject: [PATCH 3/6] github comments --- .../processor/internals/StreamsPartitionAssignor.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 2afec9bdab357..b6bda28a2f555 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -725,8 +725,8 @@ private boolean assignTasksToClients(final Set allSourceTopics, statefulTasks, assignmentConfigs); - log.info("Assigned tasks to clients as: {}{}.", Utils.NL, - clientStates.entrySet().stream() + log.info("Assigned tasks {} including stateful {} to clients as: \n{}.", + allTasks, statefulTasks, clientStates.entrySet().stream() .map(entry -> entry.getKey() + "=" + entry.getValue().currentAssignment()) .collect(Collectors.joining(Utils.NL))); @@ -934,7 +934,7 @@ private Map computeNewAssignment(final Map Date: Sun, 12 Jul 2020 16:34:52 -0700 Subject: [PATCH 4/6] add unit tests --- .../internals/StreamsPartitionAssignor.java | 12 +-- .../internals/assignment/ClientState.java | 43 ++++----- .../assignment/AssignmentTestUtils.java | 7 ++ .../internals/assignment/ClientStateTest.java | 87 +++++++++++++++---- 4 files changed, 104 insertions(+), 45 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 8581dbc1a1825..9b8b501e2aeae 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -974,11 +974,11 @@ private Map computeNewAssignment(final Set statefulT "\tassigned active {}\n" + "\trevoking active {}" + "\tassigned standby {}\n", clientId, - clientMetadata.state.prevOwnedActiveByConsumer(), + clientMetadata.state.prevOwnedActiveTasksByConsumer(), clientMetadata.state.prevOwnedStandbyByConsumer(), - clientMetadata.state.assignedActiveByConsumer(), - clientMetadata.state.revokingActiveByConsumer(), - clientMetadata.state.assignedStandbyByConsumer()); + clientMetadata.state.assignedActiveTasksByConsumer(), + clientMetadata.state.revokingActiveTasksByConsumer(), + clientMetadata.state.assignedStandbyTasksByConsumer()); } if (rebalanceRequired) { @@ -1158,7 +1158,7 @@ private Map> buildStandbyTaskMap(final String consum if (allStatefulTasks.contains(task)) { log.info("Adding removed stateful active task {} as a standby for {} before it is revoked in followup rebalance", task, consumer); - + // This has no effect on the assignment, as we'll never consult the ClientState again, but // it does perform a useful assertion that the it's legal to assign this task as a standby to this instance clientState.assignStandbyToConsumer(task, consumer); @@ -1288,7 +1288,7 @@ static Map> assignTasksToThreads(final Collection s private static SortedSet getPreviousTasksByLag(final ClientState state, final String consumer) { final SortedSet prevTasksByLag = new TreeSet<>(comparingLong(state::lagFor).thenComparing(TaskId::compareTo)); - prevTasksByLag.addAll(state.previousTasksForConsumer(consumer)); + prevTasksByLag.addAll(state.prevOwnedStatefulTasksByConsumer(consumer)); return prevTasksByLag; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java index 58f2574e28158..1c248e56bbd36 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java @@ -22,11 +22,9 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.ArrayList; import java.util.Collection; import java.util.Comparator; import java.util.HashSet; -import java.util.List; import java.util.Map; import java.util.Set; import java.util.stream.Collectors; @@ -51,10 +49,12 @@ public class ClientState { private final Set prevStandbyTasks; private final Map> consumerToPreviousStatefulTaskIds; - private final Map> consumerToPreviousActiveTaskIds; - private final Map> consumerToAssignedActiveTaskIds; - private final Map> consumerToAssignedStandbyTaskIds; - private final Map> consumerToRevokingActiveTaskIds; + // the following four maps are used only for logging purposes; + // TODO: we could consider merging them with other book-keeping maps + private final Map> consumerToPreviousActiveTaskIds; + private final Map> consumerToAssignedActiveTaskIds; + private final Map> consumerToAssignedStandbyTaskIds; + private final Map> consumerToRevokingActiveTaskIds; private final Map ownedPartitions; private final Map taskOffsetSums; // contains only stateful tasks we previously owned private final Map taskLagTotals; // contains lag for all stateful tasks in the app topology @@ -131,27 +131,27 @@ public void assignActiveTasks(final Collection tasks) { } public void assignActiveToConsumer(final TaskId task, final String consumer) { - consumerToAssignedActiveTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + consumerToAssignedActiveTaskIds.computeIfAbsent(consumer, k -> new HashSet<>()).add(task); } public void assignStandbyToConsumer(final TaskId task, final String consumer) { - consumerToAssignedStandbyTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + consumerToAssignedStandbyTaskIds.computeIfAbsent(consumer, k -> new HashSet<>()).add(task); } public void revokeActiveFromConsumer(final TaskId task, final String consumer) { - consumerToRevokingActiveTaskIds.getOrDefault(consumer, new ArrayList<>()).add(task); + consumerToRevokingActiveTaskIds.computeIfAbsent(consumer, k -> new HashSet<>()).add(task); } - public Map> prevOwnedActiveByConsumer() { + public Map> prevOwnedActiveTasksByConsumer() { return consumerToPreviousActiveTaskIds; } - public Map> prevOwnedStandbyByConsumer() { + public Map> prevOwnedStandbyByConsumer() { // standbys are just those stateful tasks minus active tasks - final Map> consumerToPreviousStandbyTaskIds = new TreeMap<>(); + final Map> consumerToPreviousStandbyTaskIds = new TreeMap<>(); for (final Map.Entry> entry: consumerToPreviousStatefulTaskIds.entrySet()) { - final List standbyTaskIds = new ArrayList<>(entry.getValue()); + final Set standbyTaskIds = new HashSet<>(entry.getValue()); if (consumerToPreviousActiveTaskIds.containsKey(entry.getKey())) standbyTaskIds.removeAll(consumerToPreviousActiveTaskIds.get(entry.getKey())); consumerToPreviousStandbyTaskIds.put(entry.getKey(), standbyTaskIds); @@ -160,15 +160,20 @@ public Map> prevOwnedStandbyByConsumer() { return consumerToPreviousStandbyTaskIds; } - public Map> assignedActiveByConsumer() { + // including both active and standby tasks + public Set prevOwnedStatefulTasksByConsumer(final String memberId) { + return consumerToPreviousStatefulTaskIds.get(memberId); + } + + public Map> assignedActiveTasksByConsumer() { return consumerToAssignedActiveTaskIds; } - public Map> revokingActiveByConsumer() { + public Map> revokingActiveTasksByConsumer() { return consumerToRevokingActiveTaskIds; } - public Map> assignedStandbyByConsumer() { + public Map> assignedStandbyTasksByConsumer() { return consumerToAssignedStandbyTaskIds; } @@ -338,10 +343,6 @@ public Set statelessActiveTasks() { return activeTasks.stream().filter(task -> !isStateful(task)).collect(Collectors.toSet()); } - public Set previousTasksForConsumer(final String memberId) { - return consumerToPreviousStatefulTaskIds.get(memberId); - } - boolean hasUnfulfilledQuota(final int tasksPerThread) { return activeTasks.size() < capacity * tasksPerThread; } @@ -399,7 +400,7 @@ private void initializePrevActiveTasksFromOwnedPartitions(final Map()).add(task); + consumerToPreviousActiveTaskIds.computeIfAbsent(partitionEntry.getValue(), k -> new HashSet<>()).add(task); } else { LOG.error("No task found for topic partition {}", tp); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/AssignmentTestUtils.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/AssignmentTestUtils.java index f0139ffdcfd82..d1da94f5b6249 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/AssignmentTestUtils.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/AssignmentTestUtils.java @@ -53,6 +53,13 @@ public final class AssignmentTestUtils { public static final UUID UUID_5 = uuidForInt(5); public static final UUID UUID_6 = uuidForInt(6); + public static final TopicPartition TP_0_0 = new TopicPartition("topic0", 0); + public static final TopicPartition TP_0_1 = new TopicPartition("topic0", 1); + public static final TopicPartition TP_0_2 = new TopicPartition("topic0", 2); + public static final TopicPartition TP_1_0 = new TopicPartition("topic1", 0); + public static final TopicPartition TP_1_1 = new TopicPartition("topic1", 1); + public static final TopicPartition TP_1_2 = new TopicPartition("topic1", 2); + public static final TaskId TASK_0_0 = new TaskId(0, 0); public static final TaskId TASK_0_1 = new TaskId(0, 1); public static final TaskId TASK_0_2 = new TaskId(0, 2); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java index d50c00eae2617..674695a6573a9 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java @@ -32,6 +32,12 @@ import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TASK_0_1; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TASK_0_2; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TASK_0_3; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_0_0; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_0_1; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_0_2; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_0; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_1; +import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_2; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.UUID_1; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.hasActiveTasks; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.hasStandbyTasks; @@ -309,35 +315,80 @@ public void shouldAddTasksWithLatestOffsetToPrevActiveTasks() { @Test public void shouldReturnPreviousStatefulTasksForConsumer() { - client.addPreviousTasksAndOffsetSums("c1", Collections.singletonMap(TASK_0_1, Task.LATEST_OFFSET)); + client.addPreviousTasksAndOffsetSums("c1", mkMap( + mkEntry(TASK_0_0, 100L), + mkEntry(TASK_0_1, Task.LATEST_OFFSET) + )); client.addPreviousTasksAndOffsetSums("c2", Collections.singletonMap(TASK_0_2, 0L)); client.addPreviousTasksAndOffsetSums("c3", Collections.emptyMap()); client.initializePrevTasks(Collections.emptyMap()); - client.computeTaskLags( - UUID_1, - mkMap( - mkEntry(TASK_0_1, 1_000L), - mkEntry(TASK_0_2, 1_000L) - ) - ); - assertThat(client.previousTasksForConsumer("c1"), equalTo(mkSet(TASK_0_1))); - assertThat(client.previousTasksForConsumer("c2"), equalTo(mkSet(TASK_0_2))); - assertTrue(client.previousTasksForConsumer("c3").isEmpty()); + assertThat(client.prevOwnedStatefulTasksByConsumer("c1"), equalTo(mkSet(TASK_0_0, TASK_0_1))); + assertThat(client.prevOwnedStatefulTasksByConsumer("c2"), equalTo(mkSet(TASK_0_2))); + assertTrue(client.prevOwnedStatefulTasksByConsumer("c3").isEmpty()); } @Test - public void shouldReturnPreviousTasksForConsumer() { + public void shouldReturnPreviousActiveStandbyTasksForConsumer() { + client.addOwnedPartitions(mkSet(TP_0_1, TP_1_1), "c1"); + client.addOwnedPartitions(mkSet(TP_0_2, TP_1_2), "c2"); + client.initializePrevTasks(mkMap( + mkEntry(TP_0_0, TASK_0_0), + mkEntry(TP_0_1, TASK_0_1), + mkEntry(TP_0_2, TASK_0_2), + mkEntry(TP_1_0, TASK_0_0), + mkEntry(TP_1_1, TASK_0_1), + mkEntry(TP_1_2, TASK_0_2)) + ); + client.addPreviousTasksAndOffsetSums("c1", mkMap( - mkEntry(TASK_0_1, 100L), - mkEntry(TASK_0_2, 0L), - mkEntry(TASK_0_3, Task.LATEST_OFFSET) - )); + mkEntry(TASK_0_1, Task.LATEST_OFFSET), + mkEntry(TASK_0_0, 10L))); + client.addPreviousTasksAndOffsetSums("c2", Collections.singletonMap(TASK_0_2, 0L)); - client.initializePrevTasks(Collections.emptyMap()); + assertThat(client.prevOwnedStatefulTasksByConsumer("c1"), equalTo(mkSet(TASK_0_1, TASK_0_0))); + assertThat(client.prevOwnedStatefulTasksByConsumer("c2"), equalTo(mkSet(TASK_0_2))); + assertThat(client.prevOwnedActiveTasksByConsumer(), equalTo( + mkMap( + mkEntry("c1", Collections.singleton(TASK_0_1)), + mkEntry("c2", Collections.singleton(TASK_0_2)) + )) + ); + assertThat(client.prevOwnedStandbyByConsumer(), equalTo( + mkMap( + mkEntry("c1", Collections.singleton(TASK_0_0)), + mkEntry("c2", Collections.emptySet()) + )) + ); + } - assertThat(client.previousTasksForConsumer("c1"), equalTo(mkSet(TASK_0_3, TASK_0_2, TASK_0_1))); + @Test + public void shouldReturnAssignedTasksForConsumer() { + client.assignActiveToConsumer(TASK_0_0, "c1"); + // calling it multiple tasks should be idempotent + client.assignActiveToConsumer(TASK_0_0, "c1"); + client.assignActiveToConsumer(TASK_0_1, "c1"); + client.assignActiveToConsumer(TASK_0_2, "c2"); + + client.assignStandbyToConsumer(TASK_0_2, "c1"); + client.assignStandbyToConsumer(TASK_0_0, "c2"); + // calling it multiple tasks should be idempotent + client.assignStandbyToConsumer(TASK_0_0, "c2"); + + client.revokeActiveFromConsumer(TASK_0_1, "c1"); + // calling it multiple tasks should be idempotent + client.revokeActiveFromConsumer(TASK_0_1, "c1"); + + assertThat(client.assignedActiveTasksByConsumer(), equalTo(mkMap( + mkEntry("c1", mkSet(TASK_0_0, TASK_0_1)), + mkEntry("c2", mkSet(TASK_0_2)) + ))); + assertThat(client.assignedStandbyTasksByConsumer(), equalTo(mkMap( + mkEntry("c1", mkSet(TASK_0_2)), + mkEntry("c2", mkSet(TASK_0_0)) + ))); + assertThat(client.revokingActiveTasksByConsumer(), equalTo(Collections.singletonMap("c1", mkSet(TASK_0_1)))); } @Test From d7e54ddb17ec093f619ed5f89d97a23cdfb36b3e Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Tue, 14 Jul 2020 14:26:09 -0700 Subject: [PATCH 5/6] checkstyle unit tests --- .../streams/processor/internals/assignment/ClientStateTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java index 674695a6573a9..952a9904c6aa1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/ClientStateTest.java @@ -38,7 +38,6 @@ import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_0; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_1; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.TP_1_2; -import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.UUID_1; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.hasActiveTasks; import static org.apache.kafka.streams.processor.internals.assignment.AssignmentTestUtils.hasStandbyTasks; import static org.apache.kafka.streams.processor.internals.assignment.SubscriptionInfo.UNKNOWN_OFFSET_SUM; From 581980bfe301fa6f09a6fa8f7375d78955a09bbf Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 16 Jul 2020 18:29:01 -0700 Subject: [PATCH 6/6] github comments --- .../internals/StreamsPartitionAssignor.java | 21 ++++++----- .../internals/assignment/ClientState.java | 37 ++++++------------- 2 files changed, 22 insertions(+), 36 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 9b8b501e2aeae..aed7d968954ae 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -969,16 +969,17 @@ private Map computeNewAssignment(final Set statefulT } log.info("Client {} per-consumer assignment:\n" + - "\tprev owned active {}\n" + - "\tprev owned standby {}\n" + - "\tassigned active {}\n" + - "\trevoking active {}" + - "\tassigned standby {}\n", clientId, - clientMetadata.state.prevOwnedActiveTasksByConsumer(), - clientMetadata.state.prevOwnedStandbyByConsumer(), - clientMetadata.state.assignedActiveTasksByConsumer(), - clientMetadata.state.revokingActiveTasksByConsumer(), - clientMetadata.state.assignedStandbyTasksByConsumer()); + "\tprev owned active {}\n" + + "\tprev owned standby {}\n" + + "\tassigned active {}\n" + + "\trevoking active {}" + + "\tassigned standby {}\n", + clientId, + clientMetadata.state.prevOwnedActiveTasksByConsumer(), + clientMetadata.state.prevOwnedStandbyByConsumer(), + clientMetadata.state.assignedActiveTasksByConsumer(), + clientMetadata.state.revokingActiveTasksByConsumer(), + clientMetadata.state.assignedStandbyTasksByConsumer()); } if (rebalanceRequired) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java index 1c248e56bbd36..9be0b24c170c2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/ClientState.java @@ -43,22 +43,23 @@ public class ClientState { private static final Logger LOG = LoggerFactory.getLogger(ClientState.class); public static final Comparator TOPIC_PARTITION_COMPARATOR = comparing(TopicPartition::topic).thenComparing(TopicPartition::partition); - private final Set activeTasks; - private final Set standbyTasks; + private final Set activeTasks = new TreeSet<>(); + private final Set standbyTasks = new TreeSet<>(); private final Set prevActiveTasks; private final Set prevStandbyTasks; - private final Map> consumerToPreviousStatefulTaskIds; - // the following four maps are used only for logging purposes; - // TODO: we could consider merging them with other book-keeping maps - private final Map> consumerToPreviousActiveTaskIds; - private final Map> consumerToAssignedActiveTaskIds; - private final Map> consumerToAssignedStandbyTaskIds; - private final Map> consumerToRevokingActiveTaskIds; - private final Map ownedPartitions; private final Map taskOffsetSums; // contains only stateful tasks we previously owned private final Map taskLagTotals; // contains lag for all stateful tasks in the app topology + private final Map ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); + private final Map> consumerToPreviousStatefulTaskIds = new TreeMap<>(); + // the following four maps are used only for logging purposes; + // TODO KAFKA-10283: we could consider merging them with other book-keeping maps at client-levels + // so that they would not be inconsistent + private final Map> consumerToPreviousActiveTaskIds = new TreeMap<>(); + private final Map> consumerToAssignedActiveTaskIds = new TreeMap<>(); + private final Map> consumerToAssignedStandbyTaskIds = new TreeMap<>(); + private final Map> consumerToRevokingActiveTaskIds = new TreeMap<>(); private int capacity; @@ -67,16 +68,8 @@ public ClientState() { } ClientState(final int capacity) { - activeTasks = new TreeSet<>(); - standbyTasks = new TreeSet<>(); prevActiveTasks = new TreeSet<>(); prevStandbyTasks = new TreeSet<>(); - consumerToPreviousStatefulTaskIds = new TreeMap<>(); - consumerToPreviousActiveTaskIds = new TreeMap<>(); - consumerToAssignedActiveTaskIds = new TreeMap<>(); - consumerToAssignedStandbyTaskIds = new TreeMap<>(); - consumerToRevokingActiveTaskIds = new TreeMap<>(); - ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = new TreeMap<>(); taskLagTotals = new TreeMap<>(); this.capacity = capacity; @@ -87,16 +80,8 @@ public ClientState(final Set previousActiveTasks, final Set previousStandbyTasks, final Map taskLagTotals, final int capacity) { - activeTasks = new TreeSet<>(); - standbyTasks = new TreeSet<>(); prevActiveTasks = unmodifiableSet(new TreeSet<>(previousActiveTasks)); prevStandbyTasks = unmodifiableSet(new TreeSet<>(previousStandbyTasks)); - consumerToPreviousStatefulTaskIds = new TreeMap<>(); - consumerToPreviousActiveTaskIds = new TreeMap<>(); - consumerToAssignedActiveTaskIds = new TreeMap<>(); - consumerToAssignedStandbyTaskIds = new TreeMap<>(); - consumerToRevokingActiveTaskIds = new TreeMap<>(); - ownedPartitions = new TreeMap<>(TOPIC_PARTITION_COMPARATOR); taskOffsetSums = emptyMap(); this.taskLagTotals = unmodifiableMap(taskLagTotals); this.capacity = capacity;