Skip to content

Commit c765056

Browse files
committed
[fix][broker] Defer replication read processing from producer callbacks
Coalesce producer acknowledgments into one queued owner turn, including send-failure cancellation and rewind. Preserve ownership through executor rejection cleanup and verify that producer callbacks never drain cursor work under the producer monitor. Assisted-by: Codex
1 parent a50b164 commit c765056

2 files changed

Lines changed: 144 additions & 12 deletions

File tree

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

Lines changed: 37 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -315,14 +315,35 @@ protected void readMoreEntries() {
315315
}
316316

317317
private void requestReadProcessing(boolean requestRead) {
318+
requestReadProcessing(requestRead, false);
319+
}
320+
321+
private void requestReadProcessing(boolean requestRead, boolean runAsync) {
318322
synchronized (inFlightTasks) {
319323
readRequested |= requestRead;
320324
if (processingReads) {
321325
return;
322326
}
323327
processingReads = true;
324328
}
325-
processReads();
329+
if (runAsync) {
330+
try {
331+
// ACKs may hold the geo producer monitor on its IO thread. Retain ownership while queued,
332+
// so concurrent callbacks only publish work rather than starting another drain.
333+
brokerService.executor().execute(this::processReads);
334+
} catch (RuntimeException e) {
335+
log.error().exception(e).log("Failed to schedule replication read processing");
336+
// Retain ownership through termination: its cleanup request must not start an inline drain
337+
// on the ACK thread either. A late read callback will settle any still-pending result.
338+
try {
339+
terminate();
340+
} finally {
341+
discardPendingReadResults();
342+
}
343+
}
344+
} else {
345+
processReads();
346+
}
326347
}
327348

