From 386709b85074ba2e41282335489e1843c3a61605 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Mon, 11 Mar 2019 11:44:59 -0700 Subject: [PATCH 1/2] KAFKA-7565: Better messaging for invalid fetch response Users have reported that when consumer poll wake up is used, it is possible to receive fetch responses that don't match the copied topic partitions collection for the session when the fetch request was created. This commit improves the error handling here by throwing an IllegalStateException instead of a NullPointerException. And by generating a message for the exception that includes a bit of more information. --- .../clients/consumer/internals/Fetcher.java | 20 ++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java index 9009ffe092a9f..6435fb1275736 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java @@ -241,13 +241,19 @@ public void onSuccess(ClientResponse resp) { for (Map.Entry> entry : response.responseData().entrySet()) { TopicPartition partition = entry.getKey(); - long fetchOffset = data.sessionPartitions().get(partition).fetchOffset; - FetchResponse.PartitionData fetchData = entry.getValue(); - - log.debug("Fetch {} at offset {} for partition {} returned fetch data {}", - isolationLevel, fetchOffset, partition, fetchData); - completedFetches.add(new CompletedFetch(partition, fetchOffset, fetchData, metricAggregator, - resp.requestHeader().apiVersion())); + FetchRequest.PartitionData requestData = data.sessionPartitions().get(partition); + if (requestData == null) { + // Received fetch response for missing session partition + throw new IllegalStateException("Received response for partition " + partition + " which is not part of the session"); + } else { + long fetchOffset = requestData.fetchOffset; + FetchResponse.PartitionData fetchData = entry.getValue(); + + log.debug("Fetch {} at offset {} for partition {} returned fetch data {}", + isolationLevel, fetchOffset, partition, fetchData); + completedFetches.add(new CompletedFetch(partition, fetchOffset, fetchData, metricAggregator, + resp.requestHeader().apiVersion())); + } } sensors.fetchLatency.record(resp.requestLatencyMs()); From 4f78e775c00a00235d853b09de0146797dc2fa8b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Armando=20Garc=C3=ADa=20Sancio?= Date: Mon, 11 Mar 2019 16:33:10 -0700 Subject: [PATCH 2/2] Create a better message for the exception --- .../kafka/clients/consumer/internals/Fetcher.java | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java index 6435fb1275736..30d209f2693ae 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Fetcher.java @@ -68,6 +68,7 @@ import org.apache.kafka.common.utils.Timer; import org.apache.kafka.common.utils.Utils; import org.slf4j.Logger; +import org.slf4j.helpers.MessageFormatter; import java.io.Closeable; import java.nio.ByteBuffer; @@ -243,8 +244,19 @@ public void onSuccess(ClientResponse resp) { TopicPartition partition = entry.getKey(); FetchRequest.PartitionData requestData = data.sessionPartitions().get(partition); if (requestData == null) { + String message; + if (data.metadata().isFull()) { + message = MessageFormatter.arrayFormat( + "Response for missing full request partition: partition={}; metadata={}", + new Object[]{partition, data.metadata()}).getMessage(); + } else { + message = MessageFormatter.arrayFormat( + "Response for missing session request partition: partition={}; metadata={}; toSend={}; toForget={}", + new Object[]{partition, data.metadata(), data.toSend(), data.toForget()}).getMessage(); + } + // Received fetch response for missing session partition - throw new IllegalStateException("Received response for partition " + partition + " which is not part of the session"); + throw new IllegalStateException(message); } else { long fetchOffset = requestData.fetchOffset; FetchResponse.PartitionData fetchData = entry.getValue();