From e225e23ec498621924504659ca3087893da81fac Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Sat, 15 Dec 2018 04:05:28 -0800 Subject: [PATCH 1/2] MINOR: Include additional detail in fetch error message --- .../org/apache/kafka/common/requests/FetchRequest.java | 2 +- core/src/main/scala/kafka/api/Request.scala | 10 ++++++++++ core/src/main/scala/kafka/server/ReplicaManager.scala | 8 ++++++-- 3 files changed, 17 insertions(+), 3 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/FetchRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/FetchRequest.java index 8d94bfd2134ea..b3443a1f70447 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/FetchRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/FetchRequest.java @@ -233,7 +233,7 @@ public PartitionData(long fetchOffset, long logStartOffset, int maxBytes, Option @Override public String toString() { - return "(offset=" + fetchOffset + + return "(fetchOffset=" + fetchOffset + ", logStartOffset=" + logStartOffset + ", maxBytes=" + maxBytes + ", currentLeaderEpoch=" + currentLeaderEpoch + diff --git a/core/src/main/scala/kafka/api/Request.scala b/core/src/main/scala/kafka/api/Request.scala index b6ec2735e9df6..bfa4af2c44d5f 100644 --- a/core/src/main/scala/kafka/api/Request.scala +++ b/core/src/main/scala/kafka/api/Request.scala @@ -24,4 +24,14 @@ object Request { // Broker ids are non-negative int. def isValidBrokerId(brokerId: Int): Boolean = brokerId >= 0 + + def describeReplicaId(replicaId: Int): String = { + replicaId match { + case OrdinaryConsumerId => "consumer" + case DebuggingConsumerId => "debug consumer" + case FutureLocalReplicaId => "future local replica" + case id if isValidBrokerId(id) => s"replica [$replicaId]" + case id => s"invalid replica [$id]" + } + } } diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 2f6430223b855..4cc3febc29775 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -886,13 +886,13 @@ class ReplicaManager(val config: KafkaConfig, brokerTopicStats.topicStats(tp.topic).totalFetchRequestRate.mark() brokerTopicStats.allTopicsStats.totalFetchRequestRate.mark() + val adjustedMaxBytes = math.min(fetchInfo.maxBytes, limitBytes) try { trace(s"Fetching log segment for partition $tp, offset $offset, partition fetch size $partitionFetchSize, " + s"remaining response limit $limitBytes" + (if (minOneMessage) s", ignoring response/partition size limits" else "")) val partition = getPartitionOrException(tp, expectLeader = fetchOnlyFromLeader) - val adjustedMaxBytes = math.min(fetchInfo.maxBytes, limitBytes) val fetchTimeMs = time.milliseconds // Try the read first, this tells us whether we need all of adjustedFetchSize for this partition @@ -946,7 +946,11 @@ class ReplicaManager(val config: KafkaConfig, case e: Throwable => brokerTopicStats.topicStats(tp.topic).failedFetchRequestRate.mark() brokerTopicStats.allTopicsStats.failedFetchRequestRate.mark() - error(s"Error processing fetch operation on partition $tp, offset $offset", e) + + val fetchSource = Request.describeReplicaId(replicaId) + error(s"Error processing fetch with max size $adjustedMaxBytes from $fetchSource " + + s"on partition $tp: $fetchInfo", e) + LogReadResult(info = FetchDataInfo(LogOffsetMetadata.UnknownOffsetMetadata, MemoryRecords.EMPTY), highWatermark = -1L, leaderLogStartOffset = -1L, From aeec58e8aa05579b68027999fb2003a9a67ba699 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 17 Dec 2018 09:49:55 -0800 Subject: [PATCH 2/2] Use match variable in describe for consistency --- core/src/main/scala/kafka/api/Request.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/api/Request.scala b/core/src/main/scala/kafka/api/Request.scala index bfa4af2c44d5f..653b5f653ac52 100644 --- a/core/src/main/scala/kafka/api/Request.scala +++ b/core/src/main/scala/kafka/api/Request.scala @@ -30,7 +30,7 @@ object Request { case OrdinaryConsumerId => "consumer" case DebuggingConsumerId => "debug consumer" case FutureLocalReplicaId => "future local replica" - case id if isValidBrokerId(id) => s"replica [$replicaId]" + case id if isValidBrokerId(id) => s"replica [$id]" case id => s"invalid replica [$id]" } }