Skip to content
24 changes: 20 additions & 4 deletions core/src/main/java/kafka/server/share/DelayedShareFetch.java
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.requests.FetchRequest;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.raft.errors.NotLeaderException;
import org.apache.kafka.server.LogReadResult;
import org.apache.kafka.server.metrics.KafkaMetricsGroup;
import org.apache.kafka.server.purgatory.DelayedOperation;
Expand Down Expand Up @@ -368,6 +369,14 @@ public boolean tryComplete() {
"topic partitions {}", shareFetch.groupId(), shareFetch.memberId(),
sharePartitions.keySet());
}
// At this point, there could be delayed requests sitting in the purgatory which are waiting on
// DelayedShareFetchPartitionKeys corresponding to partitions, whose leader has been changed to a different broker.
// In that case, such partitions would not be able to get acquired, and the tryComplete will keep on returning false.
// Eventually the operation will get timed out and completed, but it might not get removed from the purgatory.
// This has been eventually left it like this because the purging mechanism will trigger whenever the number of completed
// but still being watched operations is larger than the purge interval. This purge interval is defined by the config
// share.fetch.purgatory.purge.interval.requests and is 1000 by default, thereby ensuring that such stale operations do not
// grow indefinitely.
return false;
} catch (Exception e) {
log.error("Error processing delayed share fetch request", e);
Expand Down Expand Up @@ -757,30 +766,37 @@ private void processRemoteFetchOrException(
* Case a: The partition is in an offline log directory on this broker
* Case b: This broker does not know the partition it tries to fetch
* Case c: This broker is no longer the leader of the partition it tries to fetch
* Case d: All remote storage read requests completed
* Case d: This broker is no longer the leader or follower of the partition it tries to fetch
* Case e: All remote storage read requests completed
* @return boolean representing whether the remote fetch is completed or not.
*/
private boolean maybeCompletePendingRemoteFetch() {
boolean canComplete = false;

for (TopicIdPartition topicIdPartition : pendingRemoteFetchesOpt.get().fetchOffsetMetadataMap().keySet()) {
try {
replicaManager.getPartitionOrException(topicIdPartition.topicPartition());
Partition partition = replicaManager.getPartitionOrException(topicIdPartition.topicPartition());
if (!partition.isLeader()) {
throw new NotLeaderException("Broker is no longer the leader of topicPartition: " + topicIdPartition);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Don't you need to handle the exception below?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the review. Yep, I have added another catch block here for the new NotLeaderException thrown.

}
} catch (KafkaStorageException e) { // Case a
log.debug("TopicPartition {} is in an offline log directory, satisfy {} immediately", topicIdPartition, shareFetch.fetchParams());
canComplete = true;
} catch (UnknownTopicOrPartitionException e) { // Case b
log.debug("Broker no longer knows of topicPartition {}, satisfy {} immediately", topicIdPartition, shareFetch.fetchParams());
canComplete = true;
} catch (NotLeaderOrFollowerException e) { // Case c
} catch (NotLeaderException e) { // Case c
log.debug("Broker is no longer the leader of topicPartition {}, satisfy {} immediately", topicIdPartition, shareFetch.fetchParams());
canComplete = true;
} catch (NotLeaderOrFollowerException e) { // Case d
log.debug("Broker is no longer the leader or follower of topicPartition {}, satisfy {} immediately", topicIdPartition, shareFetch.fetchParams());
canComplete = true;
}
if (canComplete)
break;
}

if (canComplete || pendingRemoteFetchesOpt.get().isDone()) { // Case d
if (canComplete || pendingRemoteFetchesOpt.get().isDone()) { // Case e
return forceComplete();
} else
return false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -774,6 +774,7 @@ private static void removeSharePartitionFromCache(
if (sharePartition != null) {
sharePartition.markFenced();
replicaManager.removeListener(sharePartitionKey.topicIdPartition().topicPartition(), sharePartition.listener());
replicaManager.completeDelayedShareFetchRequest(new DelayedShareFetchGroupKey(sharePartitionKey.groupId(), sharePartitionKey.topicIdPartition()));
}
}

Expand Down
172 changes: 172 additions & 0 deletions core/src/test/java/kafka/server/share/DelayedShareFetchTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -153,13 +153,24 @@ public void testDelayedShareFetchTryCompleteReturnsFalseDueToNonAcquirablePartit
when(sp0.canAcquireRecords()).thenReturn(false);
when(sp1.canAcquireRecords()).thenReturn(false);

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

ReplicaManager replicaManager = mock(ReplicaManager.class);
when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);

ShareGroupMetrics shareGroupMetrics = new ShareGroupMetrics(new MockTime());
Uuid fetchId = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
.withShareGroupMetrics(shareGroupMetrics)
.withFetchId(fetchId)
.withReplicaManager(replicaManager)
.build());

when(sp0.maybeAcquireFetchLock(fetchId)).thenReturn(true);
Expand Down Expand Up @@ -218,6 +229,15 @@ public void testTryCompleteWhenMinBytesNotSatisfiedOnFirstFetch() {

PartitionMaxBytesStrategy partitionMaxBytesStrategy = mockPartitionMaxBytes(Set.of(tp0));

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);

Time time = mock(Time.class);
when(time.hiResClockMs()).thenReturn(100L).thenReturn(110L);
ShareGroupMetrics shareGroupMetrics = new ShareGroupMetrics(time);
Expand Down Expand Up @@ -287,6 +307,15 @@ public void testTryCompleteWhenMinBytesNotSatisfiedOnSubsequentFetch() {
mockTopicIdPartitionFetchBytes(replicaManager, tp0, hwmOffsetMetadata);
BiConsumer<SharePartitionKey, Throwable> exceptionHandler = mockExceptionHandler();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);

Uuid fetchId = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
Expand Down Expand Up @@ -580,6 +609,19 @@ public void testForceCompleteTriggersDelayedActionsQueue() {
List<DelayedOperationKey> delayedShareFetchWatchKeys = new ArrayList<>();
topicIdPartitions1.forEach(topicIdPartition -> delayedShareFetchWatchKeys.add(new DelayedShareFetchGroupKey(groupId, topicIdPartition.topicId(), topicIdPartition.partition())));

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

Partition p2 = mock(Partition.class);
when(p2.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);
when(replicaManager.getPartitionOrException(tp2.topicPartition())).thenReturn(p2);

Uuid fetchId1 = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch1 = DelayedShareFetchTest.DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch1)
Expand Down Expand Up @@ -737,6 +779,12 @@ public void testExceptionInMinBytesCalculation() {
when(time.hiResClockMs()).thenReturn(100L).thenReturn(110L).thenReturn(170L);
ShareGroupMetrics shareGroupMetrics = new ShareGroupMetrics(time);
Uuid fetchId = Uuid.randomUuid();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);

DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
Expand Down Expand Up @@ -881,10 +929,18 @@ public void testLocksReleasedAcquireException() {
BROKER_TOPIC_STATS);

Uuid fetchId = Uuid.randomUuid();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

ReplicaManager replicaManager = mock(ReplicaManager.class);
when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);

DelayedShareFetch delayedShareFetch = DelayedShareFetchTest.DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
.withFetchId(fetchId)
.withReplicaManager(replicaManager)
.build();

when(sp0.maybeAcquireFetchLock(fetchId)).thenReturn(true);
Expand Down Expand Up @@ -1263,6 +1319,19 @@ public void testRemoteStorageFetchTryCompleteReturnsFalse() {
when(remoteLogManager.asyncRead(any(), any())).thenReturn(mock(Future.class));
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

Partition p2 = mock(Partition.class);
when(p2.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);
when(replicaManager.getPartitionOrException(tp2.topicPartition())).thenReturn(p2);

Uuid fetchId = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
Expand All @@ -1288,6 +1357,70 @@ public void testRemoteStorageFetchTryCompleteReturnsFalse() {
delayedShareFetch.lock().unlock();
}

@Test
public void testRemoteStorageFetchPartitionLeaderChanged() {
ReplicaManager replicaManager = mock(ReplicaManager.class);
TopicIdPartition tp0 = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition("foo", 0));

SharePartition sp0 = mock(SharePartition.class);

when(sp0.canAcquireRecords()).thenReturn(true);

LinkedHashMap<TopicIdPartition, SharePartition> sharePartitions = new LinkedHashMap<>();
sharePartitions.put(tp0, sp0);

ShareFetch shareFetch = new ShareFetch(FETCH_PARAMS, "grp", Uuid.randomUuid().toString(),
new CompletableFuture<>(), List.of(tp0), BATCH_SIZE, MAX_FETCH_RECORDS,
BROKER_TOPIC_STATS);

when(sp0.nextFetchOffset()).thenReturn(10L);

// Fetch offset does not match with the cached entry for sp0, hence, a replica manager fetch will happen for sp0.
when(sp0.fetchOffsetMetadata(anyLong())).thenReturn(Optional.empty());

// Mocking remote storage read result for tp0.
doAnswer(invocation -> buildLocalAndRemoteFetchResult(Set.of(), Set.of(tp0))).when(replicaManager).readFromLog(any(), any(), any(ReplicaQuota.class), anyBoolean());

// Remote fetch related mocks. Remote fetch object does not complete within tryComplete in this mock.
RemoteLogManager remoteLogManager = mock(RemoteLogManager.class);
when(remoteLogManager.asyncRead(any(), any())).thenReturn(mock(Future.class));
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(false);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);

Uuid fetchId = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
.withReplicaManager(replicaManager)
.withPartitionMaxBytesStrategy(mockPartitionMaxBytes(Set.of(tp0)))
.withFetchId(fetchId)
.build());

