From 5b54660982848352815f3e1fdd682b306d5fd71e Mon Sep 17 00:00:00 2001 From: m1a2st Date: Tue, 26 Aug 2025 18:18:47 +0800 Subject: [PATCH 1/2] update the nonEmpty check --- core/src/main/scala/kafka/server/ReplicaManager.scala | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 496e50208db89..498721607d620 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -788,7 +788,9 @@ class ReplicaManager(val config: KafkaConfig, hasCustomErrorMessage = customException.isDefined ) } - val entriesWithoutErrorsPerPartition = entriesPerPartition.filter { case (key, _) => !errorResults.contains(key) } + val entriesWithoutErrorsPerPartition = + if (errorResults.nonEmpty) entriesPerPartition.filter { case (key, _) => !errorResults.contains(key) } + else entriesPerPartition val preAppendPartitionResponses = buildProducePartitionStatus(errorResults).map { case (k, status) => k -> status.responseStatus } From 7a17476043a218a264433d36e1071bae78d321a0 Mon Sep 17 00:00:00 2001 From: m1a2st Date: Tue, 26 Aug 2025 20:18:49 +0800 Subject: [PATCH 2/2] addressed by comment --- core/src/main/scala/kafka/server/ReplicaManager.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 498721607d620..202590ca6f4d9 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -788,6 +788,8 @@ class ReplicaManager(val config: KafkaConfig, hasCustomErrorMessage = customException.isDefined ) } + // In non-transaction paths, errorResults is typically empty, so we can + // directly use entriesPerPartition instead of creating a new filtered collection val entriesWithoutErrorsPerPartition = if (errorResults.nonEmpty) entriesPerPartition.filter { case (key, _) => !errorResults.contains(key) } else entriesPerPartition