This is an automated email from the ASF dual-hosted git repository. lhotari pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/pulsar.git
commit 86256f15cfa415fe3febc47bf0bd840fa48242ad Author: void-ptr974 <[email protected]> AuthorDate: Tue Jun 30 00:23:12 2026 +0800 [fix][broker] Fix replicator getting stuck under rate limiter throttling and honor readBatchSize/maxReadSizeBytes on the default read path (#26005) Co-authored-by: Lari Hotari <[email protected]> (cherry picked from commit abca10f9ce93da79a2842455dafda707042d42fa) --- .../persistent/GeoPersistentReplicator.java | 6 +- .../service/persistent/PersistentReplicator.java | 185 ++++++-------- .../broker/service/OneWayReplicatorTest.java | 80 +++++++ .../PersistentReplicatorInflightTaskTest.java | 265 +++++++++++++-------- 4 files changed, 321 insertions(+), 215 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/GeoPersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/GeoPersistentReplicator.java index 7345e938ad4..2d20ccb7970 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/GeoPersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/GeoPersistentReplicator.java @@ -255,9 +255,9 @@ public class GeoPersistentReplicator extends PersistentReplicator { * Explain the result of the race-condition between: * - {@link #readMoreEntries} * - {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} - * Since {@link #acquirePermitsIfNotFetchingSchema} and - * {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} acquire the - * same lock, it is safe. + * Since the read scheduling path in {@link #readMoreEntries()} and + * {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} update in-flight + * read state under the same lock, it is safe. */ beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding.Fetching_Schema); inFlightTask.incCompletedEntries(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 392f5eb98db..ed599cd2205 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -38,7 +38,6 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; -import lombok.AllArgsConstructor; import lombok.Data; import lombok.Getter; import org.apache.bookkeeper.mledger.AsyncCallbacks; @@ -91,7 +90,7 @@ public abstract class PersistentReplicator extends AbstractReplicator protected Optional<DispatchRateLimiter> dispatchRateLimiter = Optional.empty(); private final Object dispatchRateLimiterLock = new Object(); - private int readBatchSize; + private volatile int readBatchSize; private final int readMaxSizeBytes; private final int producerQueueThreshold; @@ -138,10 +137,8 @@ public abstract class PersistentReplicator extends AbstractReplicator this.expiryMonitor = new PersistentMessageExpiryMonitor(localTopic, Codec.decode(cursor.getName()), cursor, null); - readBatchSize = Math.min( - producerQueueSize, - localTopic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadBatchSize()); - readMaxSizeBytes = localTopic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadSizeBytes(); + readBatchSize = getMaxReadBatchSize(); + readMaxSizeBytes = brokerService.pulsar().getConfiguration().getDispatcherMaxReadSizeBytes(); producerQueueThreshold = (int) (producerQueueSize * 0.9); this.initializeDispatchRateLimiterIfNeeded(); @@ -149,6 +146,10 @@ public abstract class PersistentReplicator extends AbstractReplicator startProducer(); } + private int getMaxReadBatchSize() { + return Math.min(producerQueueSize, brokerService.pulsar().getConfiguration().getDispatcherMaxReadBatchSize()); + } + @Override protected void setProducerAndTriggerReadEntries(Producer<byte[]> producer) { /** @@ -212,72 +213,62 @@ public abstract class PersistentReplicator extends AbstractReplicator this.cursor.setInactive(); } - @Data - @AllArgsConstructor - private static class AvailablePermits { - private int messages; - private long bytes; - - /** - * messages, bytes - * 0, O: Producer queue is full, no permits. - * -1, -1: Rate Limiter reaches limit. - * >0, >0: available permits for read entries. - */ - public boolean isExceeded() { - return messages == -1 && bytes == -1; - } - + private record ReadLimits(int messages, long bytes) { public boolean isReadable() { return messages > 0 && bytes > 0; } } /** - * Calculate available permits for read entries. + * Calculate read limits for a read operation. Takes the rate limiter into account if it's enabled. + * Also limits to current readBatchSize and readMaxSizeBytes. */ - private AvailablePermits getRateLimiterAvailablePermits(int availablePermits) { + private ReadLimits getReadLimits(int permits) { // return 0, if Producer queue is full, it will pause read entries. - if (availablePermits <= 0) { + if (permits <= 0) { if (log.isDebugEnabled()) { - log.debug("[{}] Producer queue is full, availablePermits: {}, pause reading", - replicatorId, availablePermits); + log.debug("[{}] Producer queue is full, permits: {}, pause reading", replicatorId, permits); } - return new AvailablePermits(0, 0); + return new ReadLimits(0, 0); } - long availablePermitsOnMsg = -1; - long availablePermitsOnByte = -1; + long readLimitOnMsg; + long readLimitOnByte; // handle rate limit if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { DispatchRateLimiter rateLimiter = dispatchRateLimiter.get(); - // if dispatch-rate is in msg then read only msg according to available permit - availablePermitsOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); - availablePermitsOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); - // no permits from rate limit - if (availablePermitsOnByte == 0 || availablePermitsOnMsg == 0) { + // rateLimiter returns -1 if there is no rate limit configured + readLimitOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); + readLimitOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); + // no permits from rate limit when either limit is 0 + if (readLimitOnByte == 0 || readLimitOnMsg == 0) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded topic replicator message-rate {}/{}," - + " schedule after a {}", + log.debug("[{}] Message-read exceeded topic replicator rate limit," + + " dispatchRateOnMsg: {}, dispatchRateOnByte: {}," + + " readLimitOnMsg: {}, readLimitOnByte: {}", replicatorId, rateLimiter.getDispatchRateOnMsg(), rateLimiter.getDispatchRateOnByte(), - MESSAGE_RATE_BACKOFF_MS); + readLimitOnMsg, + readLimitOnByte); } - return new AvailablePermits(-1, -1); + return new ReadLimits(-1, -1); } + // use given permits if no rate limit configured, otherwise limit to returned rate limiter permits + readLimitOnMsg = readLimitOnMsg == -1 ? permits : Math.min(permits, readLimitOnMsg); + // use readMaxSizeBytes if no rate limit configured, otherwise limit to returned rate limiter permits + readLimitOnByte = readLimitOnByte == -1 ? readMaxSizeBytes : Math.min(readMaxSizeBytes, readLimitOnByte); + } else { + readLimitOnMsg = permits; + readLimitOnByte = readMaxSizeBytes; } - availablePermitsOnMsg = - availablePermitsOnMsg == -1 ? availablePermits : Math.min(availablePermits, availablePermitsOnMsg); - availablePermitsOnMsg = Math.min(availablePermitsOnMsg, readBatchSize); - - availablePermitsOnByte = - availablePermitsOnByte == -1 ? readMaxSizeBytes : Math.min(readMaxSizeBytes, availablePermitsOnByte); + // limit messages to current read batch size + readLimitOnMsg = Math.min(readLimitOnMsg, readBatchSize); - return new AvailablePermits((int) availablePermitsOnMsg, availablePermitsOnByte); + return new ReadLimits((int) readLimitOnMsg, readLimitOnByte); } public void disconnectIfNoTrafficAndBacklog() { @@ -308,8 +299,39 @@ public abstract class PersistentReplicator extends AbstractReplicator if (state.equals(Terminated) || state.equals(Terminating)) { return; } - // Acquire permits and check state of producer. - InFlightTask newInFlightTask = acquirePermitsIfNotFetchingSchema(); + InFlightTask newInFlightTask = null; + ReadLimits readLimits = null; + synchronized (inFlightTasks) { + if (hasPendingRead()) { + log.debug("Skip the reading because there is a pending read task"); + } else if (waitForCursorRewindingRefCnf > 0) { + log.debug("Skip the reading due to new detected schema"); + } else if (state != Started) { + log.debug("Skip the reading because producer has not started"); + } else { + int permits = getPermitsIfNoPendingRead(); + if (permits > 0) { + if (!isWritable()) { + log.debug("Throttling replication traffic to a single message permit because producer is not " + + "writable"); + // Minimize the read size if the producer is disconnected or the window is already full. + permits = 1; + } + + readLimits = getReadLimits(permits); + if (readLimits.isReadable()) { + newInFlightTask = createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), + readLimits.messages); + } else { + // no rate limiter permits from rate limit + if (log.isDebugEnabled()) { + log.debug("[{}] Throttling replication traffic. Messages To Read {}, Bytes To Read {}", + replicatorId, readLimits.messages, readLimits.bytes); + } + } + } + } + } if (newInFlightTask == null) { // no permits from rate limit if (log.isDebugEnabled()) { @@ -318,48 +340,14 @@ public abstract class PersistentReplicator extends AbstractReplicator if (!hasPendingRead()) { topic.getBrokerService().executor().schedule( () -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); - return; - } else { - return; - } - } - // If disabled RateLimiter. - if (!dispatchRateLimiter.isPresent() || !dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, -1, this, - newInFlightTask/* Context object */, topic.getMaxReadPosition()); - return; - } - // No permits of RateLimiter. - AvailablePermits availablePermits = getRateLimiterAvailablePermits(newInFlightTask.readingEntries); - if (!availablePermits.isReadable()) { - // no rate limiter permits from rate limit - if (log.isDebugEnabled()) { - log.debug("[{}] Throttling replication traffic. Messages To Read {}, Bytes To Read {}", - replicatorId, availablePermits.getMessages(), availablePermits.getBytes()); } - topic.getBrokerService().executor().schedule( - () -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); return; } - // Has permits of RateLimiter. - int messagesToRead = availablePermits.getMessages(); - long bytesToRead = availablePermits.getBytes(); - if (!isWritable()) { - if (log.isDebugEnabled()) { - log.debug("[{}] Throttling replication traffic because producer is not writable", replicatorId); - } - // Minimize the read size if the producer is disconnected or the window is already full - messagesToRead = 1; - } - // Update acquired permits exceeds limitation. - if (messagesToRead < newInFlightTask.readingEntries) { - newInFlightTask.setReadingEntries(messagesToRead); - } if (log.isDebugEnabled()) { - log.debug("[{}] Schedule read of {} messages or {} bytes", replicatorId, newInFlightTask.readingEntries, - bytesToRead); + log.debug("[{}] Schedule read of {} messages or {} bytes", replicatorId, + newInFlightTask.readingEntries, readLimits.bytes); } - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, bytesToRead, this, + cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, readLimits.bytes, this, newInFlightTask/* Context object */, topic.getMaxReadPosition()); } @@ -414,7 +402,7 @@ public abstract class PersistentReplicator extends AbstractReplicator inFlightTask.setEntries(entries); // After the replicator starts, the speed will be gradually increased. - int maxReadBatchSize = topic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadBatchSize(); + int maxReadBatchSize = getMaxReadBatchSize(); if (readBatchSize < maxReadBatchSize) { int newReadBatchSize = Math.min(readBatchSize * 2, maxReadBatchSize); if (log.isDebugEnabled()) { @@ -573,7 +561,7 @@ public abstract class PersistentReplicator extends AbstractReplicator } // Reduce read batch size to avoid flooding bookies with retries - readBatchSize = topic.getBrokerService().pulsar().getConfiguration().getDispatcherMinReadBatchSize(); + readBatchSize = brokerService.pulsar().getConfiguration().getDispatcherMinReadBatchSize(); long waitTimeMillis = readFailureBackoff.next().toMillis(); @@ -930,29 +918,6 @@ public abstract class PersistentReplicator extends AbstractReplicator } } - protected InFlightTask acquirePermitsIfNotFetchingSchema() { - synchronized (inFlightTasks) { - if (hasPendingRead()) { - log.info("[{}] Skip the reading because there is a pending read task", replicatorId); - return null; - } - if (waitForCursorRewindingRefCnf > 0) { - log.info("[{}] Skip the reading due to new detected schema", replicatorId); - return null; - } - if (state != Started) { - log.info("[{}] Skip the reading because producer has not started [{}]", replicatorId, state); - return null; - } - // Guarantee that there is a unique cursor reading task. - int permits = getPermitsIfNoPendingRead(); - if (permits == 0) { - return null; - } - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), permits); - } - } - protected int getPermitsIfNoPendingRead() { synchronized (inFlightTasks) { for (InFlightTask task : inFlightTasks) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java index 673dd082ba4..67151479ed2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java @@ -111,6 +111,7 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.DispatchRate; import org.apache.pulsar.common.policies.data.HierarchyTopicPolicies; import org.apache.pulsar.common.policies.data.PublishRate; import org.apache.pulsar.common.policies.data.ReplicatorStats; @@ -2213,6 +2214,85 @@ public class OneWayReplicatorTest extends OneWayReplicatorTestBase { ensureNoBacklogByInflightTask(getReplicator(topicName)); } + @DataProvider + public Object[][] replicatorDispatchRateLimits() { + return new Object[][] { + {1, -1L}, + {-1, 1L} + }; + } + + @Test(timeOut = 90_000, dataProvider = "replicatorDispatchRateLimits") + public void testReplicatorContinuesAfterRateLimiterHasNoPermits(int messageRate, long byteRate) throws Exception { + final String topicName = BrokerTestUtil.newUniqueName("persistent://" + replicatedNamespace + "/tp_"); + final String subscriptionName = "sub"; + final List<String> messages = Arrays.asList("msg-0", "msg-1", "msg-2"); + DispatchRate dispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(messageRate) + .dispatchThrottlingRateInByte(byteRate) + .ratePeriodInSecond(2) + .build(); + Producer<String> producer = null; + Consumer<String> consumer = null; + boolean topicCreated = false; + boolean dispatchRateConfigured = false; + try { + admin1.topics().createNonPartitionedTopic(topicName); + topicCreated = true; + admin1.topicPolicies().setReplicatorDispatchRate(topicName, dispatchRate); + dispatchRateConfigured = true; + GeoPersistentReplicator replicator = getReplicator(topicName); + Awaitility.await().untilAsserted(() -> { + assertTrue(replicator.getRateLimiter().isPresent()); + assertEquals(replicator.getRateLimiter().get().getDispatchRateOnMsg(), messageRate); + assertEquals(replicator.getRateLimiter().get().getDispatchRateOnByte(), byteRate); + }); + consumer = client2.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subscriptionName) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscribe(); + producer = client1.newProducer(Schema.STRING) + .topic(topicName) + .enableBatching(false) + .create(); + + for (String message : messages) { + producer.send(message); + } + + Set<String> expected = new HashSet<>(messages); + Set<String> received = new HashSet<>(); + Consumer<String> subscribedConsumer = consumer; + Awaitility.await().atMost(Duration.ofSeconds(60)).untilAsserted(() -> { + Message<String> message = subscribedConsumer.receive(1, TimeUnit.SECONDS); + if (message != null) { + received.add(message.getValue()); + subscribedConsumer.acknowledge(message); + } + assertEquals(received, expected); + }); + waitForReplicationTaskFinish(topicName); + ensureNoBacklogByInflightTask(replicator); + } finally { + if (producer != null) { + producer.close(); + } + if (consumer != null) { + consumer.close(); + } + if (dispatchRateConfigured) { + admin1.topicPolicies().removeReplicatorDispatchRate(topicName); + } + if (topicCreated) { + admin1.topics().setReplicationClusters(topicName, Arrays.asList(cluster1)); + waitReplicatorStopped(topicName, false); + admin1.topics().delete(topicName, true); + admin2.topics().delete(topicName, true); + } + } + } + @DataProvider public Object[] isPartitioned() { return new Object[]{ diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java index 83529969df4..773f2933ce7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java @@ -18,44 +18,64 @@ */ package org.apache.pulsar.broker.service.persistent; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.same; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; +import io.netty.channel.EventLoopGroup; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.LinkedList; import java.util.List; +import java.util.Optional; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerTest; import org.apache.pulsar.broker.BrokerTestUtil; +import org.apache.pulsar.broker.PulsarServerException; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.AbstractReplicator; -import org.apache.pulsar.broker.service.BrokerServiceInternalMethodInvoker; +import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.OneWayReplicatorTestBase; import org.apache.pulsar.broker.service.persistent.PersistentReplicator.InFlightTask; import org.apache.pulsar.broker.service.persistent.PersistentReplicator.ProducerSendCallback; import org.apache.pulsar.broker.service.persistent.PersistentReplicator.ReasonOfWaitForCursorRewinding; +import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.ProducerBuilder; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.impl.PulsarClientImpl; import org.awaitility.Awaitility; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @Slf4j @@ -199,6 +219,63 @@ public class PersistentReplicatorInflightTaskTest extends OneWayReplicatorTestBa } } + @DataProvider + public Object[][] readSchedulingLimits() { + return new Object[][] { + {"message permits exhausted", 0, -1L, true, false, 0, 0L}, + {"byte permits exhausted", -1, 0L, true, false, 0, 0L}, + {"message permits limit read batch", 5, -1L, true, true, 5, 1024L}, + {"byte permits limit read size", -1, 512L, true, true, 100, 512L}, + {"non-writable producer limits read batch", 5, 512L, false, true, 1, 512L} + }; + } + + @Test(dataProvider = "readSchedulingLimits") + public void testReadMoreEntriesSchedulesCursorReadWithReadLimits(String scenario, + long availableMessages, + long availableBytes, + boolean writable, + boolean expectRead, + int expectedMessages, + long expectedBytes) throws Exception { + TestReplicatorFixture fixture = newTestReplicatorFixture(writable); + PersistentReplicator replicator = fixture.replicator; + DispatchRateLimiter rateLimiter = mock(DispatchRateLimiter.class); + when(rateLimiter.isDispatchRateLimitingEnabled()).thenReturn(true); + when(rateLimiter.getAvailableDispatchRateLimitOnMsg()).thenReturn(availableMessages); + when(rateLimiter.getAvailableDispatchRateLimitOnByte()).thenReturn(availableBytes); + replicator.dispatchRateLimiter = Optional.of(rateLimiter); + + replicator.readMoreEntries(); + + if (expectRead) { + assertEquals(replicator.inFlightTasks.size(), 1, scenario); + InFlightTask inFlightTask = replicator.inFlightTasks.peek(); + verify(fixture.cursor).asyncReadEntriesOrWait(eq(expectedMessages), eq(expectedBytes), + same(replicator), same(inFlightTask), any(Position.class)); + assertEquals(inFlightTask.getReadingEntries(), expectedMessages, scenario); + verify(fixture.executor, never()).schedule(any(Runnable.class), anyLong(), any(TimeUnit.class)); + } else { + verify(fixture.cursor, never()).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any()); + assertTrue(replicator.inFlightTasks.isEmpty(), scenario); + verify(fixture.executor).schedule(any(Runnable.class), eq((long) PersistentTopic.MESSAGE_RATE_BACKOFF_MS), + eq(TimeUnit.MILLISECONDS)); + } + } + + @Test + public void testReadMoreEntriesSkipsReadWhenPendingReadExists() throws Exception { + TestReplicatorFixture fixture = newTestReplicatorFixture(true); + PersistentReplicator replicator = fixture.replicator; + replicator.inFlightTasks.add(new InFlightTask(PositionFactory.create(1, 1), 5, replicator.getReplicatorId())); + + replicator.readMoreEntries(); + + verify(fixture.cursor, never()).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any()); + verify(fixture.executor, never()).schedule(any(Runnable.class), anyLong(), any(TimeUnit.class)); + assertEquals(replicator.inFlightTasks.size(), 1); + } + @Test public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { log.info("Starting testCreateOrRecycleInFlightTaskIntoQueue"); @@ -389,111 +466,94 @@ public class PersistentReplicatorInflightTaskTest extends OneWayReplicatorTestBa } } - @Test - public void testAcquirePermitsIfNotFetchingSchema() throws Exception { - log.info("Starting testAcquirePermitsIfNotFetchingSchema"); - // Get the replicator for the test topic - PersistentReplicator replicator = getReplicator(topicName); - Assert.assertNotNull(replicator, "Replicator should not be null"); + @SuppressWarnings("unchecked") + private TestReplicatorFixture newTestReplicatorFixture(boolean writable) throws Exception { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setClusterName("local"); + configuration.setReplicationProducerQueueSize(1000); + configuration.setDispatcherMaxReadBatchSize(100); + configuration.setDispatcherMaxReadSizeBytes(1024); + + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getConfiguration()).thenReturn(configuration); + when(pulsar.getConfig()).thenReturn(configuration); + when(pulsar.getClient()).thenReturn(mock(PulsarClientImpl.class)); + when(pulsar.getAdminClient()).thenReturn(mock(PulsarAdmin.class)); + + BrokerService brokerService = mock(BrokerService.class); + EventLoopGroup executor = mock(EventLoopGroup.class); + when(brokerService.pulsar()).thenReturn(pulsar); + when(brokerService.getPulsar()).thenReturn(pulsar); + when(brokerService.executor()).thenReturn(executor); + + ProducerBuilder<byte[]> producerBuilder = mock(ProducerBuilder.class); + when(producerBuilder.topic(anyString())).thenReturn(producerBuilder); + when(producerBuilder.messageRoutingMode(any())).thenReturn(producerBuilder); + when(producerBuilder.enableBatching(anyBoolean())).thenReturn(producerBuilder); + when(producerBuilder.sendTimeout(anyInt(), any(TimeUnit.class))).thenReturn(producerBuilder); + when(producerBuilder.maxPendingMessages(anyInt())).thenReturn(producerBuilder); + when(producerBuilder.producerName(anyString())).thenReturn(producerBuilder); + + PulsarClientImpl replicationClient = mock(PulsarClientImpl.class); + when(replicationClient.newProducer(any(Schema.class))).thenReturn(producerBuilder); + + PersistentTopic topic = mock(PersistentTopic.class); + when(topic.getName()).thenReturn("persistent://prop/ns/test-read-scheduling"); + when(topic.getReplicatorPrefix()).thenReturn("pulsar.repl"); + when(topic.getBrokerService()).thenReturn(brokerService); + when(topic.getMaxReadPosition()).thenReturn(PositionFactory.create(1, 100)); + + ManagedCursor cursor = mock(ManagedCursor.class); + when(cursor.getName()).thenReturn("pulsar.repl.remote"); + when(cursor.getReadPosition()).thenReturn(PositionFactory.create(1, 1)); + + TestPersistentReplicator replicator = new TestPersistentReplicator(topic, cursor, brokerService, + replicationClient, mock(PulsarAdmin.class), writable); + return new TestReplicatorFixture(replicator, cursor, executor); + } - // Get access to the inFlightTasks list for setup - LinkedList<InFlightTask> inFlightTasks = replicator.inFlightTasks; - Assert.assertNotNull(inFlightTasks, "InFlightTasks list should not be null"); + private static class TestReplicatorFixture { + final TestPersistentReplicator replicator; + final ManagedCursor cursor; + final EventLoopGroup executor; - // Save original tasks and clear for testing - List<InFlightTask> originalTasks = new ArrayList<>(inFlightTasks); - inFlightTasks.clear(); + TestReplicatorFixture(TestPersistentReplicator replicator, ManagedCursor cursor, EventLoopGroup executor) { + this.replicator = replicator; + this.cursor = cursor; + this.executor = executor; + } + } - // Save original state - int originalWaitForCursorRewinding = replicator.waitForCursorRewindingRefCnf; - AbstractReplicator.State originalState = replicator.getState(); + private static class TestPersistentReplicator extends PersistentReplicator { + private final boolean writable; - try { - // Test Case 1: Normal case - no pending read, not waiting for cursor rewinding, state is Started - // Should return a new InFlightTask - // First, check the current permits available - int expectedPermits = replicator.getPermitsIfNoPendingRead(); - Assert.assertTrue(expectedPermits > 0, "Should have available permits for the test"); - InFlightTask task1 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNotNull(task1, "Should return a new InFlightTask in normal case"); - Assert.assertNotNull(task1.getReadPos(), "Task should have a read position"); - Assert.assertEquals(task1.getReadingEntries(), expectedPermits, - "Task readingEntries should equal the number of permits available"); - Assert.assertTrue(inFlightTasks.contains(task1), - "Task should be added to the inFlightTasks list"); - - // Test Case 2: With pending read - should return null - inFlightTasks.clear(); - Position position1 = PositionFactory.create(1, 1); - InFlightTask pendingReadTask = new InFlightTask(position1, 5, ""); - // Don't set readoutEntries to simulate pending read - inFlightTasks.add(pendingReadTask); - InFlightTask task2 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNull(task2, "Should return null when there is a pending read"); + TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, BrokerService brokerService, + PulsarClientImpl replicationClient, PulsarAdmin replicationAdmin, boolean writable) + throws PulsarServerException { + super("local", topic, cursor, "remote", topic.getName(), brokerService, replicationClient, + replicationAdmin); + this.writable = writable; + this.state = State.Started; + } - // Test Case 3: With waitForCursorRewinding=true - should return null - inFlightTasks.clear(); - replicator.waitForCursorRewindingRefCnf = 1; - InFlightTask task3 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNull(task3, "Should return null when waiting for cursor rewinding"); - // Reset for next test - replicator.waitForCursorRewindingRefCnf = 0; - - // Test Case 4: With state != Started - should return null - // We need to use reflection to modify the state since it's protected by AtomicReferenceFieldUpdater - BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, AbstractReplicator.State.Starting); - InFlightTask task4 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNull(task4, "Should return null when state is not Started"); - // Reset state for next test - BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, AbstractReplicator.State.Started); - - // Test Case 5: With limited permits - verify readingEntries is set correctly - inFlightTasks.clear(); - // Add a task with some in-flight messages to reduce available permits - Position positionLimited = PositionFactory.create(10, 10); - InFlightTask limitedTask = new InFlightTask(positionLimited, 5, ""); - // Add enough entries to leave just a small number of permits (e.g., 10) - List<Entry> limitedEntries = new ArrayList<>(); - int entriesCount = 990; - for (int j = 0; j < entriesCount; j++) { - limitedEntries.add(mock(Entry.class)); - } - limitedTask.setEntries(limitedEntries); - inFlightTasks.add(limitedTask); - // Check that we have limited permits available - int limitedPermits = replicator.getPermitsIfNoPendingRead(); - Assert.assertTrue(limitedPermits > 0 && limitedPermits < 20, - "Should have a small number of permits available for testing"); - // Now acquire permits and verify readingEntries matches the limited permits - InFlightTask task5 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNotNull(task5, "Should return a task with limited permits"); - Assert.assertEquals(task5.getReadingEntries(), limitedPermits, - "Task readingEntries should equal the limited number of permits available"); - - // Test Case 6: With permits=0 - should return null - inFlightTasks.clear(); - // Add tasks that will make getPermitsIfNoPendingRead() return 0 - // We need enough in-flight messages to equal producerQueueSize - for (int i = 0; i < 10; i++) { - Position position = PositionFactory.create(i, i); - InFlightTask task = new InFlightTask(position, 5, ""); - List<Entry> entries = new ArrayList<>(); - for (int j = 0; j < 100; j++) { - entries.add(mock(Entry.class)); - } - task.setEntries(entries); - inFlightTasks.add(task); - } - InFlightTask task6 = replicator.acquirePermitsIfNotFetchingSchema(); - Assert.assertNull(task6, "Should return null when permits is 0"); - log.info("Completed testAcquirePermitsIfNotFetchingSchema"); - } finally { - // Restore original state - replicator.waitForCursorRewindingRefCnf = originalWaitForCursorRewinding; - BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, originalState); - // Restore original tasks - inFlightTasks.clear(); - inFlightTasks.addAll(originalTasks); + @Override + protected void startProducer() { + // No-op for scheduling behavior tests. + } + + @Override + protected String getProducerName() { + return "test-replicator"; + } + + @Override + protected boolean isWritable() { + return writable; + } + + @Override + protected boolean replicateEntries(List<Entry> entries, InFlightTask inFlightTask) { + return true; } } @@ -503,7 +563,8 @@ public class PersistentReplicatorInflightTaskTest extends OneWayReplicatorTestBa }); replicator.beforeTerminateOrCursorRewinding(PersistentReplicator.ReasonOfWaitForCursorRewinding.Disconnecting); replicator.doRewindCursor(false); - InFlightTask inFlightTask = replicator.createOrRecycleInFlightTaskIntoQueue(PositionFactory.create(1, 1), 1); + InFlightTask inFlightTask = + replicator.createOrRecycleInFlightTaskIntoQueue(PositionFactory.create(1, 1), 1); return () -> { inFlightTask.setEntries(Collections.emptyList()); replicator.readMoreEntries();