// All the topic partitions are acquirable.
when(sp0.maybeAcquireFetchLock(fetchId)).thenReturn(true);

// Mock the behaviour of replica manager such that remote storage fetch completion timer task completes on adding it to the watch queue.
doAnswer(invocationOnMock -> {
TimerTask timerTask = invocationOnMock.getArgument(0);
timerTask.run();
return null;
}).when(replicaManager).addShareFetchTimerRequest(any());

assertFalse(delayedShareFetch.isCompleted());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shouldn't the delayed share fetch complete?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the review. The DelayedShareFetch should complete, but only after tryComplete is called. I have added an assertTrue for this after tryComplete is called to make the test better.

assertTrue(delayedShareFetch.tryComplete());
assertTrue(delayedShareFetch.isCompleted());
// Remote fetch object gets created for delayed share fetch object.
assertNotNull(delayedShareFetch.pendingRemoteFetches());
// Verify the locks are released for local log read topic partitions tp0.
Mockito.verify(delayedShareFetch, times(1)).releasePartitionLocks(Set.of(tp0));
assertTrue(delayedShareFetch.lock().tryLock());
delayedShareFetch.lock().unlock();
}

@Test
public void testRemoteStorageFetchTryCompleteThrowsException() {
ReplicaManager replicaManager = mock(ReplicaManager.class);
Expand Down Expand Up @@ -1516,6 +1649,16 @@ public void testRemoteStorageFetchRequestCompletionOnFutureCompletionFailure() {
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Uuid fetchId = Uuid.randomUuid();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);

DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
Expand Down Expand Up @@ -1586,6 +1729,12 @@ public void testRemoteStorageFetchRequestCompletionOnFutureCompletionSuccessfull
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Uuid fetchId = Uuid.randomUuid();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);

DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
Expand Down Expand Up @@ -1679,6 +1828,19 @@ public void testRemoteStorageFetchRequestCompletionAlongWithLocalLogRead() {
}).when(remoteLogManager).asyncRead(any(), any());
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

Partition p2 = mock(Partition.class);
when(p2.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);
when(replicaManager.getPartitionOrException(tp2.topicPartition())).thenReturn(p2);

Uuid fetchId = Uuid.randomUuid();
DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
Expand Down Expand Up @@ -1761,6 +1923,16 @@ public void testRemoteStorageFetchHappensForAllTopicPartitions() {
when(replicaManager.remoteLogManager()).thenReturn(Option.apply(remoteLogManager));

Uuid fetchId = Uuid.randomUuid();

Partition p0 = mock(Partition.class);
when(p0.isLeader()).thenReturn(true);

Partition p1 = mock(Partition.class);
when(p1.isLeader()).thenReturn(true);

when(replicaManager.getPartitionOrException(tp0.topicPartition())).thenReturn(p0);
when(replicaManager.getPartitionOrException(tp1.topicPartition())).thenReturn(p1);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Did we add a new test case that validates the added functionality?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the review. I have made the required updates and added a test case for the remote fetch throwing a NotLeaderException as well


DelayedShareFetch delayedShareFetch = spy(DelayedShareFetchBuilder.builder()
.withShareFetchData(shareFetch)
.withSharePartitions(sharePartitions)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2622,7 +2622,8 @@ public void testSharePartitionPartialInitializationFailure() throws Exception {
assertEquals(Errors.FENCED_STATE_EPOCH.code(), partitionDataMap.get(tp2).errorCode());
assertEquals("Fenced state epoch", partitionDataMap.get(tp2).errorMessage());

Mockito.verify(replicaManager, times(0)).completeDelayedShareFetchRequest(any());
Mockito.verify(replicaManager, times(1)).completeDelayedShareFetchRequest(
new DelayedShareFetchGroupKey(groupId, tp2));
Mockito.verify(replicaManager, times(1)).readFromLog(
any(), any(), any(ReplicaQuota.class), anyBoolean());
// Should have 1 fetch recorded and 1 failure as single topic has multiple partition fetch
Expand Down