From 0a7e2a6e820fae2632a98705494cd1a374d46f9f Mon Sep 17 00:00:00 2001 From: David Mao Date: Thu, 6 Feb 2020 15:43:20 -0800 Subject: [PATCH 1/2] KAFKA-9507 AdminClient should check for missing committed offsets --- .../org/apache/kafka/clients/admin/KafkaAdminClient.java | 6 +++++- .../kafka/clients/admin/ListConsumerGroupOffsetsResult.java | 1 + .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 6 +++++- 3 files changed, 11 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 6b8e962f97955..d3f9151528e72 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3100,7 +3100,11 @@ void handleResponse(AbstractResponse abstractResponse) { final Long offset = partitionData.offset; final String metadata = partitionData.metadata; final Optional leaderEpoch = partitionData.leaderEpoch; - groupOffsetsListing.put(topicPartition, new OffsetAndMetadata(offset, leaderEpoch, metadata)); + if (offset < 0) { + groupOffsetsListing.put(topicPartition, null); + } else { + groupOffsetsListing.put(topicPartition, new OffsetAndMetadata(offset, leaderEpoch, metadata)); + } } else { log.warn("Skipping return offset for {} due to error {}.", topicPartition, error); } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java index ea5193402bd91..27966c8c3743c 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java @@ -40,6 +40,7 @@ public class ListConsumerGroupOffsetsResult { /** * Return a future which yields a map of topic partitions to OffsetAndMetadata objects. + * If the partition does not have a committed offset, the corresponding value will be null */ public KafkaFuture> partitionsToOffsetAndMetadata() { return future; diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index cba9b48d3829d..4e6ed47089d66 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1494,6 +1494,7 @@ public void testDescribeConsumerGroupOffsets() throws Exception { TopicPartition myTopicPartition0 = new TopicPartition("my_topic", 0); TopicPartition myTopicPartition1 = new TopicPartition("my_topic", 1); TopicPartition myTopicPartition2 = new TopicPartition("my_topic", 2); + TopicPartition myTopicPartition3 = new TopicPartition("my_topic", 3); final Map responseData = new HashMap<>(); responseData.put(myTopicPartition0, new OffsetFetchResponse.PartitionData(10, @@ -1502,15 +1503,18 @@ public void testDescribeConsumerGroupOffsets() throws Exception { Optional.empty(), "", Errors.NONE)); responseData.put(myTopicPartition2, new OffsetFetchResponse.PartitionData(20, Optional.empty(), "", Errors.NONE)); + responseData.put(myTopicPartition3, new OffsetFetchResponse.PartitionData(-1, + Optional.empty(), "", Errors.NONE)); env.kafkaClient().prepareResponse(new OffsetFetchResponse(Errors.NONE, responseData)); final ListConsumerGroupOffsetsResult result = env.adminClient().listConsumerGroupOffsets("group-0"); final Map partitionToOffsetAndMetadata = result.partitionsToOffsetAndMetadata().get(); - assertEquals(3, partitionToOffsetAndMetadata.size()); + assertEquals(4, partitionToOffsetAndMetadata.size()); assertEquals(10, partitionToOffsetAndMetadata.get(myTopicPartition0).offset()); assertEquals(0, partitionToOffsetAndMetadata.get(myTopicPartition1).offset()); assertEquals(20, partitionToOffsetAndMetadata.get(myTopicPartition2).offset()); + assertNull(partitionToOffsetAndMetadata.get(myTopicPartition3)); } } From 9847c4aec837a89172fc9a09f7fefc7de6af9e51 Mon Sep 17 00:00:00 2001 From: David Mao Date: Thu, 6 Feb 2020 16:11:02 -0800 Subject: [PATCH 2/2] addresses PR comments --- .../java/org/apache/kafka/clients/admin/KafkaAdminClient.java | 1 + .../kafka/clients/admin/ListConsumerGroupOffsetsResult.java | 2 +- .../org/apache/kafka/clients/admin/KafkaAdminClientTest.java | 3 ++- 3 files changed, 4 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index d3f9151528e72..52332f8820607 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3100,6 +3100,7 @@ void handleResponse(AbstractResponse abstractResponse) { final Long offset = partitionData.offset; final String metadata = partitionData.metadata; final Optional leaderEpoch = partitionData.leaderEpoch; + // Negative offset indicates that the group has no committed offset for this partition if (offset < 0) { groupOffsetsListing.put(topicPartition, null); } else { diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java index 27966c8c3743c..48f4531418110 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult.java @@ -40,7 +40,7 @@ public class ListConsumerGroupOffsetsResult { /** * Return a future which yields a map of topic partitions to OffsetAndMetadata objects. - * If the partition does not have a committed offset, the corresponding value will be null + * If the group does not have a committed offset for this partition, the corresponding value in the returned map will be null. */ public KafkaFuture> partitionsToOffsetAndMetadata() { return future; diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 4e6ed47089d66..ed3e8a6781ac7 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1503,7 +1503,7 @@ public void testDescribeConsumerGroupOffsets() throws Exception { Optional.empty(), "", Errors.NONE)); responseData.put(myTopicPartition2, new OffsetFetchResponse.PartitionData(20, Optional.empty(), "", Errors.NONE)); - responseData.put(myTopicPartition3, new OffsetFetchResponse.PartitionData(-1, + responseData.put(myTopicPartition3, new OffsetFetchResponse.PartitionData(OffsetFetchResponse.INVALID_OFFSET, Optional.empty(), "", Errors.NONE)); env.kafkaClient().prepareResponse(new OffsetFetchResponse(Errors.NONE, responseData)); @@ -1514,6 +1514,7 @@ public void testDescribeConsumerGroupOffsets() throws Exception { assertEquals(10, partitionToOffsetAndMetadata.get(myTopicPartition0).offset()); assertEquals(0, partitionToOffsetAndMetadata.get(myTopicPartition1).offset()); assertEquals(20, partitionToOffsetAndMetadata.get(myTopicPartition2).offset()); + assertTrue(partitionToOffsetAndMetadata.containsKey(myTopicPartition3)); assertNull(partitionToOffsetAndMetadata.get(myTopicPartition3)); } }