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..653b5f653ac52 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 [$id]" + 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,