diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java index 4f5aa3e4726e4..4d272f34ca4af 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java @@ -275,10 +275,18 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H } private boolean isMuted(TopicPartition tp, long now) { - boolean result = muted.containsKey(tp) && muted.get(tp) > now; - if (!result) + // Take care to avoid unnecessary map look-ups because this method is a hotspot if producing to a + // large number of partitions + Long throttleUntilTime = muted.get(tp); + if (throttleUntilTime == null) + return false; + + if (now >= throttleUntilTime) { muted.remove(tp); - return result; + return false; + } + + return true; } public void resetNextBatchExpiryTime() {