328349
private void processReads() {
@@ -596,15 +617,20 @@ protected static final class ProducerSendCallback implements SendCallback {
596617

597618
@Override
598619
public void sendComplete(Throwable exception, OpSendMsgStats opSendMsgStats) {
599-
if (exception != null && !(exception instanceof PulsarClientException.InvalidMessageException)) {
620+
boolean failed = exception != null && !(exception instanceof PulsarClientException.InvalidMessageException);
621+
if (failed) {
600622
replicator.log.error()
601623
.attr("inFlightTasks", replicator.inFlightTasks)
602624
.attr("pendingQueueSize", replicator.producer.getPendingQueueSize())
603625
.exception(exception)
604626
.log("Error producing on remote broker");
605-
// cursor should be rewound since it was incremented when readMoreEntries
606-
replicator.beforeTerminateOrCursorRewinding(ReasonOfWaitForCursorRewinding.Failed_Publishing);
607-
replicator.doRewindCursor(false);
627+
synchronized (replicator.inFlightTasks) {
628+
// Unlike asynchronous schema lookup, this recovery has no outstanding stage to wait for.
629+
// Publish cancellation and rewind together, leaving cursor work to the owner.
630+
replicator.inFlightTasks.forEach(task -> task.skipReadResultDueToCursorRewind = true);
631+
replicator.cancelReadRequested = true;
632+
replicator.rewindRequested = true;
633+
}
608634
// The failed send has completed from the producer queue perspective. The cursor rewind
609635
// makes the entry readable again, so this in-flight task must release its permit.
610636
inFlightTask.incCompletedEntries();
@@ -629,18 +655,20 @@ public void sendComplete(Throwable exception, OpSendMsgStats opSendMsgStats) {
629655
pendingRead = replicator.hasPendingRead();
630656
permits = pendingRead ? 0 : replicator.getPermitsIfNoPendingRead();
631657
}
632-
if (pendingRead) {
633-
replicator.readMoreEntries();
634-
} else if (replicator.producerQueueSize - permits < replicator.producerQueueThreshold) {
658+
boolean requestRead = pendingRead;
659+
if (!pendingRead && replicator.producerQueueSize - permits < replicator.producerQueueThreshold) {
635660
if (replicator.producerQueueSize == permits || replicator.producer.isWritable()) {
636-
replicator.readMoreEntries();
661+
requestRead = true;
637662
} else {
638663
replicator.log.debug()
639664
.attr("pending", replicator.producerQueueSize - permits)
640665
.attr("isWritable", replicator.producer.isWritable())
641666
.log("Not resuming reads");
642667
}
643668
}
669+
if (requestRead || failed) {
670+
replicator.requestReadProcessing(requestRead, true);
671+
}
644672

645673
recycle();
646674
}

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

Lines changed: 107 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import static org.mockito.Mockito.doThrow;
2828
import static org.mockito.Mockito.inOrder;
2929
import static org.mockito.Mockito.mock;
30+
import static org.mockito.Mockito.never;
3031
import static org.mockito.Mockito.times;
3132
import static org.mockito.Mockito.verify;
3233
import static org.mockito.Mockito.when;
@@ -35,6 +36,7 @@
3536
import java.util.Collections;
3637
import java.util.List;
3738
import java.util.Queue;
39+
import java.util.concurrent.CompletableFuture;
3840
import java.util.concurrent.ConcurrentLinkedQueue;
3941
import java.util.concurrent.CountDownLatch;
4042
import java.util.concurrent.RejectedExecutionException;
@@ -55,6 +57,7 @@
5557
import org.apache.pulsar.broker.service.persistent.PersistentReplicator.InFlightTask;
5658
import org.apache.pulsar.client.admin.PulsarAdmin;
5759
import org.apache.pulsar.client.api.ProducerBuilder;
60+
import org.apache.pulsar.client.api.PulsarClientException;
5861
import org.apache.pulsar.client.api.Schema;
5962
import org.apache.pulsar.client.impl.ProducerImpl;
6063
import org.apache.pulsar.client.impl.PulsarClientImpl;
@@ -463,17 +466,119 @@ public void testProducerAckDuringReservedReadRetriesFailedReadBeforeTimerWithout
463466

464467
reservedRead.callback.readEntriesFailed(new ManagedLedgerException.TooManyRequestsException("read failed"),
465468
reservedRead.context);
469+
fixture.runQueuedWork();
466470
ScheduledWork fallbackRetry = fixture.takeScheduledWork();
467471
synchronized (requests) {
468472
// The ACK demand is consumed after the failed reservation settles, before its fallback timer runs.
469473
assertThat(requests).hasSize(3);
470474
}
471475

476+
ReadRequest immediateRetry = requests.get(2);
477+
immediateRetry.callback.readEntriesFailed(new ManagedLedgerException.TooManyRequestsException("retry failed"),
478+
immediateRetry.context);
479+
fixture.runQueuedWork();
480+
assertThat(requests).hasSize(3);
481+
assertThat(fixture.scheduledWork).isEmpty();
482+
472483
fallbackRetry.command.run();
473484
synchronized (requests) {
474-
// The fallback request cannot overlap the retry it found pending.
475-
assertThat(requests).hasSize(3);
485+
// One ACK permits only one immediate retry. A second failure must wait for the existing timer.
486+
assertThat(requests).hasSize(4);
487+
}
488+
}
489+
490+
@DataProvider
491+
public Object[][] sendCompletionOutcomes() {
492+
return new Object[][] {{false}, {true}};
493+
}
494+
495+
@Test(dataProvider = "sendCompletionOutcomes")
496+
public void testSendCompletionQueuesReadProcessingOutsideProducerMonitor(boolean failedSend) throws Exception {
497+
ServiceConfiguration configuration = new ServiceConfiguration();
498+
configuration.setReplicationProducerQueueSize(2);
499+
TestReplicatorFixture fixture = newTestReplicatorFixture(configuration);
500+
TestPersistentReplicator replicator = fixture.replicator;
501+
ProducerImpl<?> producer = mock(ProducerImpl.class);
502+
when(producer.isWritable()).thenReturn(true);
503+
replicator.setProducerForTest(producer);
504+
List<ReadRequest> requests = new ArrayList<>();
505+
doAnswer(invocation -> {
506+
assertThat(Thread.holdsLock(producer)).isFalse();
507+
requests.add(new ReadRequest(invocation.getArgument(2), invocation.getArgument(3)));
508+
return null;
509+
}).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
510+
doAnswer(invocation -> {
511+
assertThat(Thread.holdsLock(producer)).isFalse();
512+
return null;
513+
}).when(fixture.cursor).rewind();
514+
doAnswer(invocation -> {
515+
assertThat(Thread.holdsLock(producer)).isFalse();
516+
return false;
517+
}).when(fixture.cursor).cancelPendingReadRequest();
518+
List<PersistentReplicator.ProducerSendCallback> callbacks = new ArrayList<>();
519+
replicator.entryObserver = (entry, task, entries) -> callbacks.add(
520+
PersistentReplicator.ProducerSendCallback.create(replicator, entry, null, task));
521+
522+
replicator.readMoreEntries();
523+
ReadRequest read = requests.get(0);
524+
read.callback.readEntriesComplete(List.of(entry(0), entry(1)), read.context);
525+
assertThat(requests).hasSize(1);
526+
assertThat(callbacks).hasSize(2);
527+
synchronized (producer) {
528+
callbacks.get(0).sendComplete(failedSend ? new PulsarClientException("send failed") : null, null);
529+
callbacks.get(1).sendComplete(null, null);
530+
assertThat(requests).hasSize(1);
531+
verify(fixture.cursor, never()).rewind();
532+
verify(fixture.cursor, never()).cancelPendingReadRequest();
533+
assertThat(fixture.queuedWork).hasSize(1);
534+
}
535+
fixture.runQueuedWork();
536+
assertThat(requests).hasSize(2);
537+
if (failedSend) {
538+
verify(fixture.cursor).rewind();
539+
assertThat(replicator.waitForCursorRewindingRefCnf).isZero();
540+
}
541+
}
542+
543+
@Test
544+
public void testRejectedAckHandoffTerminatesWithoutDrainingOnAckThread() throws Exception {
545+
ServiceConfiguration configuration = new ServiceConfiguration();
546+
configuration.setReplicationProducerQueueSize(2);
547+
TestReplicatorFixture fixture = newTestReplicatorFixture(configuration);
548+
TestPersistentReplicator replicator = fixture.replicator;
549+
ProducerImpl<?> producer = mock(ProducerImpl.class);
550+
when(producer.isWritable()).thenReturn(true);
551+
when(producer.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
552+
replicator.setProducerForTest(producer);
553+
List<ReadRequest> requests = new ArrayList<>();
554+
doAnswer(invocation -> {
555+
requests.add(new ReadRequest(invocation.getArgument(2), invocation.getArgument(3)));
556+
return null;
557+
}).when(fixture.cursor).asyncReadEntriesOrWait(anyInt(), anyLong(), any(), any(), any());
558+
List<PersistentReplicator.ProducerSendCallback> callbacks = new ArrayList<>();
559+
replicator.entryObserver = (entry, task, entries) -> callbacks.add(
560+
PersistentReplicator.ProducerSendCallback.create(replicator, entry, null, task));
561+
replicator.readMoreEntries();
562+
Entry first = entry(0);
563+
Entry second = entry(1);
564+
ReadRequest read = requests.get(0);
565+
read.callback.readEntriesComplete(List.of(first, second), read.context);
566+
doThrow(new RejectedExecutionException("ACK handoff rejected"))
567+
.when(fixture.executor).execute(any(Runnable.class));
568+
569+
synchronized (producer) {
570+
callbacks.get(0).sendComplete(null, null);
571+
callbacks.get(1).sendComplete(null, null);
476572
}
573+
574+
assertThat(replicator.getState()).isEqualTo(State.Terminated);
575+
assertThat(requests).hasSize(1);
576+
assertThat(fixture.queuedWork).isEmpty();
577+
verify(fixture.cursor, never()).rewind();
578+
verify(fixture.cursor, never()).cancelPendingReadRequest();
579+
verify(first).release();
580+
verify(second).release();
581+
verify(producer).closeAsync();
477582
}
478583

479584
@Test
@@ -607,7 +712,6 @@ private static TestReplicatorFixture newTestReplicatorFixture() throws Exception
607712
@SuppressWarnings("unchecked")
608713
private static TestReplicatorFixture newTestReplicatorFixture(ServiceConfiguration configuration) throws Exception {
609714
configuration.setClusterName("local");
610-
configuration.setReplicationProducerQueueSize(1000);
611715
configuration.setDispatcherMaxReadBatchSize(1000);
612716
configuration.setDispatcherMaxReadSizeBytes(1024 * 1024);
613717

0 commit comments

Comments
 (0)