Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -140,15 +140,25 @@ public List<SubscriptionEvent> poll(
refreshAndGetTopicOwnership(topicName, queues, consumerId);
final List<ConsensusPrefetchingQueue> assignedQueues =
getAssignedQueues(queues, consumerId, ownershipSnapshot);
if (assignedQueues.isEmpty()) {
final List<ConsensusPrefetchingQueue> pollQueues =
new ArrayList<>(buildPollOrderForAssignedQueues(assignedQueues, topicName));
final int assignedQueueCount = pollQueues.size();
pollQueues.addAll(buildFallbackPollOrder(queues, assignedQueues));
if (pollQueues.isEmpty()) {
continue;
}

final List<ConsensusPrefetchingQueue> 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;
}
Expand Down Expand Up @@ -185,7 +195,7 @@ public List<SubscriptionEvent> poll(
eventsToPoll.add(event);
totalSize += currentSize;

if (totalSize >= maxBytes) {
if (isFallbackQueue || totalSize >= maxBytes) {
break;
}
}
Expand Down Expand Up @@ -616,6 +626,21 @@ private List<ConsensusPrefetchingQueue> buildPollOrderForAssignedQueues(
return orderedQueues;
}

private List<ConsensusPrefetchingQueue> buildFallbackPollOrder(
final List<ConsensusPrefetchingQueue> queues,
final List<ConsensusPrefetchingQueue> 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<ConsensusPrefetchingQueue> queues, final SubscriptionCommitContext commitContext) {
final String regionId = commitContext.getRegionId();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,15 @@
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;

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
Expand Down Expand Up @@ -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<SubscriptionEvent> 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<SubscriptionEvent> 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);
Expand Down
Loading