From 25ed1ac707003d1cab006382f9e742adf4f88391 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 12 Jun 2026 09:40:16 +0800 Subject: [PATCH 01/20] [fix][broker] Prevent replicator from getting stuck when dispatch rate limiter has no permits --- .../persistent/PersistentReplicator.java | 118 ++++++++++-------- .../PersistentReplicatorInflightTaskTest.java | 35 ++++++ 2 files changed, 98 insertions(+), 55 deletions(-) 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 8433d9b9c7ba2..5b1773b60829d 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 @@ -240,6 +240,17 @@ public boolean isReadable() { } } + @VisibleForTesting + static class ReadEntriesRequest { + private final InFlightTask inFlightTask; + private final long bytesToRead; + + ReadEntriesRequest(InFlightTask inFlightTask, long bytesToRead) { + this.inFlightTask = inFlightTask; + this.bytesToRead = bytesToRead; + } + } + /** * Calculate available permits for read entries. */ @@ -287,9 +298,8 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - // Acquire permits and check state of producer. - InFlightTask newInFlightTask = acquirePermitsIfNotFetchingSchema(); - if (newInFlightTask == null) { + ReadEntriesRequest request = acquireReadEntriesRequestIfNeeded(); + if (request == null) { // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); if (!hasPendingRead()) { @@ -300,44 +310,60 @@ protected void readMoreEntries() { 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 - log.debug() - .attr("messages", availablePermits.getMessages()) - .attr("bytes", availablePermits.getBytes()) - .log("Throttling replication traffic"); - 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()) { - log.debug("Throttling replication traffic because producer is not writable"); - // 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); - } + InFlightTask newInFlightTask = request.inFlightTask; log.debug() .attr("readingEntries", newInFlightTask.readingEntries) - .attr("bytesToRead", bytesToRead) + .attr("bytesToRead", request.bytesToRead) .log("Scheduling read"); - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, bytesToRead, this, + cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, request.bytesToRead, this, newInFlightTask/* Context object */, topic.getMaxReadPosition()); } + @VisibleForTesting + ReadEntriesRequest acquireReadEntriesRequestIfNeeded() { + synchronized (inFlightTasks) { + if (hasPendingRead()) { + log.info("Skip the reading because there is a pending read task"); + return null; + } + if (waitForCursorRewindingRefCnf > 0) { + log.info("Skip the reading due to new detected schema"); + return null; + } + if (state != Started) { + log.info("Skip the reading because producer has not started"); + return null; + } + int permits = getPermitsIfNoPendingRead(); + if (permits == 0) { + return null; + } + + int messagesToRead = permits; + long bytesToRead = -1; + if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { + AvailablePermits availablePermits = getRateLimiterAvailablePermits(permits); + if (!availablePermits.isReadable()) { + // no rate limiter permits from rate limit + log.debug() + .attr("messages", availablePermits.getMessages()) + .attr("bytes", availablePermits.getBytes()) + .log("Throttling replication traffic"); + return null; + } + messagesToRead = availablePermits.getMessages(); + bytesToRead = availablePermits.getBytes(); + if (!isWritable()) { + log.debug("Throttling replication traffic because producer is not writable"); + // Minimize the read size if the producer is disconnected or the window is already full + messagesToRead = 1; + } + } + return new ReadEntriesRequest( + createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead), bytesToRead); + } + } + @Override public void readEntriesComplete(List entries, Object ctx) { log.debug() @@ -871,26 +897,8 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } protected InFlightTask acquirePermitsIfNotFetchingSchema() { - synchronized (inFlightTasks) { - if (hasPendingRead()) { - log.info("Skip the reading because there is a pending read task"); - return null; - } - if (waitForCursorRewindingRefCnf > 0) { - log.info("Skip the reading due to new detected schema"); - return null; - } - if (state != Started) { - log.info("Skip the reading because producer has not started"); - return null; - } - // Guarantee that there is a unique cursor reading task. - int permits = getPermitsIfNoPendingRead(); - if (permits == 0) { - return null; - } - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), permits); - } + ReadEntriesRequest request = acquireReadEntriesRequestIfNeeded(); + return request == null ? null : request.inFlightTask; } protected int getPermitsIfNoPendingRead() { 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 478437ffde015..7c964fcd76cde 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 @@ -21,6 +21,7 @@ import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import java.util.ArrayList; @@ -28,6 +29,7 @@ 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; @@ -167,6 +169,39 @@ public void testReadEntriesFailedCompletesInFlightTaskAfterReplicatorTerminated( } } + @Test + public void testRateLimiterWithoutPermitsDoesNotCreateInFlightTask() throws Exception { + assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(0, -1); + assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(-1, 0); + } + + private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long availableMessages, + long availableBytes) throws Exception { + PersistentReplicator replicator = getReplicator(topicName); + + LinkedList inFlightTasks = replicator.inFlightTasks; + List originalTasks = new ArrayList<>(inFlightTasks); + Optional originalRateLimiter = replicator.dispatchRateLimiter; + inFlightTasks.clear(); + + 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); + + try { + Assert.assertNull(replicator.acquireReadEntriesRequestIfNeeded()); + Assert.assertTrue(inFlightTasks.isEmpty()); + Assert.assertFalse(replicator.hasPendingRead()); + assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); + } finally { + inFlightTasks.clear(); + inFlightTasks.addAll(originalTasks); + replicator.dispatchRateLimiter = originalRateLimiter; + } + } + @Test public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { log.info("Starting testCreateOrRecycleInFlightTaskIntoQueue"); From 91d578da9973fdef3cfa39a9f540cb512ff26147 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Thu, 18 Jun 2026 13:07:16 +0800 Subject: [PATCH 02/20] Address replicator rate limiter permit calculation --- .../service/persistent/PersistentReplicator.java | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) 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 5b1773b60829d..366900d31d8ac 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 @@ -335,14 +335,19 @@ ReadEntriesRequest acquireReadEntriesRequestIfNeeded() { return null; } int permits = getPermitsIfNoPendingRead(); - if (permits == 0) { + if (permits <= 0) { return null; } int messagesToRead = permits; long bytesToRead = -1; if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - AvailablePermits availablePermits = getRateLimiterAvailablePermits(permits); + if (!isWritable()) { + log.debug("Throttling replication traffic because producer is not writable"); + // Minimize the read size if the producer is disconnected or the window is already full + messagesToRead = 1; + } + AvailablePermits availablePermits = getRateLimiterAvailablePermits(messagesToRead); if (!availablePermits.isReadable()) { // no rate limiter permits from rate limit log.debug() @@ -353,11 +358,6 @@ ReadEntriesRequest acquireReadEntriesRequestIfNeeded() { } messagesToRead = availablePermits.getMessages(); bytesToRead = availablePermits.getBytes(); - if (!isWritable()) { - log.debug("Throttling replication traffic because producer is not writable"); - // Minimize the read size if the producer is disconnected or the window is already full - messagesToRead = 1; - } } return new ReadEntriesRequest( createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead), bytesToRead); From bf4999d57a6a45cdf73d6441086f52c97d28b8df Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Tue, 23 Jun 2026 09:19:29 +0800 Subject: [PATCH 03/20] Keep replicator read limits in in-flight task --- .../persistent/PersistentReplicator.java | 44 ++++++++----------- .../PersistentReplicatorInflightTaskTest.java | 11 +++-- 2 files changed, 26 insertions(+), 29 deletions(-) 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 366900d31d8ac..a572699fa2f96 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 @@ -240,17 +240,6 @@ public boolean isReadable() { } } - @VisibleForTesting - static class ReadEntriesRequest { - private final InFlightTask inFlightTask; - private final long bytesToRead; - - ReadEntriesRequest(InFlightTask inFlightTask, long bytesToRead) { - this.inFlightTask = inFlightTask; - this.bytesToRead = bytesToRead; - } - } - /** * Calculate available permits for read entries. */ @@ -298,8 +287,8 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - ReadEntriesRequest request = acquireReadEntriesRequestIfNeeded(); - if (request == null) { + InFlightTask newInFlightTask = acquireInFlightTaskIfNeeded(); + if (newInFlightTask == null) { // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); if (!hasPendingRead()) { @@ -310,17 +299,16 @@ protected void readMoreEntries() { return; } } - InFlightTask newInFlightTask = request.inFlightTask; log.debug() .attr("readingEntries", newInFlightTask.readingEntries) - .attr("bytesToRead", request.bytesToRead) + .attr("bytesToRead", newInFlightTask.bytesToRead) .log("Scheduling read"); - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, request.bytesToRead, this, + cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, newInFlightTask.bytesToRead, this, newInFlightTask/* Context object */, topic.getMaxReadPosition()); } @VisibleForTesting - ReadEntriesRequest acquireReadEntriesRequestIfNeeded() { + InFlightTask acquireInFlightTaskIfNeeded() { synchronized (inFlightTasks) { if (hasPendingRead()) { log.info("Skip the reading because there is a pending read task"); @@ -359,8 +347,7 @@ ReadEntriesRequest acquireReadEntriesRequestIfNeeded() { messagesToRead = availablePermits.getMessages(); bytesToRead = availablePermits.getBytes(); } - return new ReadEntriesRequest( - createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead), bytesToRead); + return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead, bytesToRead); } } @@ -822,6 +809,7 @@ public ManagedCursor getCursor() { protected static class InFlightTask { Position readPos; int readingEntries; + long bytesToRead; volatile List entries; volatile int completedEntries; volatile boolean skipReadResultDueToCursorRewind; @@ -838,17 +826,23 @@ public synchronized void incCompletedEntries() { } } - synchronized void recycle(Position readStart, int readingEntries) { + synchronized void recycle(Position readStart, int readingEntries, long bytesToRead) { this.readPos = readStart; this.readingEntries = readingEntries; + this.bytesToRead = bytesToRead; this.entries = null; this.completedEntries = 0; this.skipReadResultDueToCursorRewind = false; } public InFlightTask(Position readPos, int readingEntries, String replicatorId) { + this(readPos, readingEntries, -1, replicatorId); + } + + public InFlightTask(Position readPos, int readingEntries, long bytesToRead, String replicatorId) { this.readPos = readPos; this.readingEntries = readingEntries; + this.bytesToRead = bytesToRead; this.replicatorId = replicatorId; } @@ -868,6 +862,7 @@ public String toString() { + "{replicatorId=" + replicatorId + ", readPos=" + readPos + ", readingEntries=" + readingEntries + + ", bytesToRead=" + bytesToRead + ", readoutEntries=" + (entries == null ? "-1" : entries.size()) + ", completedEntries=" + completedEntries + ", skipReadResultDueToCursorRewound=" + skipReadResultDueToCursorRewind @@ -876,7 +871,7 @@ public String toString() { } @VisibleForTesting - InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries) { + InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries, long bytesToRead) { synchronized (inFlightTasks) { // Reuse projects that has done. if (!inFlightTasks.isEmpty()) { @@ -884,21 +879,20 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE if (first.isDone()) { // Remove from the first index, and add to the latest index. inFlightTasks.poll(); - first.recycle(readPos, readingEntries); + first.recycle(readPos, readingEntries, bytesToRead); inFlightTasks.add(first); return first; } } // New project if nothing can be reused. - InFlightTask task = new InFlightTask(readPos, readingEntries, replicatorId); + InFlightTask task = new InFlightTask(readPos, readingEntries, bytesToRead, replicatorId); inFlightTasks.add(task); return task; } } protected InFlightTask acquirePermitsIfNotFetchingSchema() { - ReadEntriesRequest request = acquireReadEntriesRequestIfNeeded(); - return request == null ? null : request.inFlightTask; + return acquireInFlightTaskIfNeeded(); } protected int getPermitsIfNoPendingRead() { 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 7c964fcd76cde..2b4eb006f15e6 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 @@ -191,7 +191,7 @@ private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long avail replicator.dispatchRateLimiter = Optional.of(rateLimiter); try { - Assert.assertNull(replicator.acquireReadEntriesRequestIfNeeded()); + Assert.assertNull(replicator.acquireInFlightTaskIfNeeded()); Assert.assertTrue(inFlightTasks.isEmpty()); Assert.assertFalse(replicator.hasPendingRead()); assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); @@ -221,35 +221,38 @@ public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { // Test Case 1: Create a new task when the queue is empty Position position1 = PositionFactory.create(1, 1); Assert.assertNotNull(position1, "Position should not be null"); - InFlightTask task1 = replicator.createOrRecycleInFlightTaskIntoQueue(position1, 10); + InFlightTask task1 = replicator.createOrRecycleInFlightTaskIntoQueue(position1, 10, -1); // Verify a new task was created and added to the queue Assert.assertNotNull(task1, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 1, "Queue should have one task"); Assert.assertEquals(task1.getReadPos(), position1, "Task should have the correct position"); Assert.assertEquals(task1.getReadingEntries(), 10, "Task should have the correct reading entries count"); + Assert.assertEquals(task1.getBytesToRead(), -1, "Task should have the correct byte read limit"); // Mark the task as done to test recycling task1.setEntries(Collections.emptyList()); // Test Case 2: Recycle an existing task Position position2 = PositionFactory.create(2, 2); Assert.assertNotNull(position2, "Position should not be null"); - InFlightTask task2 = replicator.createOrRecycleInFlightTaskIntoQueue(position2, 20); + InFlightTask task2 = replicator.createOrRecycleInFlightTaskIntoQueue(position2, 20, 1024); // Verify the task was recycled Assert.assertNotNull(task2, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 1, "Queue should still have one task"); Assert.assertEquals(task2.getReadPos(), position2, "Task should have the updated position"); Assert.assertEquals(task2.getReadingEntries(), 20, "Task should have the updated reading entries count"); + Assert.assertEquals(task2.getBytesToRead(), 1024, "Task should have the updated byte read limit"); // Test Case 3: Create a new task when no tasks can be recycled task2.setEntries(null); // Make the task not done Position position3 = PositionFactory.create(3, 3); Assert.assertNotNull(position3, "Position should not be null"); - InFlightTask task3 = replicator.createOrRecycleInFlightTaskIntoQueue(position3, 30); + InFlightTask task3 = replicator.createOrRecycleInFlightTaskIntoQueue(position3, 30, 2048); // Verify a new task was created Assert.assertNotNull(task3, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 2, "Queue should have two tasks"); Assert.assertEquals(task3.getReadPos(), position3, "Task should have the correct position"); Assert.assertEquals(task3.getReadingEntries(), 30, "Task should have the correct reading entries count"); + Assert.assertEquals(task3.getBytesToRead(), 2048, "Task should have the correct byte read limit"); // cleanup. log.info("Completed testCreateOrRecycleInFlightTaskIntoQueue"); From f583b3dc948bb9f39ffe282919bc295584a1685b Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Tue, 23 Jun 2026 21:09:40 +0800 Subject: [PATCH 04/20] Address replicator permit acquisition review --- .../persistent/PersistentReplicator.java | 87 +++++++++---------- .../PersistentReplicatorInflightTaskTest.java | 2 +- 2 files changed, 42 insertions(+), 47 deletions(-) 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 ba264f9070c08..d802f7773389f 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 @@ -287,7 +287,7 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - InFlightTask newInFlightTask = acquireInFlightTaskIfNeeded(); + InFlightTask newInFlightTask = acquirePermitsIfNotFetchingSchema(); if (newInFlightTask == null) { // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); @@ -307,50 +307,6 @@ protected void readMoreEntries() { newInFlightTask/* Context object */, topic.getMaxReadPosition()); } - @VisibleForTesting - InFlightTask acquireInFlightTaskIfNeeded() { - synchronized (inFlightTasks) { - if (hasPendingRead()) { - log.info("Skip the reading because there is a pending read task"); - return null; - } - if (waitForCursorRewindingRefCnf > 0) { - log.info("Skip the reading due to new detected schema"); - return null; - } - if (state != Started) { - log.info("Skip the reading because producer has not started"); - return null; - } - int permits = getPermitsIfNoPendingRead(); - if (permits <= 0) { - return null; - } - - int messagesToRead = permits; - long bytesToRead = -1; - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!isWritable()) { - log.debug("Throttling replication traffic because producer is not writable"); - // Minimize the read size if the producer is disconnected or the window is already full - messagesToRead = 1; - } - AvailablePermits availablePermits = getRateLimiterAvailablePermits(messagesToRead); - if (!availablePermits.isReadable()) { - // no rate limiter permits from rate limit - log.debug() - .attr("messages", availablePermits.getMessages()) - .attr("bytes", availablePermits.getBytes()) - .log("Throttling replication traffic"); - return null; - } - messagesToRead = availablePermits.getMessages(); - bytesToRead = availablePermits.getBytes(); - } - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead, bytesToRead); - } - } - @Override public void readEntriesComplete(List entries, Object ctx) { log.debug() @@ -895,7 +851,46 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } protected InFlightTask acquirePermitsIfNotFetchingSchema() { - return acquireInFlightTaskIfNeeded(); + synchronized (inFlightTasks) { + if (hasPendingRead()) { + log.info("Skip the reading because there is a pending read task"); + return null; + } + if (waitForCursorRewindingRefCnf > 0) { + log.info("Skip the reading due to new detected schema"); + return null; + } + if (state != Started) { + log.info("Skip the reading because producer has not started"); + return null; + } + int permits = getPermitsIfNoPendingRead(); + if (permits <= 0) { + return null; + } + + int messagesToRead = permits; + long bytesToRead = -1; + if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { + if (!isWritable()) { + log.debug("Throttling replication traffic because producer is not writable"); + // Minimize the read size if the producer is disconnected or the window is already full + messagesToRead = 1; + } + AvailablePermits availablePermits = getRateLimiterAvailablePermits(messagesToRead); + if (!availablePermits.isReadable()) { + // no rate limiter permits from rate limit + log.debug() + .attr("messages", availablePermits.getMessages()) + .attr("bytes", availablePermits.getBytes()) + .log("Throttling replication traffic"); + return null; + } + messagesToRead = availablePermits.getMessages(); + bytesToRead = availablePermits.getBytes(); + } + return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead, bytesToRead); + } } protected int getPermitsIfNoPendingRead() { 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 7f3f82a22cd2b..5467f77daacc6 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 @@ -195,7 +195,7 @@ private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long avail replicator.dispatchRateLimiter = Optional.of(rateLimiter); try { - Assert.assertNull(replicator.acquireInFlightTaskIfNeeded()); + Assert.assertNull(replicator.acquirePermitsIfNotFetchingSchema()); Assert.assertTrue(inFlightTasks.isEmpty()); Assert.assertFalse(replicator.hasPendingRead()); assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); From 6420f7130ef0495f58016cac01925c34e75ec79a Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Wed, 24 Jun 2026 10:28:21 +0800 Subject: [PATCH 05/20] Rename replicator in-flight read task helper --- .../persistent/GeoPersistentReplicator.java | 2 +- .../persistent/PersistentReplicator.java | 4 ++-- .../PersistentReplicatorInflightTaskTest.java | 20 +++++++++---------- 3 files changed, 13 insertions(+), 13 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 c6534b346c25b..24afda1ee7612 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 @@ -264,7 +264,7 @@ protected boolean replicateEntries(List entries, final InFlightTask inFli * Explain the result of the race-condition between: * - {@link #readMoreEntries} * - {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} - * Since {@link #acquirePermitsIfNotFetchingSchema} and + * Since {@link #tryReserveInFlightReadTask} and * {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} acquire the * same lock, it is safe. */ 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 d802f7773389f..1d94cc678ff21 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 @@ -287,7 +287,7 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - InFlightTask newInFlightTask = acquirePermitsIfNotFetchingSchema(); + InFlightTask newInFlightTask = tryReserveInFlightReadTask(); if (newInFlightTask == null) { // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); @@ -850,7 +850,7 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } } - protected InFlightTask acquirePermitsIfNotFetchingSchema() { + protected InFlightTask tryReserveInFlightReadTask() { synchronized (inFlightTasks) { if (hasPendingRead()) { log.info("Skip the reading because there is a pending read task"); 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 5467f77daacc6..63a7bba68968d 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 @@ -195,7 +195,7 @@ private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long avail replicator.dispatchRateLimiter = Optional.of(rateLimiter); try { - Assert.assertNull(replicator.acquirePermitsIfNotFetchingSchema()); + Assert.assertNull(replicator.tryReserveInFlightReadTask()); Assert.assertTrue(inFlightTasks.isEmpty()); Assert.assertFalse(replicator.hasPendingRead()); assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); @@ -428,8 +428,8 @@ public void testGetPermitsIfNoPendingRead() throws Exception { } @Test - public void testAcquirePermitsIfNotFetchingSchema() throws Exception { - log.info("Starting testAcquirePermitsIfNotFetchingSchema"); + public void testTryReserveInFlightReadTask() throws Exception { + log.info("Starting testTryReserveInFlightReadTask"); // Get the replicator for the test topic PersistentReplicator replicator = getReplicator(topicName); Assert.assertNotNull(replicator, "Replicator should not be null"); @@ -452,7 +452,7 @@ public void testAcquirePermitsIfNotFetchingSchema() throws Exception { // 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(); + InFlightTask task1 = replicator.tryReserveInFlightReadTask(); 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, @@ -466,13 +466,13 @@ public void testAcquirePermitsIfNotFetchingSchema() throws Exception { InFlightTask pendingReadTask = new InFlightTask(position1, 5, ""); // Don't set readoutEntries to simulate pending read inFlightTasks.add(pendingReadTask); - InFlightTask task2 = replicator.acquirePermitsIfNotFetchingSchema(); + InFlightTask task2 = replicator.tryReserveInFlightReadTask(); Assert.assertNull(task2, "Should return null when there is a pending read"); // Test Case 3: With waitForCursorRewinding=true - should return null inFlightTasks.clear(); replicator.waitForCursorRewindingRefCnf = 1; - InFlightTask task3 = replicator.acquirePermitsIfNotFetchingSchema(); + InFlightTask task3 = replicator.tryReserveInFlightReadTask(); Assert.assertNull(task3, "Should return null when waiting for cursor rewinding"); // Reset for next test replicator.waitForCursorRewindingRefCnf = 0; @@ -480,7 +480,7 @@ public void testAcquirePermitsIfNotFetchingSchema() throws Exception { // 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(); + InFlightTask task4 = replicator.tryReserveInFlightReadTask(); Assert.assertNull(task4, "Should return null when state is not Started"); // Reset state for next test BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, AbstractReplicator.State.Started); @@ -503,7 +503,7 @@ public void testAcquirePermitsIfNotFetchingSchema() throws Exception { 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(); + InFlightTask task5 = replicator.tryReserveInFlightReadTask(); 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"); @@ -522,9 +522,9 @@ public void testAcquirePermitsIfNotFetchingSchema() throws Exception { task.setEntries(entries); inFlightTasks.add(task); } - InFlightTask task6 = replicator.acquirePermitsIfNotFetchingSchema(); + InFlightTask task6 = replicator.tryReserveInFlightReadTask(); Assert.assertNull(task6, "Should return null when permits is 0"); - log.info("Completed testAcquirePermitsIfNotFetchingSchema"); + log.info("Completed testTryReserveInFlightReadTask"); } finally { // Restore original state replicator.waitForCursorRewindingRefCnf = originalWaitForCursorRewinding; From 7885d45a8697f24c4bf5334e5a4d398120f20100 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Thu, 25 Jun 2026 00:09:31 +0800 Subject: [PATCH 06/20] Refine in-flight read task helper name --- .../persistent/GeoPersistentReplicator.java | 2 +- .../persistent/PersistentReplicator.java | 5 +++-- .../PersistentReplicatorInflightTaskTest.java | 20 +++++++++---------- 3 files changed, 14 insertions(+), 13 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 24afda1ee7612..ae1576a57d3ee 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 @@ -264,7 +264,7 @@ protected boolean replicateEntries(List entries, final InFlightTask inFli * Explain the result of the race-condition between: * - {@link #readMoreEntries} * - {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} - * Since {@link #tryReserveInFlightReadTask} and + * Since {@link #maybeCreateInFlightReadTask} and * {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} acquire the * same lock, it is safe. */ 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 1d94cc678ff21..732ba61d5d1fb 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 @@ -287,7 +287,7 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - InFlightTask newInFlightTask = tryReserveInFlightReadTask(); + InFlightTask newInFlightTask = maybeCreateInFlightReadTask(); if (newInFlightTask == null) { // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); @@ -850,7 +850,8 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } } - protected InFlightTask tryReserveInFlightReadTask() { + @VisibleForTesting + InFlightTask maybeCreateInFlightReadTask() { synchronized (inFlightTasks) { if (hasPendingRead()) { log.info("Skip the reading because there is a pending read task"); 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 63a7bba68968d..98277d1422182 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 @@ -195,7 +195,7 @@ private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long avail replicator.dispatchRateLimiter = Optional.of(rateLimiter); try { - Assert.assertNull(replicator.tryReserveInFlightReadTask()); + Assert.assertNull(replicator.maybeCreateInFlightReadTask()); Assert.assertTrue(inFlightTasks.isEmpty()); Assert.assertFalse(replicator.hasPendingRead()); assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); @@ -428,8 +428,8 @@ public void testGetPermitsIfNoPendingRead() throws Exception { } @Test - public void testTryReserveInFlightReadTask() throws Exception { - log.info("Starting testTryReserveInFlightReadTask"); + public void testMaybeCreateInFlightReadTask() throws Exception { + log.info("Starting testMaybeCreateInFlightReadTask"); // Get the replicator for the test topic PersistentReplicator replicator = getReplicator(topicName); Assert.assertNotNull(replicator, "Replicator should not be null"); @@ -452,7 +452,7 @@ public void testTryReserveInFlightReadTask() throws Exception { // First, check the current permits available int expectedPermits = replicator.getPermitsIfNoPendingRead(); Assert.assertTrue(expectedPermits > 0, "Should have available permits for the test"); - InFlightTask task1 = replicator.tryReserveInFlightReadTask(); + InFlightTask task1 = replicator.maybeCreateInFlightReadTask(); 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, @@ -466,13 +466,13 @@ public void testTryReserveInFlightReadTask() throws Exception { InFlightTask pendingReadTask = new InFlightTask(position1, 5, ""); // Don't set readoutEntries to simulate pending read inFlightTasks.add(pendingReadTask); - InFlightTask task2 = replicator.tryReserveInFlightReadTask(); + InFlightTask task2 = replicator.maybeCreateInFlightReadTask(); Assert.assertNull(task2, "Should return null when there is a pending read"); // Test Case 3: With waitForCursorRewinding=true - should return null inFlightTasks.clear(); replicator.waitForCursorRewindingRefCnf = 1; - InFlightTask task3 = replicator.tryReserveInFlightReadTask(); + InFlightTask task3 = replicator.maybeCreateInFlightReadTask(); Assert.assertNull(task3, "Should return null when waiting for cursor rewinding"); // Reset for next test replicator.waitForCursorRewindingRefCnf = 0; @@ -480,7 +480,7 @@ public void testTryReserveInFlightReadTask() throws Exception { // 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.tryReserveInFlightReadTask(); + InFlightTask task4 = replicator.maybeCreateInFlightReadTask(); Assert.assertNull(task4, "Should return null when state is not Started"); // Reset state for next test BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, AbstractReplicator.State.Started); @@ -503,7 +503,7 @@ public void testTryReserveInFlightReadTask() throws Exception { 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.tryReserveInFlightReadTask(); + InFlightTask task5 = replicator.maybeCreateInFlightReadTask(); 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"); @@ -522,9 +522,9 @@ public void testTryReserveInFlightReadTask() throws Exception { task.setEntries(entries); inFlightTasks.add(task); } - InFlightTask task6 = replicator.tryReserveInFlightReadTask(); + InFlightTask task6 = replicator.maybeCreateInFlightReadTask(); Assert.assertNull(task6, "Should return null when permits is 0"); - log.info("Completed testTryReserveInFlightReadTask"); + log.info("Completed testMaybeCreateInFlightReadTask"); } finally { // Restore original state replicator.waitForCursorRewindingRefCnf = originalWaitForCursorRewinding; From ce540eea2316f403199699ee6c4ad2a4db598704 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 24 Jun 2026 22:47:47 +0300 Subject: [PATCH 07/20] Resolve merge conflict --- .../persistent/PersistentReplicatorInflightTaskTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 8f67570d6eeb4..98479060280e7 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 @@ -541,7 +541,8 @@ public static Runnable pauseReplicator(PersistentReplicator replicator) { }); 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, -1); return () -> { inFlightTask.setEntries(Collections.emptyList()); replicator.readMoreEntries(); From 5ac3d8492f0b8d497ee138e188984237ecb2168b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 00:44:35 +0300 Subject: [PATCH 08/20] Reduce duplication --- .../persistent/PersistentReplicator.java | 38 +++++++++---------- 1 file changed, 18 insertions(+), 20 deletions(-) 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 322bda5a3a712..4bec0483cb852 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 @@ -935,27 +935,25 @@ InFlightTask maybeCreateInFlightReadTask() { return null; } - int messagesToRead = permits; - long bytesToRead = -1; - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!isWritable()) { - log.debug("Throttling replication traffic because producer is not writable"); - // Minimize the read size if the producer is disconnected or the window is already full - messagesToRead = 1; - } - AvailablePermits availablePermits = getRateLimiterAvailablePermits(messagesToRead); - if (!availablePermits.isReadable()) { - // no rate limiter permits from rate limit - log.debug() - .attr("messages", availablePermits.getMessages()) - .attr("bytes", availablePermits.getBytes()) - .log("Throttling replication traffic"); - return null; - } - messagesToRead = availablePermits.getMessages(); - bytesToRead = availablePermits.getBytes(); + if (!isWritable()) { + log.debug("Throttling replication traffic because producer is not writable"); + // Minimize the read size if the producer is disconnected or the window is already full + permits = 1; + } + + AvailablePermits availablePermits = getRateLimiterAvailablePermits(permits); + + if (!availablePermits.isReadable()) { + // no rate limiter permits from rate limit + log.debug() + .attr("messages", availablePermits.getMessages()) + .attr("bytes", availablePermits.getBytes()) + .log("Throttling replication traffic"); + return null; } - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), messagesToRead, bytesToRead); + + return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), availablePermits.getMessages(), + availablePermits.getBytes()); } } From e0273556ce7619cc7af48000f62a627d96c74524 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 01:04:12 +0300 Subject: [PATCH 09/20] improve logging and comments --- .../service/persistent/PersistentReplicator.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) 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 4bec0483cb852..1f63653e0c5e6 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 @@ -259,15 +259,17 @@ private AvailablePermits getRateLimiterAvailablePermits(int availablePermits) { if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { DispatchRateLimiter rateLimiter = dispatchRateLimiter.get(); // if dispatch-rate is in msg then read only msg according to available permit + // rateLimiter returns -1 if there is no rate limit configured availablePermitsOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); availablePermitsOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); - // no permits from rate limit + // no permits from rate limit when either limit is 0 if (availablePermitsOnByte == 0 || availablePermitsOnMsg == 0) { log.debug() .attr("dispatchRateOnMsg", rateLimiter.getDispatchRateOnMsg()) .attr("dispatchRateOnByte", rateLimiter.getDispatchRateOnByte()) - .attr("backoffMs", MESSAGE_RATE_BACKOFF_MS) - .log("Message-read exceeded topic replicator message-rate, scheduling after a delay"); + .attr("availablePermitsOnMsg", availablePermitsOnMsg) + .attr("availablePermitsOnByte", availablePermitsOnByte) + .log("Message-read exceeded topic replicator rate limit"); return new AvailablePermits(-1, -1); } } @@ -936,7 +938,7 @@ InFlightTask maybeCreateInFlightReadTask() { } if (!isWritable()) { - log.debug("Throttling replication traffic because producer is not writable"); + 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; } From a5987d9d20432eec84b3f21a93a5b929b74ff8dc Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 01:23:12 +0300 Subject: [PATCH 10/20] Rename AvailablePermits to ReadLimits since the previous name was overlapping with "permits" --- .../persistent/PersistentReplicator.java | 63 +++++++------------ 1 file changed, 23 insertions(+), 40 deletions(-) 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 1f63653e0c5e6..e04282f55fcc6 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 @@ -39,7 +39,6 @@ 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; @@ -218,70 +217,54 @@ protected void disableReplicatorRead() { 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 configured max read batch size and max read size. */ - private AvailablePermits getRateLimiterAvailablePermits(int availablePermits) { + private ReadLimits getReadLimits(int availablePermits) { // return 0, if Producer queue is full, it will pause read entries. if (availablePermits <= 0) { log.debug() .attr("availablePermits", availablePermits) .log("Producer queue is full, pausing reads"); - return new AvailablePermits(0, 0); + return new ReadLimits(0, 0); } - long availablePermitsOnMsg = -1; - long availablePermitsOnByte = -1; + long readLimitOnMsg = -1; + long readLimitOnByte = -1; // 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 // rateLimiter returns -1 if there is no rate limit configured - availablePermitsOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); - availablePermitsOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); + readLimitOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); + readLimitOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); // no permits from rate limit when either limit is 0 - if (availablePermitsOnByte == 0 || availablePermitsOnMsg == 0) { + if (readLimitOnByte == 0 || readLimitOnMsg == 0) { log.debug() .attr("dispatchRateOnMsg", rateLimiter.getDispatchRateOnMsg()) .attr("dispatchRateOnByte", rateLimiter.getDispatchRateOnByte()) - .attr("availablePermitsOnMsg", availablePermitsOnMsg) - .attr("availablePermitsOnByte", availablePermitsOnByte) + .attr("readLimitOnMsg", readLimitOnMsg) + .attr("readLimitOnByte", readLimitOnByte) .log("Message-read exceeded topic replicator rate limit"); - return new AvailablePermits(-1, -1); + return new ReadLimits(-1, -1); } } - availablePermitsOnMsg = - availablePermitsOnMsg == -1 ? availablePermits : Math.min(availablePermits, availablePermitsOnMsg); - availablePermitsOnMsg = Math.min(availablePermitsOnMsg, readBatchSize); + readLimitOnMsg = readLimitOnMsg == -1 ? availablePermits : Math.min(availablePermits, readLimitOnMsg); + readLimitOnMsg = Math.min(readLimitOnMsg, readBatchSize); - availablePermitsOnByte = - availablePermitsOnByte == -1 ? readMaxSizeBytes : Math.min(readMaxSizeBytes, availablePermitsOnByte); + readLimitOnByte = readLimitOnByte == -1 ? readMaxSizeBytes : Math.min(readMaxSizeBytes, readLimitOnByte); - return new AvailablePermits((int) availablePermitsOnMsg, availablePermitsOnByte); + return new ReadLimits((int) readLimitOnMsg, readLimitOnByte); } public void disconnectIfNoTrafficAndBacklog() { @@ -943,19 +926,19 @@ InFlightTask maybeCreateInFlightReadTask() { permits = 1; } - AvailablePermits availablePermits = getRateLimiterAvailablePermits(permits); + ReadLimits readLimits = getReadLimits(permits); - if (!availablePermits.isReadable()) { + if (!readLimits.isReadable()) { // no rate limiter permits from rate limit log.debug() - .attr("messages", availablePermits.getMessages()) - .attr("bytes", availablePermits.getBytes()) + .attr("messages", readLimits.messages) + .attr("bytes", readLimits.bytes) .log("Throttling replication traffic"); return null; } - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), availablePermits.getMessages(), - availablePermits.getBytes()); + return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), readLimits.messages, + readLimits.bytes); } } From b6a10c7a9b570bb4e4aba19bffb9fdc80899689b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 01:44:09 +0300 Subject: [PATCH 11/20] Fix readBatchSize logic and improve readability --- .../persistent/PersistentReplicator.java | 31 ++++++++++++------- 1 file changed, 19 insertions(+), 12 deletions(-) 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 e04282f55fcc6..0b5ef880eb160 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 @@ -92,7 +92,7 @@ public abstract class PersistentReplicator extends AbstractReplicator protected Optional dispatchRateLimiter = Optional.empty(); private final Object dispatchRateLimiterLock = new Object(); - private int readBatchSize; + private volatile int readBatchSize; private final int readMaxSizeBytes; private final int producerQueueThreshold; @@ -140,9 +140,7 @@ public PersistentReplicator(String localCluster, PersistentTopic localTopic, Man this.expiryMonitor = new PersistentMessageExpiryMonitor(localTopic, Codec.decode(cursor.getName()), cursor, null); - readBatchSize = Math.min( - producerQueueSize, - localTopic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadBatchSize()); + readBatchSize = getMaxReadBatchSize(); readMaxSizeBytes = localTopic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadSizeBytes(); producerQueueThreshold = (int) (producerQueueSize * 0.9); @@ -151,6 +149,10 @@ public PersistentReplicator(String localCluster, PersistentTopic localTopic, Man startProducer(); } + private int getMaxReadBatchSize() { + return Math.min(producerQueueSize, brokerService.pulsar().getConfiguration().getDispatcherMaxReadBatchSize()); + } + @Override protected void setProducerAndTriggerReadEntries(Producer producer) { /** @@ -225,7 +227,7 @@ public boolean isReadable() { /** * Calculate read limits for a read operation. Takes the rate limiter into account if it's enabled. - * Also limits to configured max read batch size and max read size. + * Also limits to current readBatchSize and readMaxSizeBytes. */ private ReadLimits getReadLimits(int availablePermits) { @@ -237,8 +239,8 @@ private ReadLimits getReadLimits(int availablePermits) { return new ReadLimits(0, 0); } - long readLimitOnMsg = -1; - long readLimitOnByte = -1; + long readLimitOnMsg; + long readLimitOnByte; // handle rate limit if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { @@ -257,13 +259,18 @@ private ReadLimits getReadLimits(int availablePermits) { .log("Message-read exceeded topic replicator rate limit"); return new ReadLimits(-1, -1); } + // use available permits if no rate limit configured, otherwise limit to returned rate limiter permits + readLimitOnMsg = readLimitOnMsg == -1 ? availablePermits : Math.min(availablePermits, 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 = availablePermits; + readLimitOnByte = readMaxSizeBytes; } - readLimitOnMsg = readLimitOnMsg == -1 ? availablePermits : Math.min(availablePermits, readLimitOnMsg); + // limit messages to current read batch size readLimitOnMsg = Math.min(readLimitOnMsg, readBatchSize); - readLimitOnByte = readLimitOnByte == -1 ? readMaxSizeBytes : Math.min(readMaxSizeBytes, readLimitOnByte); - return new ReadLimits((int) readLimitOnMsg, readLimitOnByte); } @@ -366,7 +373,7 @@ public void readEntriesComplete(List entries, Object ctx) { 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); log.debug() @@ -523,7 +530,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { } // 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(); From a4eff4b30a41b4fbe3a846ea17560e7a0bbecca9 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 02:19:28 +0300 Subject: [PATCH 12/20] Use a single concept "permits" instead of a separate "availablePermits" --- .../service/persistent/PersistentReplicator.java | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) 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 0b5ef880eb160..4c635d10abeaa 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 @@ -229,12 +229,12 @@ public boolean isReadable() { * 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 ReadLimits getReadLimits(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) { log.debug() - .attr("availablePermits", availablePermits) + .attr("permits", permits) .log("Producer queue is full, pausing reads"); return new ReadLimits(0, 0); } @@ -260,11 +260,11 @@ private ReadLimits getReadLimits(int availablePermits) { return new ReadLimits(-1, -1); } // use available permits if no rate limit configured, otherwise limit to returned rate limiter permits - readLimitOnMsg = readLimitOnMsg == -1 ? availablePermits : Math.min(availablePermits, readLimitOnMsg); + 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 = availablePermits; + readLimitOnMsg = permits; readLimitOnByte = readMaxSizeBytes; } @@ -309,10 +309,8 @@ protected void readMoreEntries() { if (!hasPendingRead()) { topic.getBrokerService().executor().schedule( () -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); - return; - } else { - return; } + return; } log.debug() .attr("readingEntries", newInFlightTask.readingEntries) From c2446247beb702f7faaf428ce55de2944179c95e Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 02:21:54 +0300 Subject: [PATCH 13/20] Fix test to take readBatchSize into account --- .../persistent/PersistentReplicatorInflightTaskTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 98479060280e7..7498fd3ca8bb5 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 @@ -455,8 +455,8 @@ public void testMaybeCreateInFlightReadTask() throws Exception { InFlightTask task1 = replicator.maybeCreateInFlightReadTask(); 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.assertEquals(task1.getReadingEntries(), 100, + "Task readingEntries should equal to readBatchSize"); Assert.assertTrue(inFlightTasks.contains(task1), "Task should be added to the inFlightTasks list"); From b3bc95afb8a3fd7b416ea6b816bed8f7953ca386 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 02:31:36 +0300 Subject: [PATCH 14/20] Update comment --- .../pulsar/broker/service/persistent/PersistentReplicator.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 4c635d10abeaa..17d2ce7483702 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 @@ -245,7 +245,6 @@ private ReadLimits getReadLimits(int permits) { // 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 // rateLimiter returns -1 if there is no rate limit configured readLimitOnMsg = rateLimiter.getAvailableDispatchRateLimitOnMsg(); readLimitOnByte = rateLimiter.getAvailableDispatchRateLimitOnByte(); @@ -259,7 +258,7 @@ private ReadLimits getReadLimits(int permits) { .log("Message-read exceeded topic replicator rate limit"); return new ReadLimits(-1, -1); } - // use available permits if no rate limit configured, otherwise limit to returned rate limiter permits + // 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); From 8b7c5264ef9eb82d9cec230f144443df5387daf8 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 25 Jun 2026 02:32:33 +0300 Subject: [PATCH 15/20] Improve consistency --- .../pulsar/broker/service/persistent/PersistentReplicator.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 17d2ce7483702..981d26ecced10 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 @@ -141,7 +141,7 @@ public PersistentReplicator(String localCluster, PersistentTopic localTopic, Man Codec.decode(cursor.getName()), cursor, null); readBatchSize = getMaxReadBatchSize(); - readMaxSizeBytes = localTopic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadSizeBytes(); + readMaxSizeBytes = brokerService.pulsar().getConfiguration().getDispatcherMaxReadSizeBytes(); producerQueueThreshold = (int) (producerQueueSize * 0.9); this.initializeDispatchRateLimiterIfNeeded(); From 733affb52f863b8d1fd38b8c0eb896baaf266f4b Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Thu, 25 Jun 2026 22:54:33 +0800 Subject: [PATCH 16/20] Address remaining replicator review comments --- .../persistent/PersistentReplicator.java | 22 +++---- .../broker/service/OneWayReplicatorTest.java | 61 +++++++++++++++++++ .../PersistentReplicatorInflightTaskTest.java | 22 ++++--- 3 files changed, 85 insertions(+), 20 deletions(-) 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 981d26ecced10..8c64798fe45f7 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 @@ -313,9 +313,9 @@ protected void readMoreEntries() { } log.debug() .attr("readingEntries", newInFlightTask.readingEntries) - .attr("bytesToRead", newInFlightTask.bytesToRead) + .attr("maxBytesToRead", newInFlightTask.maxBytesToRead) .log("Scheduling read"); - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, newInFlightTask.bytesToRead, this, + cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, newInFlightTask.maxBytesToRead, this, newInFlightTask/* Context object */, topic.getMaxReadPosition()); } @@ -822,7 +822,7 @@ public ManagedCursor getCursor() { protected static class InFlightTask { Position readPos; int readingEntries; - long bytesToRead; + long maxBytesToRead; volatile List entries; volatile int completedEntries; volatile boolean skipReadResultDueToCursorRewind; @@ -839,10 +839,10 @@ public synchronized void incCompletedEntries() { } } - synchronized void recycle(Position readStart, int readingEntries, long bytesToRead) { + synchronized void recycle(Position readStart, int readingEntries, long maxBytesToRead) { this.readPos = readStart; this.readingEntries = readingEntries; - this.bytesToRead = bytesToRead; + this.maxBytesToRead = maxBytesToRead; this.entries = null; this.completedEntries = 0; this.skipReadResultDueToCursorRewind = false; @@ -852,10 +852,10 @@ public InFlightTask(Position readPos, int readingEntries, String replicatorId) { this(readPos, readingEntries, -1, replicatorId); } - public InFlightTask(Position readPos, int readingEntries, long bytesToRead, String replicatorId) { + public InFlightTask(Position readPos, int readingEntries, long maxBytesToRead, String replicatorId) { this.readPos = readPos; this.readingEntries = readingEntries; - this.bytesToRead = bytesToRead; + this.maxBytesToRead = maxBytesToRead; this.replicatorId = replicatorId; } @@ -875,7 +875,7 @@ public String toString() { + "{replicatorId=" + replicatorId + ", readPos=" + readPos + ", readingEntries=" + readingEntries - + ", bytesToRead=" + bytesToRead + + ", maxBytesToRead=" + maxBytesToRead + ", readoutEntries=" + (entries == null ? "-1" : entries.size()) + ", completedEntries=" + completedEntries + ", skipReadResultDueToCursorRewound=" + skipReadResultDueToCursorRewind @@ -884,7 +884,7 @@ public String toString() { } @VisibleForTesting - InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries, long bytesToRead) { + InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries, long maxBytesToRead) { synchronized (inFlightTasks) { // Reuse projects that has done. if (!inFlightTasks.isEmpty()) { @@ -892,13 +892,13 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE if (first.isDone()) { // Remove from the first index, and add to the latest index. inFlightTasks.poll(); - first.recycle(readPos, readingEntries, bytesToRead); + first.recycle(readPos, readingEntries, maxBytesToRead); inFlightTasks.add(first); return first; } } // New project if nothing can be reused. - InFlightTask task = new InFlightTask(readPos, readingEntries, bytesToRead, replicatorId); + InFlightTask task = new InFlightTask(readPos, readingEntries, maxBytesToRead, replicatorId); inFlightTasks.add(task); return task; } 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 b37fda1fa29f9..d86732aee7913 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 @@ -119,6 +119,7 @@ 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; @@ -2431,6 +2432,66 @@ public void testReplicatorsInflightTaskListIsEmptyAfterReplicationFinished() thr ensureNoBacklogByInflightTask(getReplicator(topicName)); } + @Test(timeOut = 90_000) + public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Exception { + final String topicName = BrokerTestUtil.newUniqueName("persistent://" + replicatedNamespace + "/tp_"); + final String subscriptionName = "sub"; + final List messages = Arrays.asList("msg-0", "msg-1", "msg-2"); + DispatchRate dispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(1) + .dispatchThrottlingRateInByte(-1) + .ratePeriodInSecond(2) + .build(); + Producer producer = null; + Consumer consumer = null; + boolean topicCreated = false; + try { + admin1.namespaces().setReplicatorDispatchRate(replicatedNamespace, dispatchRate); + admin1.topics().createNonPartitionedTopic(topicName); + topicCreated = true; + waitReplicatorStarted(topicName); + 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 received = new HashSet<>(); + for (int i = 0; i < messages.size(); i++) { + Message message = consumer.receive(30, TimeUnit.SECONDS); + assertNotNull(message); + received.add(message.getValue()); + consumer.acknowledge(message); + } + + assertEquals(received, new HashSet<>(messages)); + waitForReplicationTaskFinish(topicName); + ensureNoBacklogByInflightTask(getReplicator(topicName)); + } finally { + if (producer != null) { + producer.close(); + } + if (consumer != null) { + consumer.close(); + } + admin1.namespaces().setReplicatorDispatchRate(replicatedNamespace, null); + 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 7498fd3ca8bb5..7f07dd8416f73 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 @@ -58,6 +58,7 @@ import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @CustomLog @@ -173,14 +174,17 @@ public void testReadEntriesFailedCompletesInFlightTaskAfterReplicatorTerminated( } } - @Test - public void testRateLimiterWithoutPermitsDoesNotCreateInFlightTask() throws Exception { - assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(0, -1); - assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(-1, 0); + @DataProvider + public Object[][] rateLimiterWithoutPermits() { + return new Object[][] { + {0, -1}, + {-1, 0} + }; } - private void assertRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long availableMessages, - long availableBytes) throws Exception { + @Test(dataProvider = "rateLimiterWithoutPermits") + public void testRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long availableMessages, + long availableBytes) throws Exception { PersistentReplicator replicator = getReplicator(topicName); LinkedList inFlightTasks = replicator.inFlightTasks; @@ -259,7 +263,7 @@ public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { Assert.assertEquals(inFlightTasks.size(), 1, "Queue should have one task"); Assert.assertEquals(task1.getReadPos(), position1, "Task should have the correct position"); Assert.assertEquals(task1.getReadingEntries(), 10, "Task should have the correct reading entries count"); - Assert.assertEquals(task1.getBytesToRead(), -1, "Task should have the correct byte read limit"); + Assert.assertEquals(task1.getMaxBytesToRead(), -1, "Task should have the correct byte read limit"); // Mark the task as done to test recycling task1.setEntries(Collections.emptyList()); @@ -272,7 +276,7 @@ public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { Assert.assertEquals(inFlightTasks.size(), 1, "Queue should still have one task"); Assert.assertEquals(task2.getReadPos(), position2, "Task should have the updated position"); Assert.assertEquals(task2.getReadingEntries(), 20, "Task should have the updated reading entries count"); - Assert.assertEquals(task2.getBytesToRead(), 1024, "Task should have the updated byte read limit"); + Assert.assertEquals(task2.getMaxBytesToRead(), 1024, "Task should have the updated byte read limit"); // Test Case 3: Create a new task when no tasks can be recycled task2.setEntries(null); // Make the task not done @@ -284,7 +288,7 @@ public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { Assert.assertEquals(inFlightTasks.size(), 2, "Queue should have two tasks"); Assert.assertEquals(task3.getReadPos(), position3, "Task should have the correct position"); Assert.assertEquals(task3.getReadingEntries(), 30, "Task should have the correct reading entries count"); - Assert.assertEquals(task3.getBytesToRead(), 2048, "Task should have the correct byte read limit"); + Assert.assertEquals(task3.getMaxBytesToRead(), 2048, "Task should have the correct byte read limit"); // cleanup. log.info("Completed testCreateOrRecycleInFlightTaskIntoQueue"); From b9531c61a9b868636db4a654d116c34dbac2becb Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 26 Jun 2026 00:46:14 +0800 Subject: [PATCH 17/20] Address replicator read limit review comments --- .../persistent/PersistentReplicator.java | 100 +++++++++--------- .../broker/service/OneWayReplicatorTest.java | 16 ++- .../PersistentReplicatorInflightTaskTest.java | 70 ++++++------ 3 files changed, 96 insertions(+), 90 deletions(-) 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 8c64798fe45f7..33f67cf1ef76e 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 @@ -219,7 +219,8 @@ protected void disableReplicatorRead() { this.cursor.setInactive(); } - private record ReadLimits(int messages, long bytes) { + @VisibleForTesting + record ReadLimits(int messages, long bytes) { public boolean isReadable() { return messages > 0 && bytes > 0; } @@ -301,9 +302,15 @@ protected void readMoreEntries() { if (state.equals(Terminated) || state.equals(Terminating)) { return; } - InFlightTask newInFlightTask = maybeCreateInFlightReadTask(); + InFlightTask newInFlightTask = null; + ReadLimits readLimits = null; + synchronized (inFlightTasks) { + readLimits = maybeGetReadLimitsForNextReadInLock(); + if (readLimits != null) { + newInFlightTask = createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), readLimits.messages); + } + } if (newInFlightTask == null) { - // no permits from rate limit log.debug("Not scheduling read due to pending read or no permits"); if (!hasPendingRead()) { topic.getBrokerService().executor().schedule( @@ -313,9 +320,8 @@ protected void readMoreEntries() { } log.debug() .attr("readingEntries", newInFlightTask.readingEntries) - .attr("maxBytesToRead", newInFlightTask.maxBytesToRead) .log("Scheduling read"); - cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, newInFlightTask.maxBytesToRead, this, + cursor.asyncReadEntriesOrWait(newInFlightTask.readingEntries, readLimits.bytes, this, newInFlightTask/* Context object */, topic.getMaxReadPosition()); } @@ -822,7 +828,6 @@ public ManagedCursor getCursor() { protected static class InFlightTask { Position readPos; int readingEntries; - long maxBytesToRead; volatile List entries; volatile int completedEntries; volatile boolean skipReadResultDueToCursorRewind; @@ -839,23 +844,17 @@ public synchronized void incCompletedEntries() { } } - synchronized void recycle(Position readStart, int readingEntries, long maxBytesToRead) { + synchronized void recycle(Position readStart, int readingEntries) { this.readPos = readStart; this.readingEntries = readingEntries; - this.maxBytesToRead = maxBytesToRead; this.entries = null; this.completedEntries = 0; this.skipReadResultDueToCursorRewind = false; } public InFlightTask(Position readPos, int readingEntries, String replicatorId) { - this(readPos, readingEntries, -1, replicatorId); - } - - public InFlightTask(Position readPos, int readingEntries, long maxBytesToRead, String replicatorId) { this.readPos = readPos; this.readingEntries = readingEntries; - this.maxBytesToRead = maxBytesToRead; this.replicatorId = replicatorId; } @@ -875,7 +874,6 @@ public String toString() { + "{replicatorId=" + replicatorId + ", readPos=" + readPos + ", readingEntries=" + readingEntries - + ", maxBytesToRead=" + maxBytesToRead + ", readoutEntries=" + (entries == null ? "-1" : entries.size()) + ", completedEntries=" + completedEntries + ", skipReadResultDueToCursorRewound=" + skipReadResultDueToCursorRewind @@ -884,7 +882,7 @@ public String toString() { } @VisibleForTesting - InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries, long maxBytesToRead) { + InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingEntries) { synchronized (inFlightTasks) { // Reuse projects that has done. if (!inFlightTasks.isEmpty()) { @@ -892,58 +890,60 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE if (first.isDone()) { // Remove from the first index, and add to the latest index. inFlightTasks.poll(); - first.recycle(readPos, readingEntries, maxBytesToRead); + first.recycle(readPos, readingEntries); inFlightTasks.add(first); return first; } } // New project if nothing can be reused. - InFlightTask task = new InFlightTask(readPos, readingEntries, maxBytesToRead, replicatorId); + InFlightTask task = new InFlightTask(readPos, readingEntries, replicatorId); inFlightTasks.add(task); return task; } } @VisibleForTesting - InFlightTask maybeCreateInFlightReadTask() { + ReadLimits maybeGetReadLimitsForNextRead() { synchronized (inFlightTasks) { - if (hasPendingRead()) { - log.info("Skip the reading because there is a pending read task"); - return null; - } - if (waitForCursorRewindingRefCnf > 0) { - log.info("Skip the reading due to new detected schema"); - return null; - } - if (state != Started) { - log.info("Skip the reading because producer has not started"); - return null; - } - int permits = getPermitsIfNoPendingRead(); - if (permits <= 0) { - return null; - } + return maybeGetReadLimitsForNextReadInLock(); + } + } - 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; - } + private ReadLimits maybeGetReadLimitsForNextReadInLock() { + if (hasPendingRead()) { + log.info("Skip the reading because there is a pending read task"); + return null; + } + if (waitForCursorRewindingRefCnf > 0) { + log.info("Skip the reading due to new detected schema"); + return null; + } + if (state != Started) { + log.info("Skip the reading because producer has not started"); + return null; + } + int permits = getPermitsIfNoPendingRead(); + if (permits <= 0) { + return null; + } - ReadLimits readLimits = getReadLimits(permits); + 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; + } - if (!readLimits.isReadable()) { - // no rate limiter permits from rate limit - log.debug() - .attr("messages", readLimits.messages) - .attr("bytes", readLimits.bytes) - .log("Throttling replication traffic"); - return null; - } + ReadLimits readLimits = getReadLimits(permits); - return createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), readLimits.messages, - readLimits.bytes); + if (!readLimits.isReadable()) { + // no rate limiter permits from rate limit + log.debug() + .attr("messages", readLimits.messages) + .attr("bytes", readLimits.bytes) + .log("Throttling replication traffic"); + return null; } + return readLimits; } protected int getPermitsIfNoPendingRead() { 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 d86732aee7913..1a18496b1778c 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 @@ -2445,11 +2445,17 @@ public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Excepti Producer producer = null; Consumer consumer = null; boolean topicCreated = false; + boolean dispatchRateConfigured = false; try { - admin1.namespaces().setReplicatorDispatchRate(replicatedNamespace, dispatchRate); admin1.topics().createNonPartitionedTopic(topicName); topicCreated = true; - waitReplicatorStarted(topicName); + admin1.topicPolicies().setReplicatorDispatchRate(topicName, dispatchRate); + dispatchRateConfigured = true; + GeoPersistentReplicator replicator = getReplicator(topicName); + Awaitility.await().untilAsserted(() -> { + assertTrue(replicator.getRateLimiter().isPresent()); + assertEquals(replicator.getRateLimiter().get().getDispatchRateOnMsg(), 1); + }); consumer = client2.newConsumer(Schema.STRING) .topic(topicName) .subscriptionName(subscriptionName) @@ -2474,7 +2480,7 @@ public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Excepti assertEquals(received, new HashSet<>(messages)); waitForReplicationTaskFinish(topicName); - ensureNoBacklogByInflightTask(getReplicator(topicName)); + ensureNoBacklogByInflightTask(replicator); } finally { if (producer != null) { producer.close(); @@ -2482,7 +2488,9 @@ public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Excepti if (consumer != null) { consumer.close(); } - admin1.namespaces().setReplicatorDispatchRate(replicatedNamespace, null); + if (dispatchRateConfigured) { + admin1.topicPolicies().removeReplicatorDispatchRate(topicName); + } if (topicCreated) { admin1.topics().setReplicationClusters(topicName, Arrays.asList(cluster1)); waitReplicatorStopped(topicName, false); 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 7f07dd8416f73..f15c758143b7b 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 @@ -48,6 +48,7 @@ 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.broker.service.persistent.PersistentReplicator.ReadLimits; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; @@ -183,8 +184,8 @@ public Object[][] rateLimiterWithoutPermits() { } @Test(dataProvider = "rateLimiterWithoutPermits") - public void testRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long availableMessages, - long availableBytes) throws Exception { + public void testRateLimiterWithoutPermitsDoesNotReturnReadLimits(long availableMessages, + long availableBytes) throws Exception { PersistentReplicator replicator = getReplicator(topicName); LinkedList inFlightTasks = replicator.inFlightTasks; @@ -199,7 +200,7 @@ public void testRateLimiterWithoutPermitsDoesNotCreateInFlightTask(long availabl replicator.dispatchRateLimiter = Optional.of(rateLimiter); try { - Assert.assertNull(replicator.maybeCreateInFlightReadTask()); + Assert.assertNull(replicator.maybeGetReadLimitsForNextRead()); Assert.assertTrue(inFlightTasks.isEmpty()); Assert.assertFalse(replicator.hasPendingRead()); assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); @@ -257,38 +258,35 @@ public void testCreateOrRecycleInFlightTaskIntoQueue() throws Exception { // Test Case 1: Create a new task when the queue is empty Position position1 = PositionFactory.create(1, 1); Assert.assertNotNull(position1, "Position should not be null"); - InFlightTask task1 = replicator.createOrRecycleInFlightTaskIntoQueue(position1, 10, -1); + InFlightTask task1 = replicator.createOrRecycleInFlightTaskIntoQueue(position1, 10); // Verify a new task was created and added to the queue Assert.assertNotNull(task1, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 1, "Queue should have one task"); Assert.assertEquals(task1.getReadPos(), position1, "Task should have the correct position"); Assert.assertEquals(task1.getReadingEntries(), 10, "Task should have the correct reading entries count"); - Assert.assertEquals(task1.getMaxBytesToRead(), -1, "Task should have the correct byte read limit"); // Mark the task as done to test recycling task1.setEntries(Collections.emptyList()); // Test Case 2: Recycle an existing task Position position2 = PositionFactory.create(2, 2); Assert.assertNotNull(position2, "Position should not be null"); - InFlightTask task2 = replicator.createOrRecycleInFlightTaskIntoQueue(position2, 20, 1024); + InFlightTask task2 = replicator.createOrRecycleInFlightTaskIntoQueue(position2, 20); // Verify the task was recycled Assert.assertNotNull(task2, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 1, "Queue should still have one task"); Assert.assertEquals(task2.getReadPos(), position2, "Task should have the updated position"); Assert.assertEquals(task2.getReadingEntries(), 20, "Task should have the updated reading entries count"); - Assert.assertEquals(task2.getMaxBytesToRead(), 1024, "Task should have the updated byte read limit"); // Test Case 3: Create a new task when no tasks can be recycled task2.setEntries(null); // Make the task not done Position position3 = PositionFactory.create(3, 3); Assert.assertNotNull(position3, "Position should not be null"); - InFlightTask task3 = replicator.createOrRecycleInFlightTaskIntoQueue(position3, 30, 2048); + InFlightTask task3 = replicator.createOrRecycleInFlightTaskIntoQueue(position3, 30); // Verify a new task was created Assert.assertNotNull(task3, "Task should not be null"); Assert.assertEquals(inFlightTasks.size(), 2, "Queue should have two tasks"); Assert.assertEquals(task3.getReadPos(), position3, "Task should have the correct position"); Assert.assertEquals(task3.getReadingEntries(), 30, "Task should have the correct reading entries count"); - Assert.assertEquals(task3.getMaxBytesToRead(), 2048, "Task should have the correct byte read limit"); // cleanup. log.info("Completed testCreateOrRecycleInFlightTaskIntoQueue"); @@ -432,8 +430,8 @@ public void testGetPermitsIfNoPendingRead() throws Exception { } @Test - public void testMaybeCreateInFlightReadTask() throws Exception { - log.info("Starting testMaybeCreateInFlightReadTask"); + public void testMaybeGetReadLimitsForNextRead() throws Exception { + log.info("Starting testMaybeGetReadLimitsForNextRead"); // Get the replicator for the test topic PersistentReplicator replicator = getReplicator(topicName); Assert.assertNotNull(replicator, "Replicator should not be null"); @@ -452,17 +450,17 @@ public void testMaybeCreateInFlightReadTask() throws Exception { 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 + // Should return read limits int expectedPermits = replicator.getPermitsIfNoPendingRead(); Assert.assertTrue(expectedPermits > 0, "Should have available permits for the test"); - InFlightTask task1 = replicator.maybeCreateInFlightReadTask(); - 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(), 100, - "Task readingEntries should equal to readBatchSize"); - Assert.assertTrue(inFlightTasks.contains(task1), - "Task should be added to the inFlightTasks list"); + ReadLimits readLimits1 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNotNull(readLimits1, "Should return read limits in normal case"); + Assert.assertEquals(readLimits1.messages(), 100, + "Read limit should equal readBatchSize"); + Assert.assertEquals(readLimits1.bytes(), pulsar1.getConfig().getDispatcherMaxReadSizeBytes(), + "Byte read limit should equal dispatcherMaxReadSizeBytes"); + Assert.assertTrue(inFlightTasks.isEmpty(), + "Getting read limits should not add a task to the inFlightTasks list"); // Test Case 2: With pending read - should return null inFlightTasks.clear(); @@ -470,26 +468,26 @@ public void testMaybeCreateInFlightReadTask() throws Exception { InFlightTask pendingReadTask = new InFlightTask(position1, 5, ""); // Don't set readoutEntries to simulate pending read inFlightTasks.add(pendingReadTask); - InFlightTask task2 = replicator.maybeCreateInFlightReadTask(); - Assert.assertNull(task2, "Should return null when there is a pending read"); + ReadLimits readLimits2 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNull(readLimits2, "Should return null when there is a pending read"); // Test Case 3: With waitForCursorRewinding=true - should return null inFlightTasks.clear(); replicator.waitForCursorRewindingRefCnf = 1; - InFlightTask task3 = replicator.maybeCreateInFlightReadTask(); - Assert.assertNull(task3, "Should return null when waiting for cursor rewinding"); + ReadLimits readLimits3 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNull(readLimits3, "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.maybeCreateInFlightReadTask(); - Assert.assertNull(task4, "Should return null when state is not Started"); + ReadLimits readLimits4 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNull(readLimits4, "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 + // Test Case 5: With limited permits - verify message read limits are set correctly inFlightTasks.clear(); // Add a task with some in-flight messages to reduce available permits Position positionLimited = PositionFactory.create(10, 10); @@ -506,11 +504,11 @@ public void testMaybeCreateInFlightReadTask() throws Exception { 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.maybeCreateInFlightReadTask(); - 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"); + // Now verify read limits match the limited permits + ReadLimits readLimits5 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNotNull(readLimits5, "Should return read limits with limited permits"); + Assert.assertEquals(readLimits5.messages(), limitedPermits, + "Message read limit should equal the limited number of permits available"); // Test Case 6: With permits=0 - should return null inFlightTasks.clear(); @@ -526,9 +524,9 @@ public void testMaybeCreateInFlightReadTask() throws Exception { task.setEntries(entries); inFlightTasks.add(task); } - InFlightTask task6 = replicator.maybeCreateInFlightReadTask(); - Assert.assertNull(task6, "Should return null when permits is 0"); - log.info("Completed testMaybeCreateInFlightReadTask"); + ReadLimits readLimits6 = replicator.maybeGetReadLimitsForNextRead(); + Assert.assertNull(readLimits6, "Should return null when permits is 0"); + log.info("Completed testMaybeGetReadLimitsForNextRead"); } finally { // Restore original state replicator.waitForCursorRewindingRefCnf = originalWaitForCursorRewinding; @@ -546,7 +544,7 @@ public static Runnable pauseReplicator(PersistentReplicator replicator) { replicator.beforeTerminateOrCursorRewinding(PersistentReplicator.ReasonOfWaitForCursorRewinding.Disconnecting); replicator.doRewindCursor(false); InFlightTask inFlightTask = - replicator.createOrRecycleInFlightTaskIntoQueue(PositionFactory.create(1, 1), 1, -1); + replicator.createOrRecycleInFlightTaskIntoQueue(PositionFactory.create(1, 1), 1); return () -> { inFlightTask.setEntries(Collections.emptyList()); replicator.readMoreEntries(); From 8c4e5733e48862d9cf8b874c20c3608488d34456 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 26 Jun 2026 08:38:24 +0800 Subject: [PATCH 18/20] Remove test-only read limits helper --- .../persistent/GeoPersistentReplicator.java | 6 +- .../persistent/PersistentReplicator.java | 10 +- .../PersistentReplicatorInflightTaskTest.java | 149 ------------------ 3 files changed, 4 insertions(+), 161 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 52e11ece25921..c73a7d04bc15a 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 @@ -263,9 +263,9 @@ protected boolean replicateEntries(List entries, final InFlightTask inFli * Explain the result of the race-condition between: * - {@link #readMoreEntries} * - {@link #beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding)} - * Since {@link #maybeCreateInFlightReadTask} 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 33f67cf1ef76e..eaf65a2ebe2e5 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 @@ -219,8 +219,7 @@ protected void disableReplicatorRead() { this.cursor.setInactive(); } - @VisibleForTesting - record ReadLimits(int messages, long bytes) { + private record ReadLimits(int messages, long bytes) { public boolean isReadable() { return messages > 0 && bytes > 0; } @@ -902,13 +901,6 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } } - @VisibleForTesting - ReadLimits maybeGetReadLimitsForNextRead() { - synchronized (inFlightTasks) { - return maybeGetReadLimitsForNextReadInLock(); - } - } - private ReadLimits maybeGetReadLimitsForNextReadInLock() { if (hasPendingRead()) { log.info("Skip the reading because there is a pending read task"); 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 f15c758143b7b..59f249bfbcb89 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 @@ -22,7 +22,6 @@ import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import java.util.ArrayList; @@ -30,7 +29,6 @@ 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; @@ -43,12 +41,10 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerTest; import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.service.AbstractReplicator; -import org.apache.pulsar.broker.service.BrokerServiceInternalMethodInvoker; 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.broker.service.persistent.PersistentReplicator.ReadLimits; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; @@ -59,7 +55,6 @@ import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; -import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @CustomLog @@ -175,42 +170,6 @@ public void testReadEntriesFailedCompletesInFlightTaskAfterReplicatorTerminated( } } - @DataProvider - public Object[][] rateLimiterWithoutPermits() { - return new Object[][] { - {0, -1}, - {-1, 0} - }; - } - - @Test(dataProvider = "rateLimiterWithoutPermits") - public void testRateLimiterWithoutPermitsDoesNotReturnReadLimits(long availableMessages, - long availableBytes) throws Exception { - PersistentReplicator replicator = getReplicator(topicName); - - LinkedList inFlightTasks = replicator.inFlightTasks; - List originalTasks = new ArrayList<>(inFlightTasks); - Optional originalRateLimiter = replicator.dispatchRateLimiter; - inFlightTasks.clear(); - - 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); - - try { - Assert.assertNull(replicator.maybeGetReadLimitsForNextRead()); - Assert.assertTrue(inFlightTasks.isEmpty()); - Assert.assertFalse(replicator.hasPendingRead()); - assertEquals(replicator.getPermitsIfNoPendingRead(), 1000); - } finally { - inFlightTasks.clear(); - inFlightTasks.addAll(originalTasks); - replicator.dispatchRateLimiter = originalRateLimiter; - } - } - @Test public void testFailedPublishCompletesInFlightTask() throws Exception { PersistentReplicator replicator = spy(getReplicator(topicName)); @@ -429,114 +388,6 @@ public void testGetPermitsIfNoPendingRead() throws Exception { } } - @Test - public void testMaybeGetReadLimitsForNextRead() throws Exception { - log.info("Starting testMaybeGetReadLimitsForNextRead"); - // Get the replicator for the test topic - PersistentReplicator replicator = getReplicator(topicName); - Assert.assertNotNull(replicator, "Replicator should not be null"); - - // Get access to the inFlightTasks list for setup - LinkedList inFlightTasks = replicator.inFlightTasks; - Assert.assertNotNull(inFlightTasks, "InFlightTasks list should not be null"); - - // Save original tasks and clear for testing - List originalTasks = new ArrayList<>(inFlightTasks); - inFlightTasks.clear(); - - // Save original state - int originalWaitForCursorRewinding = replicator.waitForCursorRewindingRefCnf; - AbstractReplicator.State originalState = replicator.getState(); - - try { - // Test Case 1: Normal case - no pending read, not waiting for cursor rewinding, state is Started - // Should return read limits - int expectedPermits = replicator.getPermitsIfNoPendingRead(); - Assert.assertTrue(expectedPermits > 0, "Should have available permits for the test"); - ReadLimits readLimits1 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNotNull(readLimits1, "Should return read limits in normal case"); - Assert.assertEquals(readLimits1.messages(), 100, - "Read limit should equal readBatchSize"); - Assert.assertEquals(readLimits1.bytes(), pulsar1.getConfig().getDispatcherMaxReadSizeBytes(), - "Byte read limit should equal dispatcherMaxReadSizeBytes"); - Assert.assertTrue(inFlightTasks.isEmpty(), - "Getting read limits should not add a task 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); - ReadLimits readLimits2 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNull(readLimits2, "Should return null when there is a pending read"); - - // Test Case 3: With waitForCursorRewinding=true - should return null - inFlightTasks.clear(); - replicator.waitForCursorRewindingRefCnf = 1; - ReadLimits readLimits3 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNull(readLimits3, "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); - ReadLimits readLimits4 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNull(readLimits4, "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 message read limits are 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 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 verify read limits match the limited permits - ReadLimits readLimits5 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNotNull(readLimits5, "Should return read limits with limited permits"); - Assert.assertEquals(readLimits5.messages(), limitedPermits, - "Message read limit 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 entries = new ArrayList<>(); - for (int j = 0; j < 100; j++) { - entries.add(mock(Entry.class)); - } - task.setEntries(entries); - inFlightTasks.add(task); - } - ReadLimits readLimits6 = replicator.maybeGetReadLimitsForNextRead(); - Assert.assertNull(readLimits6, "Should return null when permits is 0"); - log.info("Completed testMaybeGetReadLimitsForNextRead"); - } finally { - // Restore original state - replicator.waitForCursorRewindingRefCnf = originalWaitForCursorRewinding; - BrokerServiceInternalMethodInvoker.replicatorSetState(replicator, originalState); - // Restore original tasks - inFlightTasks.clear(); - inFlightTasks.addAll(originalTasks); - } - } - public static Runnable pauseReplicator(PersistentReplicator replicator) { Awaitility.await().untilAsserted(() -> { assertTrue(replicator.isConnected()); From 1f02f500778fed054b3834152d2341f7d44202ce Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 26 Jun 2026 13:00:52 +0800 Subject: [PATCH 19/20] Cover byte-rate replicator throttling --- .../broker/service/OneWayReplicatorTest.java | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) 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 1a18496b1778c..f54daebfba7fd 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 @@ -2432,14 +2432,22 @@ public void testReplicatorsInflightTaskListIsEmptyAfterReplicationFinished() thr ensureNoBacklogByInflightTask(getReplicator(topicName)); } - @Test(timeOut = 90_000) - public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Exception { + @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 messages = Arrays.asList("msg-0", "msg-1", "msg-2"); DispatchRate dispatchRate = DispatchRate.builder() - .dispatchThrottlingRateInMsg(1) - .dispatchThrottlingRateInByte(-1) + .dispatchThrottlingRateInMsg(messageRate) + .dispatchThrottlingRateInByte(byteRate) .ratePeriodInSecond(2) .build(); Producer producer = null; @@ -2454,7 +2462,8 @@ public void testReplicatorContinuesAfterRateLimiterHasNoPermits() throws Excepti GeoPersistentReplicator replicator = getReplicator(topicName); Awaitility.await().untilAsserted(() -> { assertTrue(replicator.getRateLimiter().isPresent()); - assertEquals(replicator.getRateLimiter().get().getDispatchRateOnMsg(), 1); + assertEquals(replicator.getRateLimiter().get().getDispatchRateOnMsg(), messageRate); + assertEquals(replicator.getRateLimiter().get().getDispatchRateOnByte(), byteRate); }); consumer = client2.newConsumer(Schema.STRING) .topic(topicName) From da0d4ac6c095e5fa1111686abb9c61277baf11d5 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Sat, 27 Jun 2026 13:17:17 +0800 Subject: [PATCH 20/20] [test][broker] Cover replicator read scheduling behavior --- .../persistent/PersistentReplicator.java | 68 +++---- .../broker/service/OneWayReplicatorTest.java | 18 +- .../PersistentReplicatorInflightTaskTest.java | 169 ++++++++++++++++++ 3 files changed, 207 insertions(+), 48 deletions(-) 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 eaf65a2ebe2e5..8444205be5579 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 @@ -304,9 +304,34 @@ protected void readMoreEntries() { InFlightTask newInFlightTask = null; ReadLimits readLimits = null; synchronized (inFlightTasks) { - readLimits = maybeGetReadLimitsForNextReadInLock(); - if (readLimits != null) { - newInFlightTask = createOrRecycleInFlightTaskIntoQueue(cursor.getReadPosition(), readLimits.messages); + 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 + log.debug() + .attr("messages", readLimits.messages) + .attr("bytes", readLimits.bytes) + .log("Throttling replication traffic"); + } + } } } if (newInFlightTask == null) { @@ -901,43 +926,6 @@ InFlightTask createOrRecycleInFlightTaskIntoQueue(Position readPos, int readingE } } - private ReadLimits maybeGetReadLimitsForNextReadInLock() { - if (hasPendingRead()) { - log.info("Skip the reading because there is a pending read task"); - return null; - } - if (waitForCursorRewindingRefCnf > 0) { - log.info("Skip the reading due to new detected schema"); - return null; - } - if (state != Started) { - log.info("Skip the reading because producer has not started"); - return null; - } - int permits = getPermitsIfNoPendingRead(); - if (permits <= 0) { - return null; - } - - 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 readLimits = getReadLimits(permits); - - if (!readLimits.isReadable()) { - // no rate limiter permits from rate limit - log.debug() - .attr("messages", readLimits.messages) - .attr("bytes", readLimits.bytes) - .log("Throttling replication traffic"); - return null; - } - return readLimits; - } - 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 f54daebfba7fd..7f69b48550986 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 @@ -2479,15 +2479,17 @@ public void testReplicatorContinuesAfterRateLimiterHasNoPermits(int messageRate, producer.send(message); } + Set expected = new HashSet<>(messages); Set received = new HashSet<>(); - for (int i = 0; i < messages.size(); i++) { - Message message = consumer.receive(30, TimeUnit.SECONDS); - assertNotNull(message); - received.add(message.getValue()); - consumer.acknowledge(message); - } - - assertEquals(received, new HashSet<>(messages)); + Consumer subscribedConsumer = consumer; + Awaitility.await().atMost(Duration.ofSeconds(60)).untilAsserted(() -> { + Message message = subscribedConsumer.receive(1, TimeUnit.SECONDS); + if (message != null) { + received.add(message.getValue()); + subscribedConsumer.acknowledge(message); + } + assertEquals(received, expected); + }); waitForReplicationTaskFinish(topicName); ensureNoBacklogByInflightTask(replicator); } finally { 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 59f249bfbcb89..96c2c2cbae627 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,43 +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.CustomLog; 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.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; @CustomLog @@ -198,6 +219,63 @@ public void testFailedPublishCompletesInFlightTask() throws Exception { } } + @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"); @@ -388,6 +466,97 @@ public void testGetPermitsIfNoPendingRead() throws Exception { } } + @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 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); + } + + private static class TestReplicatorFixture { + final TestPersistentReplicator replicator; + final ManagedCursor cursor; + final EventLoopGroup executor; + + TestReplicatorFixture(TestPersistentReplicator replicator, ManagedCursor cursor, EventLoopGroup executor) { + this.replicator = replicator; + this.cursor = cursor; + this.executor = executor; + } + } + + private static class TestPersistentReplicator extends PersistentReplicator { + private final boolean writable; + + 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; + } + + @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 entries, InFlightTask inFlightTask) { + return true; + } + } + public static Runnable pauseReplicator(PersistentReplicator replicator) { Awaitility.await().untilAsserted(() -> { assertTrue(replicator.isConnected());