This is an automated email from the ASF dual-hosted git repository.

chia7712 pushed a commit to branch 4.2
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/4.2 by this push:
     new 1859197bbf7 KAFKA-20746 NPE in MetadataCache#toCluster (#22800)
1859197bbf7 is described below

commit 1859197bbf7becc0692c9d546ec7b79ed49dbbdc
Author: Kuan-Po Tseng <[email protected]>
AuthorDate: Fri Jul 10 21:40:51 2026 +0800

    KAFKA-20746 NPE in MetadataCache#toCluster (#22800)
    
    Backport of https://github.com/apache/kafka/pull/22723 to 4.2.
    
    MetadataCache#toCluster builds a Cluster view for quota callbacks.
    Previously,  it filtered fenced brokers out of the broker-to-node map,
    but partition metadata  can still reference fenced brokers as replicas.
    When such replicas were converted  to PartitionInfo, the missing broker
    entry could cause a NullPointerException.
    
    Include all registered brokers in the broker-to-node map so fenced
    replicas can  be represented as offline replicas in the generated
    Cluster.
    
    Reviewers: Luke Chen <[email protected]>, Chia-Ping Tsai
     <[email protected]>
---
 .../unit/kafka/server/MetadataCacheTest.scala      | 53 ++++++++++++++++++++++
 .../org/apache/kafka/metadata/MetadataCache.java   | 10 ++--
 2 files changed, 58 insertions(+), 5 deletions(-)

diff --git a/core/src/test/scala/unit/kafka/server/MetadataCacheTest.scala 
b/core/src/test/scala/unit/kafka/server/MetadataCacheTest.scala
index d252195789a..1d6ffda5c54 100644
--- a/core/src/test/scala/unit/kafka/server/MetadataCacheTest.scala
+++ b/core/src/test/scala/unit/kafka/server/MetadataCacheTest.scala
@@ -929,6 +929,59 @@ class MetadataCacheTest {
     ), offlinePartitions(brokers, partitions))
   }
 
+  @Test
+  def testToClusterIncludesFencedReplicasAsOffline(): Unit = {
+    val delta = new MetadataDelta.Builder().build()
+    val broker0RegisterRecord = new RegisterBrokerRecord()
+      .setBrokerId(0)
+      .setFenced(false)
+      .setEndPoints(new BrokerEndpointCollection(Seq(
+        new BrokerEndpoint()
+          .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id)
+          .setPort(9092)
+          .setName("PLAINTEXT")
+          .setHost("broker-0")
+      ).iterator.asJava))
+    val broker1RegisterRecord = new RegisterBrokerRecord()
+      .setBrokerId(1)
+      .setFenced(true)
+      .setEndPoints(new BrokerEndpointCollection(Seq(
+        new BrokerEndpoint()
+          .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id)
+          .setPort(9092)
+          .setName("PLAINTEXT")
+          .setHost("broker-1")
+      ).iterator.asJava))
+
+    delta.replay(broker0RegisterRecord)
+    delta.replay(broker1RegisterRecord)
+
+    val topicName = "foo"
+    val topicId = Uuid.randomUuid()
+    val partitionId = 0
+    val broker0 = broker0RegisterRecord.brokerId()
+    val broker1 = broker1RegisterRecord.brokerId()
+    delta.replay(new TopicRecord().setTopicId(topicId).setName(topicName))
+    // Broker 1 is not in the ISR, but it remains in the replica set.
+    // This used to trigger a NullPointerException when fenced brokers were
+    // omitted from the broker map used to build replica metadata.
+    delta.replay(new PartitionRecord()
+      .setTopicId(topicId)
+      .setPartitionId(partitionId)
+      .setReplicas(asList[Integer](broker0, broker1))
+      .setLeader(broker0)
+      .setIsr(asList[Integer](broker0)))
+
+    val cluster = MetadataCache.toCluster("cluster-id", 
delta.apply(MetadataProvenance.EMPTY))
+    val partition = cluster.partition(new TopicPartition(topicName, 
partitionId))
+
+    assertNotNull(partition)
+    assertEquals(Seq(broker0, broker1), partition.replicas().map(_.id()).toSeq)
+    assertEquals(Seq(broker1), partition.offlineReplicas().map(_.id()).toSeq)
+    val fencedBroker = cluster.nodeById(broker1)
+    assertNotNull(fencedBroker)
+    assertTrue(fencedBroker.isFenced)
+  }
 
   val oldRequestControllerEpoch: Int = 122
   val newRequestControllerEpoch: Int = 123
diff --git 
a/metadata/src/main/java/org/apache/kafka/metadata/MetadataCache.java 
b/metadata/src/main/java/org/apache/kafka/metadata/MetadataCache.java
index 9fb3ec24640..5ec39080bb4 100644
--- a/metadata/src/main/java/org/apache/kafka/metadata/MetadataCache.java
+++ b/metadata/src/main/java/org/apache/kafka/metadata/MetadataCache.java
@@ -37,11 +37,11 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collection;
 import java.util.Collections;
-import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ThreadLocalRandom;
@@ -148,10 +148,8 @@ public interface MetadataCache extends ConfigRepository {
         boolean ignoreTopicsWithExceptions);
 
     static Cluster toCluster(String clusterId, MetadataImage image) {
-        Map<Integer, List<Node>> brokerToNodes = new HashMap<>();
-        image.cluster().brokers().values().stream()
-            .filter(broker -> !broker.fenced())
-            .forEach(broker -> brokerToNodes.put(broker.id(), broker.nodes()));
+        Map<Integer, List<Node>> brokerToNodes = 
image.cluster().brokers().values().stream()
+            .collect(Collectors.toMap(BrokerRegistration::id, 
BrokerRegistration::nodes));
 
         List<PartitionInfo> partitionInfos = new ArrayList<>();
         Set<String> internalTopics = new HashSet<>();
@@ -169,6 +167,7 @@ public interface MetadataCache extends ConfigRepository {
                             toArray(partition.isr, brokerToNodes),
                             getOfflineReplicas(image, partition).stream()
                                 .map(brokerToNodes::get)
+                                .filter(Objects::nonNull)
                                 .flatMap(Collection::stream)
                                 .toArray(Node[]::new)
                         ));
@@ -197,6 +196,7 @@ public interface MetadataCache extends ConfigRepository {
     private static Node[] toArray(int[] replicas, Map<Integer, List<Node>> 
brokerToNodes) {
         return Arrays.stream(replicas)
             .mapToObj(brokerToNodes::get)
+            .filter(Objects::nonNull)
             .flatMap(Collection::stream)
             .toArray(Node[]::new);
     }

Reply via email to