From fafbf81af64b01f5c6fd6b1ca57a4d77ef0f70e3 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Wed, 19 Jul 2023 23:19:10 +0530 Subject: [PATCH 1/4] KAFKA-15218: Avoid NPE thrown while deleting topic and fetch from follower concurrently --- core/src/main/scala/kafka/cluster/Partition.scala | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index e1478c62d510c..a23ce25c8f71f 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -998,7 +998,13 @@ class Partition(val topicPartition: TopicPartition, // 3. Its metadata cached broker epoch matches its Fetch request broker epoch. Or the Fetch // request broker epoch is -1 which bypasses the epoch verification. case kRaftMetadataCache: KRaftMetadataCache => - val storedBrokerEpoch = remoteReplicasMap.get(followerReplicaId).stateSnapshot.brokerEpoch + val mayBeReplica = getReplica(followerReplicaId) + // The topic is already deleted and we don't have any replica information. In this case, we can return false + // so as to avoid NPE + if (mayBeReplica.isEmpty) { + return false + } + val storedBrokerEpoch = mayBeReplica.get.stateSnapshot.brokerEpoch val cachedBrokerEpoch = kRaftMetadataCache.getAliveBrokerEpoch(followerReplicaId) !kRaftMetadataCache.isBrokerFenced(followerReplicaId) && !kRaftMetadataCache.isBrokerShuttingDown(followerReplicaId) && From 74d261e6121f67039ecc00872b0547d73baabdf7 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Thu, 20 Jul 2023 15:25:13 +0530 Subject: [PATCH 2/4] Adding warning log line --- core/src/main/scala/kafka/cluster/Partition.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index a23ce25c8f71f..5dee341fe9f54 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -1002,6 +1002,7 @@ class Partition(val topicPartition: TopicPartition, // The topic is already deleted and we don't have any replica information. In this case, we can return false // so as to avoid NPE if (mayBeReplica.isEmpty) { + warn(s"The replica state of replica ID:[$followerReplicaId] doesn't exist in the leader node. It might because the topic is already deleted.") return false } val storedBrokerEpoch = mayBeReplica.get.stateSnapshot.brokerEpoch From c6e894545d94e7b9c2abfc258bab86a53e53b87e Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Mon, 24 Jul 2023 06:23:23 +0530 Subject: [PATCH 3/4] Adding unit test --- .../unit/kafka/cluster/PartitionTest.scala | 54 ++++++++++++++++++- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index ab3d06ef65369..9bc3bf638bc62 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -37,7 +37,7 @@ import org.apache.kafka.metadata.LeaderRecoveryState import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.Test import org.mockito.ArgumentMatchers -import org.mockito.ArgumentMatchers.{any, anyString} +import org.mockito.ArgumentMatchers.{any, anyBoolean, anyInt, anyLong, anyString} import org.mockito.Mockito._ import org.mockito.invocation.InvocationOnMock @@ -55,7 +55,7 @@ import org.apache.kafka.server.common.MetadataVersion.IBP_2_6_IV0 import org.apache.kafka.server.metrics.KafkaYammerMetrics import org.apache.kafka.server.util.{KafkaScheduler, MockTime} import org.apache.kafka.storage.internals.epoch.LeaderEpochFileCache -import org.apache.kafka.storage.internals.log.{AppendOrigin, CleanerConfig, EpochEntry, FetchIsolation, FetchParams, LogAppendInfo, LogDirFailureChannel, LogReadInfo, LogStartOffsetIncrementReason, ProducerStateManager, ProducerStateManagerConfig} +import org.apache.kafka.storage.internals.log.{AppendOrigin, CleanerConfig, EpochEntry, FetchIsolation, FetchParams, LogAppendInfo, LogDirFailureChannel, LogOffsetMetadata, LogReadInfo, LogStartOffsetIncrementReason, ProducerStateManager, ProducerStateManagerConfig} import org.junit.jupiter.params.ParameterizedTest import org.junit.jupiter.params.provider.ValueSource @@ -1290,6 +1290,56 @@ class PartitionTest extends AbstractPartitionTest { ) } + @Test + def testIsReplicaIsrEligibleWithEmptyReplicaMap(): Unit = { + val mockMetadataCache = mock(classOf[KRaftMetadataCache]) + val partition = spy(new Partition(topicPartition, + replicaLagTimeMaxMs = Defaults.ReplicaLagTimeMaxMs, + interBrokerProtocolVersion = interBrokerProtocolVersion, + localBrokerId = brokerId, + () => defaultBrokerEpoch(brokerId), + time, + alterPartitionListener, + delayedOperations, + mockMetadataCache, + logManager, + alterPartitionManager)) + + when(offsetCheckpoints.fetch(ArgumentMatchers.anyString, ArgumentMatchers.eq(topicPartition))) + .thenReturn(None) + val log = logManager.getOrCreateLog(topicPartition, topicId = None) + seedLogData(log, numRecords = 6, leaderEpoch = 4) + + val controllerEpoch = 0 + val leaderEpoch = 5 + val remoteBrokerId = brokerId + 1 + val replicas = List[Integer](brokerId, remoteBrokerId).asJava + + partition.createLogIfNotExists(isNew = false, isFutureReplica = false, offsetCheckpoints, None) + + val initializeTimeMs = time.milliseconds() + assertTrue(partition.makeLeader( + new LeaderAndIsrPartitionState() + .setControllerEpoch(controllerEpoch) + .setLeader(brokerId) + .setLeaderEpoch(leaderEpoch) + .setIsr(List[Integer](brokerId).asJava) + .setPartitionEpoch(1) + .setReplicas(replicas) + .setIsNew(true), + offsetCheckpoints, None), "Expected become leader transition to succeed") + + doAnswer(_ => { + // simulate topic is deleted at the moment + partition.delete() + val replica = new Replica(remoteBrokerId, topicPartition) + partition.updateFollowerFetchState(replica, mock(classOf[LogOffsetMetadata]), 0, initializeTimeMs, 0, 0) + mock(classOf[LogReadInfo]) + }).when(partition).fetchRecords(any(), any(), anyLong(), anyInt(), anyBoolean(), anyBoolean()) + + fetchFollower(partition, replicaId = remoteBrokerId, fetchOffset = 3L) + } + @Test def testInvalidAlterPartitionRequestsAreNotRetried(): Unit = { val log = logManager.getOrCreateLog(topicPartition, topicId = None) From 999d3223792ee518fbc3cd45582c1f6f79cd688c Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Mon, 24 Jul 2023 10:20:36 +0530 Subject: [PATCH 4/4] Adding assert for no errors thrown --- core/src/test/scala/unit/kafka/cluster/PartitionTest.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 9bc3bf638bc62..3326b217b8f06 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -1337,7 +1337,7 @@ class PartitionTest extends AbstractPartitionTest { mock(classOf[LogReadInfo]) }).when(partition).fetchRecords(any(), any(), anyLong(), anyInt(), anyBoolean(), anyBoolean()) - fetchFollower(partition, replicaId = remoteBrokerId, fetchOffset = 3L) + assertDoesNotThrow(() => fetchFollower(partition, replicaId = remoteBrokerId, fetchOffset = 3L)) } @Test