Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -301,8 +301,8 @@ public Map<String, Assignment> assign(final Cluster metadata,
throw new IllegalStateException("Unknown metadata version: " + usedVersion
+ "; latest supported version: " + SubscriptionInfo.LATEST_SUPPORTED_VERSION);
}
if (info.version() < minUserMetadataVersion) {
minUserMetadataVersion = info.version();
if (usedVersion < minUserMetadataVersion) {
minUserMetadataVersion = usedVersion;
}

// create the new client metadata if necessary
Expand Down Expand Up @@ -615,7 +615,7 @@ public void onAssignment(final Assignment assignment) {
switch (usedVersion) {
case 1:
processVersionOneAssignment(info, partitions, activeTasks);
partitionsByHost = new HashMap<>();
partitionsByHost = Collections.emptyMap();
break;
case 2:
processVersionTwoAssignment(info, partitions, activeTasks, topicToPartitionInfo);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,10 +114,10 @@ public ByteBuffer encode() {

private void encodeVersionOne(final DataOutputStream out) throws IOException {
out.writeInt(1); // version
encodeVersionOneData(out);
encodeActiveAndStandbyTaskAssignment(out);
}

private void encodeVersionOneData(final DataOutputStream out) throws IOException {
private void encodeActiveAndStandbyTaskAssignment(final DataOutputStream out) throws IOException {
// encode active tasks
out.writeInt(activeTasks.size());
for (final TaskId id : activeTasks) {
Expand All @@ -137,11 +137,11 @@ private void encodeVersionOneData(final DataOutputStream out) throws IOException

private void encodeVersionTwo(final DataOutputStream out) throws IOException {
out.writeInt(2); // version
encodeVersionOneData(out);
encodeVersionTwoData(out);
encodeActiveAndStandbyTaskAssignment(out);
encodePartitionsByHost(out);
}

private void encodeVersionTwoData(final DataOutputStream out) throws IOException {
private void encodePartitionsByHost(final DataOutputStream out) throws IOException {
// encode partitions by host
out.writeInt(partitionsByHost.size());
for (final Map.Entry<HostInfo, Set<TopicPartition>> entry : partitionsByHost.entrySet()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,9 +93,7 @@ public ByteBuffer encode() {
buf = encodeVersionOne();
break;
case 2:
byte[] endPointBytes = null;
endPointBytes = prepareUserEndPoint();
buf = encodeVersionTwo(endPointBytes);
buf = encodeVersionTwo(prepareUserEndPoint());
break;
default:
throw new IllegalStateException("Unknown metadata version: " + usedVersion
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1095,6 +1095,31 @@ public void shouldThrowKafkaExceptionIfStreamThreadConfigIsNotThreadDataProvider
partitionAssignor.configure(config);
}

public void shouldReturnLowestAssignmentVersionForDifferentSubscriptionVersions() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Add @test annotation.

final Map<String, PartitionAssignor.Subscription> subscriptions = new HashMap<>();
final Set<TaskId> emptyTasks = Collections.emptySet();
subscriptions.put(
"consumer1",
new PartitionAssignor.Subscription(
Collections.singletonList("topic1"),
new SubscriptionInfo(1, UUID.randomUUID(), emptyTasks, emptyTasks, null).encode()
)
);
subscriptions.put(
"consumer2",
new PartitionAssignor.Subscription(
Collections.singletonList("topic1"),
new SubscriptionInfo(2, UUID.randomUUID(), emptyTasks, emptyTasks, null).encode()
)
);

final Map<String, PartitionAssignor.Assignment> assignment = partitionAssignor.assign(metadata, subscriptions);

assertThat(assignment.size(), equalTo(2));
assertThat(AssignmentInfo.decode(assignment.get("consumer1").userData()).version(), equalTo(1));
assertThat(AssignmentInfo.decode(assignment.get("consumer2").userData()).version(), equalTo(1));
}

private PartitionAssignor.Assignment createAssignment(final Map<HostInfo, Set<TopicPartition>> firstHostState) {
final AssignmentInfo info = new AssignmentInfo(Collections.<TaskId>emptyList(),
Collections.<TaskId, Set<TopicPartition>>emptyMap(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ public void shouldDecodePreviousVersion() throws IOException {
final AssignmentInfo decoded = AssignmentInfo.decode(encodeV1(oldVersion));
assertEquals(oldVersion.activeTasks(), decoded.activeTasks());
assertEquals(oldVersion.standbyTasks(), decoded.standbyTasks());
assertNull(decoded.partitionsByHost()); // should be empty as wasn't in V1
assertNull(decoded.partitionsByHost()); // should be null as wasn't in V1
assertEquals(1, decoded.version());
}

Expand Down