From 73d19b9dd9c5de6ac48710acb32b8d77f9dc0df7 Mon Sep 17 00:00:00 2001 From: Lucas Bradstreet Date: Sat, 10 Aug 2019 22:13:42 -0700 Subject: [PATCH 1/6] MINOR: avoid double HashMap lookup in Producer RecordAccumulator --- .../kafka/clients/producer/internals/RecordAccumulator.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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..efb3ae1cac30f 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,7 +275,8 @@ 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; + Long time = muted.get(tp); + boolean result = time != null && time > now; if (!result) muted.remove(tp); return result; From 51dfeb989c082686c199ff429e10ae7a3bed2f3c Mon Sep 17 00:00:00 2001 From: Lucas Bradstreet Date: Sat, 10 Aug 2019 22:36:21 -0700 Subject: [PATCH 2/6] Name variable better --- .../kafka/clients/producer/internals/RecordAccumulator.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 efb3ae1cac30f..7b8fcf4845e59 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,8 +275,8 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H } private boolean isMuted(TopicPartition tp, long now) { - Long time = muted.get(tp); - boolean result = time != null && time > now; + Long throttleUntilTime = muted.get(tp); + boolean result = throttleUntilTime != null && throttleUntilTime > now; if (!result) muted.remove(tp); return result; From 5fff9e26a08fde0766d133407a6770669bb42e5e Mon Sep 17 00:00:00 2001 From: Lucas Bradstreet Date: Sat, 10 Aug 2019 22:58:58 -0700 Subject: [PATCH 3/6] Improve readability of the logic --- .../kafka/clients/producer/internals/RecordAccumulator.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 7b8fcf4845e59..05ec711cd8b41 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 @@ -276,10 +276,10 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H private boolean isMuted(TopicPartition tp, long now) { Long throttleUntilTime = muted.get(tp); - boolean result = throttleUntilTime != null && throttleUntilTime > now; - if (!result) + boolean unmute = throttleUntilTime != null && now >= throttleUntilTime; + if (unmute) muted.remove(tp); - return result; + return unmute; } public void resetNextBatchExpiryTime() { From 5741f7765eff6f17d787afaab69f26c088731137 Mon Sep 17 00:00:00 2001 From: Lucas Bradstreet Date: Sat, 10 Aug 2019 23:05:16 -0700 Subject: [PATCH 4/6] Rework unmute/isMuted logic --- .../clients/producer/internals/RecordAccumulator.java | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) 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 05ec711cd8b41..05587599c13ff 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 @@ -276,10 +276,15 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H private boolean isMuted(TopicPartition tp, long now) { Long throttleUntilTime = muted.get(tp); - boolean unmute = throttleUntilTime != null && now >= throttleUntilTime; - if (unmute) + if (throttleUntilTime == null) + return false; + + if (now >= throttleUntilTime) { muted.remove(tp); - return unmute; + return false; + } + + return true; } public void resetNextBatchExpiryTime() { From 64a73beb8d09d6b179fdde0bd2959e6a480fbfaf Mon Sep 17 00:00:00 2001 From: Ismael Juma Date: Sun, 11 Aug 2019 00:00:26 -0700 Subject: [PATCH 5/6] Add comment to reduce chances of regression in the future --- .../kafka/clients/producer/internals/RecordAccumulator.java | 1 + 1 file changed, 1 insertion(+) 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 05587599c13ff..e293cab94eeb5 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,6 +275,7 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H } private boolean isMuted(TopicPartition tp, long now) { + // Take care to avoid unnecessary map look-ups because this method is a hotspot with a large number of partitions Long throttleUntilTime = muted.get(tp); if (throttleUntilTime == null) return false; From 75d4320e502a9dd0939447d01e0c1730cbe7727b Mon Sep 17 00:00:00 2001 From: Ismael Juma Date: Sun, 11 Aug 2019 00:01:42 -0700 Subject: [PATCH 6/6] Tweak wording --- .../kafka/clients/producer/internals/RecordAccumulator.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 e293cab94eeb5..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,7 +275,8 @@ private RecordAppendResult tryAppend(long timestamp, byte[] key, byte[] value, H } private boolean isMuted(TopicPartition tp, long now) { - // Take care to avoid unnecessary map look-ups because this method is a hotspot with a large number of partitions + // 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;