Skip to content

Commit a3f2c34

Browse files
committed
[fix][broker] Preserve replication recovery after scheduling failures
Terminate when retry submission fails, preserve an accepted timer if producer startup throws, and rewind through the read owner before retrying a rejected completion whose cursor may have advanced. Assisted-by: Codex
1 parent d8636e4 commit a3f2c34

2 files changed

Lines changed: 87 additions & 21 deletions

File tree

‎pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java‎

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -403,11 +403,9 @@ private void discardPendingReadResults() {
403403
private void handleReadRetrySchedulingFailure(Exception exception) {
404404
// Ownership has already been released. Never clear a newer owner's state here.
405405
log.error().exception(exception).log("Failed to schedule replication read retry");
406-
if (exception instanceof RejectedExecutionException) {
407-
// A rejected retry has no wakeup left if there are no producer ACKs in flight.
408-
// Do not leave the replicator apparently Started but unable to make progress.
409-
terminate();
410-
}
406+
// A failed retry submission has no wakeup left if there are no producer ACKs in flight.
407+
// Do not leave the replicator apparently Started but unable to make progress.
408+
terminate();
411409
}
412410

413411
/** Processes one read, result or recovery transition, with no callback invoked under the state lock. */
@@ -468,13 +466,19 @@ private boolean processRead() {
468466
if (retryDelayMillis > 0) {
469467
try {
470468
scheduleReadRetry(retryDelayMillis);
471-
if (state == Disconnected) {
472-
startProducer();
473-
}
474469
} catch (Exception e) {
475470
// Ownership was already released: a new owner might be running now. Do not let
476471
// this failure reach the owner cleanup in processReads and clear its ownership.
477472
handleReadRetrySchedulingFailure(e);
473+
return false;
474+
}
475+
if (state == Disconnected) {
476+
try {
477+
startProducer();
478+
} catch (Exception e) {
479+
// The read-retry timer was accepted; this is not a timer scheduling failure.
480+
log.error().exception(e).log("Failed to restart replication producer; retry remains scheduled");
481+
}
478482
}
479483
return false;
480484
}
@@ -521,7 +525,7 @@ private void scheduleReadRetry(long delayMillis) {
521525
}
522526
readMoreEntries();
523527
}, delayMillis, TimeUnit.MILLISECONDS);
524-
} catch (RejectedExecutionException e) {
528+
} catch (RuntimeException e) {
525529
synchronized (inFlightTasks) {
526530
readRetryScheduled = false;
527531
}
@@ -747,6 +751,13 @@ private void handleReadFailure(ManagedLedgerException exception, InFlightTask ta
747751
terminate();
748752
return;
749753
}
754+
if (exception.getCause() instanceof RejectedExecutionException) {
755+
synchronized (inFlightTasks) {
756+
// Completion may be rejected after advancing the cursor but before transferring entries.
757+
// Only this owner can restore the position before it admits another read.
758+
rewindRequested = true;
759+
}
760+
}
750761
readBatchSize = brokerService.pulsar().getConfiguration().getDispatcherMinReadBatchSize();
751762
long waitTimeMillis = delayReadRetry();
752763
if (!(exception instanceof TooManyRequestsException)) {

‎pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorReadProcessingTest.java‎

Lines changed: 67 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -386,11 +386,19 @@ public void testReadFailureWaitsForTimerButNewReadDemandCanResumeImmediately() t
386386
}
387387
}
388388

