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 3d02bfdffb29a..7e55d468d3b66 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 @@ -512,8 +512,8 @@ private void handleFetchResponse(ClientResponse resp, FetchRequest request) { this.subscriptions.fetched(tp, record.offset() + 1); this.records.add(new PartitionRecords<>(fetchOffset, tp, parsed)); this.sensors.recordsFetchLag.record(partition.highWatermark - record.offset()); - } else if (buffer.capacity() >= this.fetchSize) { - // we did not read a single message from a max fetchable buffer + } else if (buffer.limit() > 0) { + // we did not read a single message from a non-empty buffer // because that message's size is larger than fetch size, in this case // record this exception this.recordTooLargePartitions.put(tp, fetchOffset);