diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBroker.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBroker.java index 6b29f069524f..fd04724488c4 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBroker.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBroker.java @@ -140,15 +140,25 @@ public List poll( refreshAndGetTopicOwnership(topicName, queues, consumerId); final List assignedQueues = getAssignedQueues(queues, consumerId, ownershipSnapshot); - if (assignedQueues.isEmpty()) { + final List pollQueues = + new ArrayList<>(buildPollOrderForAssignedQueues(assignedQueues, topicName)); + final int assignedQueueCount = pollQueues.size(); + pollQueues.addAll(buildFallbackPollOrder(queues, assignedQueues)); + if (pollQueues.isEmpty()) { continue; } - - final List pollQueues = - buildPollOrderForAssignedQueues(assignedQueues, topicName); final int eventsBeforeTopicPoll = eventsToPoll.size(); - for (final ConsensusPrefetchingQueue consensusQueue : pollQueues) { + for (int queueIndex = 0; queueIndex < pollQueues.size(); queueIndex++) { + final boolean isFallbackQueue = queueIndex >= assignedQueueCount; + // Ownership remains the preferred polling path. Only borrow one ready event from another + // region when every assigned region is currently empty, preventing a consumer from timing + // out solely because its region has a temporary dry spell. + if (isFallbackQueue && eventsToPoll.size() > eventsBeforeTopicPoll) { + break; + } + + final ConsensusPrefetchingQueue consensusQueue = pollQueues.get(queueIndex); if (consensusQueue.isClosed()) { continue; } @@ -185,7 +195,7 @@ public List poll( eventsToPoll.add(event); totalSize += currentSize; - if (totalSize >= maxBytes) { + if (isFallbackQueue || totalSize >= maxBytes) { break; } } @@ -616,6 +626,21 @@ private List buildPollOrderForAssignedQueues( return orderedQueues; } + private List buildFallbackPollOrder( + final List queues, + final List assignedQueues) { + return queues.stream() + .filter(queue -> !queue.isClosed()) + .filter(queue -> !assignedQueues.contains(queue)) + // Do not block on another consumer's region. Borrow only events that are already ready. + .filter(queue -> queue.getPrefetchedEventCount() > 0) + .sorted( + Comparator.comparingLong(ConsensusPrefetchingQueue::getLag) + .reversed() + .thenComparing(queue -> queue.getConsensusGroupId().toString())) + .collect(Collectors.toList()); + } + private ConsensusPrefetchingQueue getQueueForCommitContext( final List queues, final SubscriptionCommitContext commitContext) { final String regionId = commitContext.getRegionId(); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBrokerPayloadLimitTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBrokerPayloadLimitTest.java index d07ec0b557f1..c2cb73d8fd2d 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBrokerPayloadLimitTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/ConsensusSubscriptionBrokerPayloadLimitTest.java @@ -35,6 +35,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertSame; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -42,6 +43,7 @@ public class ConsensusSubscriptionBrokerPayloadLimitTest { private static final String CONSUMER_GROUP_ID = "consumerGroup"; private static final String CONSUMER_ID = "consumer"; + private static final String OTHER_CONSUMER_ID = "otherConsumer"; private static final String TOPIC_NAME = "topic"; @Test @@ -73,6 +75,57 @@ public void testPollRequeuesEventThatWouldExceedPayloadLimit() throws Exception verify(secondQueue).requeue(CONSUMER_ID, secondCommitContext); } + @Test + public void testPollBorrowsReadyEventWhenAssignedQueueIsEmpty() throws Exception { + final ConsensusSubscriptionBroker broker = new ConsensusSubscriptionBroker(CONSUMER_GROUP_ID); + final ConsensusPrefetchingQueue fallbackQueue = mock(ConsensusPrefetchingQueue.class); + final ConsensusPrefetchingQueue assignedQueue = mock(ConsensusPrefetchingQueue.class); + final SubscriptionEvent fallbackEvent = mock(SubscriptionEvent.class); + + when(fallbackQueue.getConsensusGroupId()).thenReturn(new DataRegionId(1)); + when(assignedQueue.getConsensusGroupId()).thenReturn(new DataRegionId(2)); + when(fallbackQueue.getPrefetchedEventCount()).thenReturn(1); + when(fallbackQueue.poll(CONSUMER_ID, null)).thenReturn(fallbackEvent); + when(fallbackEvent.getCurrentResponseSize()).thenReturn(30); + bindQueues(broker, Arrays.asList(fallbackQueue, assignedQueue)); + + // The first consumer initially owns both regions. After the second consumer joins, the stable + // ownership balancer moves DataRegion[2] to it and keeps DataRegion[1] with the first consumer. + broker.poll(OTHER_CONSUMER_ID, Collections.singleton(TOPIC_NAME), 60L); + final List events = + broker.poll(CONSUMER_ID, Collections.singleton(TOPIC_NAME), 60L); + + assertEquals(1, events.size()); + assertSame(fallbackEvent, events.get(0)); + verify(assignedQueue).poll(CONSUMER_ID, null); + verify(fallbackQueue).poll(CONSUMER_ID, null); + } + + @Test + public void testPollKeepsOwnershipAffinityWhenAssignedQueueHasData() throws Exception { + final ConsensusSubscriptionBroker broker = new ConsensusSubscriptionBroker(CONSUMER_GROUP_ID); + final ConsensusPrefetchingQueue fallbackQueue = mock(ConsensusPrefetchingQueue.class); + final ConsensusPrefetchingQueue assignedQueue = mock(ConsensusPrefetchingQueue.class); + final SubscriptionEvent fallbackEvent = mock(SubscriptionEvent.class); + final SubscriptionEvent assignedEvent = mock(SubscriptionEvent.class); + + when(fallbackQueue.getConsensusGroupId()).thenReturn(new DataRegionId(1)); + when(assignedQueue.getConsensusGroupId()).thenReturn(new DataRegionId(2)); + when(fallbackQueue.getPrefetchedEventCount()).thenReturn(1); + when(fallbackQueue.poll(CONSUMER_ID, null)).thenReturn(fallbackEvent); + when(assignedQueue.poll(CONSUMER_ID, null)).thenReturn(assignedEvent); + when(assignedEvent.getCurrentResponseSize()).thenReturn(30); + bindQueues(broker, Arrays.asList(fallbackQueue, assignedQueue)); + + broker.poll(OTHER_CONSUMER_ID, Collections.singleton(TOPIC_NAME), 60L); + final List events = + broker.poll(CONSUMER_ID, Collections.singleton(TOPIC_NAME), 60L); + + assertEquals(1, events.size()); + assertSame(assignedEvent, events.get(0)); + verify(fallbackQueue, never()).poll(CONSUMER_ID, null); + } + @Test public void testRequeueInFlightEventsAcrossRegionQueues() throws Exception { final ConsensusSubscriptionBroker broker = new ConsensusSubscriptionBroker(CONSUMER_GROUP_ID);