389-
@Test
390-
public void testReadFailureTerminatesWhenRetrySchedulingIsRejectedWhileExecutorIsRunning() throws Exception {
389+
@DataProvider
390+
public Object[][] retrySchedulingFailures() {
391+
return new Object[][] {
392+
{new RejectedExecutionException("retry rejected")},
393+
{new IllegalStateException("retry scheduling failed")}
394+
};
395+
}
396+
397+
@Test(dataProvider = "retrySchedulingFailures")
398+
public void testReadFailureTerminatesWhenRetrySchedulingFails(RuntimeException failure) throws Exception {
391399
TestReplicatorFixture fixture = newTestReplicatorFixture();
392400
TestPersistentReplicator replicator = fixture.replicator;
393-
fixture.rejectScheduledWork.set(true);
401+
doThrow(failure).when(fixture.executor).schedule(any(Runnable.class), anyLong(), any(TimeUnit.class));
394402
when(fixture.executor.isShuttingDown()).thenReturn(false);
395403
List<ReadRequest> requests = new ArrayList<>();
396404
doAnswer(invocation -> {
@@ -417,6 +425,52 @@ public void testReadFailureTerminatesWhenRetrySchedulingIsRejectedWhileExecutorI
417425
verify(fixture.cursor).setInactive();
418426
}
419427

428+
@Test
429+
public void testProducerRestartFailureKeepsScheduledRetry() throws Exception {
430+
TestReplicatorFixture fixture = newTestReplicatorFixture();
431+
fixture.replicator.markDisconnected();
432+
fixture.replicator.producerStartFailure = new IllegalStateException("producer startup failed");
433+
434+
fixture.replicator.readMoreEntries();
435+
assertThat(fixture.replicator.getState()).isEqualTo(State.Disconnected);
436+
fixture.takeScheduledWork().command.run();
437+
assertThat(fixture.replicator.getState()).isEqualTo(State.Disconnected);
438+
assertThat(fixture.scheduledWork).hasSize(1);
439+
verify(fixture.cursor, never()).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
440+
}
441+
442+
@Test
443+
public void testRejectedReadCompletionRewindsBeforeRetriedRead() throws Exception {
444+
TestReplicatorFixture fixture = newTestReplicatorFixture();
445+
List<ReadRequest> requests = new ArrayList<>();
446+
doAnswer(invocation -> {
447+
requests.add(new ReadRequest(invocation.getArgument(2), invocation.getArgument(3)));
448+
return null;
449+
}).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
450+
Position start = PositionFactory.create(1, 1);
451+
Position advanced = PositionFactory.create(1, 3);
452+
when(fixture.cursor.getReadPosition()).thenReturn(start);
453+
fixture.replicator.readMoreEntries();
454+
when(fixture.cursor.getReadPosition()).thenReturn(advanced);
455+
doAnswer(invocation -> {
456+
when(fixture.cursor.getReadPosition()).thenReturn(start);
457+
return null;
458+
}).when(fixture.cursor).rewind();
459+
460+
ReadRequest read = requests.get(0);
461+
read.callback.readEntriesFailed(new ManagedLedgerException(new RejectedExecutionException("handoff rejected")),
462+
read.context);
463+
assertThat(requests).hasSize(1);
464+
fixture.takeScheduledWork().command.run();
465+
466+
assertThat(requests).hasSize(2);
467+
assertThat(((InFlightTask) requests.get(1).context).getReadPos()).isEqualTo(start);
468+
InOrder cursorCalls = inOrder(fixture.cursor);
469+
cursorCalls.verify(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
470+
cursorCalls.verify(fixture.cursor).rewind();
471+
cursorCalls.verify(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
472+
}
473+
420474
@Test
421475
public void testProducerAckDuringReservedReadRetriesFailedReadBeforeTimerWithoutOverlappingReads()
422476
throws Exception {
@@ -725,15 +779,11 @@ private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfigurati
725779
EventLoopGroup executor = mock(EventLoopGroup.class);
726780
Queue<Runnable> queuedWork = new ConcurrentLinkedQueue<>();
727781
Queue<ScheduledWork> scheduledWork = new ConcurrentLinkedQueue<>();
728-
AtomicBoolean rejectScheduledWork = new AtomicBoolean();
729782
doAnswer(invocation -> {
730783
queuedWork.add(invocation.getArgument(0));
731784
return null;
732785
}).when(executor).execute(any(Runnable.class));
733786
doAnswer(invocation -> {
734-
if (rejectScheduledWork.get()) {
735-
throw new RejectedExecutionException("test retry scheduler rejection");
736-
}
737787
scheduledWork.add(new ScheduledWork(invocation.getArgument(0), invocation.getArgument(1),
738788
invocation.getArgument(2)));
739789
return null;
@@ -765,7 +815,7 @@ private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfigurati
765815

766816
TestPersistentReplicator replicator = new TestPersistentReplicator(topic, cursor, brokerService,
767817
replicationClient, mock(PulsarAdmin.class));
768-
return new TestReplicatorFixture(replicator, cursor, executor, queuedWork, scheduledWork, rejectScheduledWork);
818+
return new TestReplicatorFixture(replicator, cursor, executor, queuedWork, scheduledWork);
769819
}
770820

771821
private record ReadRequest(ReadEntriesCallback callback, Object context) {
@@ -780,18 +830,15 @@ private static final class TestReplicatorFixture {
780830
private final EventLoopGroup executor;
781831
private final Queue<Runnable> queuedWork;
782832
private final Queue<ScheduledWork> scheduledWork;
783-
private final AtomicBoolean rejectScheduledWork;
784833

785834
private TestReplicatorFixture(TestPersistentReplicator replicator, ManagedCursor cursor,
786835
EventLoopGroup executor,
787-
Queue<Runnable> queuedWork, Queue<ScheduledWork> scheduledWork,
788-
AtomicBoolean rejectScheduledWork) {
836+
Queue<Runnable> queuedWork, Queue<ScheduledWork> scheduledWork) {
789837
this.replicator = replicator;
790838
this.cursor = cursor;
791839
this.executor = executor;
792840
this.queuedWork = queuedWork;
793841
this.scheduledWork = scheduledWork;
794-
this.rejectScheduledWork = rejectScheduledWork;
795842
}
796843

797844
private void runQueuedWork() {
@@ -821,6 +868,7 @@ private static final class TestPersistentReplicator extends PersistentReplicator
821868
private final AtomicBoolean rewindBeforeCursorInvocation = new AtomicBoolean();
822869
private final AtomicBoolean cancelBeforeCursorInvocation = new AtomicBoolean();
823870
private EntryObserver entryObserver;
871+
private RuntimeException producerStartFailure;
824872

825873
private TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, BrokerService brokerService,
826874
PulsarClientImpl replicationClient, PulsarAdmin replicationAdmin)
@@ -833,6 +881,9 @@ private TestPersistentReplicator(PersistentTopic topic, ManagedCursor cursor, Br
833881
@Override
834882
protected void startProducer() {
835883
// The test drives read scheduling directly.
884+
if (producerStartFailure != null) {
885+
throw producerStartFailure;
886+
}
836887
}
837888

838889
@Override
@@ -861,6 +912,10 @@ private void markTerminated() {
861912
state = State.Terminated;
862913
}
863914

915+
private void markDisconnected() {
916+
state = State.Disconnected;
917+
}
918+
864919
private void setProducerForTest(ProducerImpl<?> producer) {
865920
this.producer = producer;
866921
}

0 commit comments

Comments
 (0)