Skip to content

[fix][client] Schedule TableView tail-read retries instead of sleeping on the callback thread - #26672

Open
namest504 wants to merge 6 commits into
apache:masterfrom
namest504:fix/tableview-tail-retry
Open

namest504 wants to merge 6 commits into
apache:masterfrom
namest504:fix/tableview-tail-retry

Conversation

@namest504

@namest504 namest504 commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #26671

Main Issue: #26671

Motivation

When a tail read fails with anything other than AlreadyClosedException, AbstractTableViewImpl.readTailMessages() calls Thread.sleep(50) inside the exceptionally callback and then reads again. That callback runs on the thread that completed the failed read: normally the client's shared internal executor thread, which then blocks for 50 ms per retry, or the caller's own thread when readNextAsync() returns an already-failed future, in which case the retry recurses on the same stack indefinitely. There is no backoff: the read is retried every 50 ms for as long as it keeps failing. See #26671 for the details.

The delay itself was deliberate (#22064 added it as "a minor delay before retrying"). This change keeps that intent and only changes how the wait happens.

Modifications

  • The retry runs after a Backoff delay, built from the client's startingBackoffInterval / maxBackoffInterval, through CompletableFuture.delayedExecutor instead of sleeping. The backoff is reset once a message has been read and handled; a failure thrown while handling a message backs off like a failed read.
  • The delay is not timed on one of the client's own executors: a task queued there is discarded without notice when the executor is shut down (shutdownNow()), and with no read in flight while a retry waits, the tail reads would end there with nothing to record it, leaving a refresh waiting for a message pending until the table view is closed. The JDK's delayed executor does not depend on the client, so the retry still runs and finds the closed reader. ScalableTopicProducer times its retries the same way.
  • The tail-read loop records a stop cause when it stops: on closeAsync(), on a read that fails because the reader was closed under the table view, on a read the reader rejects because its executor was shut down under it (PulsarClientImpl#shutdown() does not close the consumers, and shared resources can be closed before the client), and on any failure once the client is closed. From then on the loop issues no read and a retry still waiting for its delay does nothing.
  • A refresh is tracked from the moment refreshAsync() is called, so a stop fails it right away whether it is waiting for a message or still looking up the last message ids, and a refresh started after the stop fails before any lookup, an empty topic included. The failure is an AlreadyClosedException wrapped in a CompletionException, the shape the failed read of a closed reader delivers, so callbacks that look at getCause() keep working.
  • The read is invoked through FutureUtil.supplySafely, so a reader that throws instead of returning a failed future is handled like a failed read.

Behaviour changes:

  • The first retry waits the client's initial backoff (default 100 ms) instead of a fixed 50 ms, and subsequent retries double up to maxBackoffInterval (default 60 s) with the usual jitter. While messages keep failing in handling, each of them delays the next read by the current backoff.
  • closeAsync() fails the refreshes that have not completed right away, on the caller's thread and before the reader is closed.
  • refreshAsync() on a table view whose tail reads have stopped fails right away with the stop cause, without looking up the last message ids.
  • A first tail read that throws no longer fails createAsync(): it is handled like any other failed read.

What this change does not cover:

  • While a retry is waiting no read is in flight, so a reader closed under the table view (the client being closed while the table view is still open) is only noticed when the retry runs, up to one backoff delay later. On master that took at most 50 ms. Closing the table view itself is immediate.
  • A retry still waiting for its delay when the table view is closed keeps the table view reachable until the delay elapses, up to maxBackoffInterval. The retry itself does nothing then. The same held for a task queued on the client's scheduler before this change.
  • A reader that keeps failing its reads while the client is open is retried at the backoff rate until the table view or the client is closed (shutdown() alone does not count: isClosed() stays false). On master it was retried every 50 ms.
  • PulsarClient#shutdown() does not close the reader: a tail read in flight then is not completed and the refreshes keep waiting, as on master.
  • A refresh waiting for a message whose handling throws still waits for the next message, as on master. That belongs with the follow-up on handling failures I mentioned in a comment below.

Verifying this change

  • Make sure that the change passes the CI checks.

TableViewImplTest covers the retry timing (testTailReadFailureSchedulesRetryInsteadOfRecursing hangs in the retry loop before the fix, testTailReadRetryDelayGrowsOnRepeatedFailures, testTailReadRetryDelayResetsAfterSuccessfulRead, testProcessingFailureBacksOffLikeAReadFailure, testReadThatThrowsIsRetriedLikeAFailedRead, testHandlingFailureThatIsARejectionIsRetried), the stops (testCloseStopsThePendingTailReadRetryAndFailsPendingRefreshes, testRetryThatRunsRightAwayDoesNotOutliveClose, testNoReadIsIssuedAfterTheTableViewClosed, testNoTailReadStartsWhenTheTableViewClosesDuringTheInitialReplay, testFirstStopCauseIsKept, testRetryWaitingWhenTheClientClosesStillFailsPendingRefreshes, testReadRejectedByItsExecutorFailsPendingRefreshes, testFailedReadIsNotRetriedOnceTheClientIsClosed) and the refresh lifecycle (testRefreshFetchingLastMessageIdsFailsWhenTheTableViewCloses, ...WhenTheReaderIsClosed, ...WhenTheReadIsRejected, testRefreshAfterTheTableViewClosedDoesNotLookUpTheLastMessageIds, testRefreshAfterARejectedReadFailsRightAway, testEmptyLookupAnswerArrivingDuringTheStopDoesNotCompleteTheRefresh, testCompletedRefreshesAreNoLongerTracked). Each fails without the code it covers.

Also ran TableViewBuilderImplTest, MappedTableViewImplTest and, in pulsar-broker, TableViewTest, ServiceUnitStateChannelTest, LoadDataStoreTest and ServiceUnitStateCompactionTest, plus checkstyle and spotless. Outside the test suite I ran the close and shutdown cases against a standalone broker with the code before and after.

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Tail-read retries of a TableView no longer block the client's internal executor thread. After the backoff delay the retry issues its read from the JDK's default async executor: a common-pool thread, or a new thread on a JVM whose common pool has a single worker. Messages are handled on the thread that completes the read, normally the consumer's internal thread as before; if the read is already complete when the retry attaches its callback, on the retry's thread. When a retry finds that the tail reads must stop, the refresh failures are delivered on its thread.

Documentation

  • doc-not-needed

Matching PR in forked repository

PR in forked repository: namest504#2

@namest504
namest504 force-pushed the fix/tableview-tail-retry branch from eb6bb2f to fa9ddf5 Compare September 21, 2026 08:33
…g on the callback thread

When a tail read failed with anything other than AlreadyClosedException,
AbstractTableViewImpl called Thread.sleep(50) inside the CompletableFuture
callback and then read again. That callback runs on the thread that
completed the failed read: the client's shared internal executor thread, so
everything pinned to it waited 50 ms per retry, or the caller's own thread
when readNextAsync() returned an already-failed future, in which case the
retry recursed on the same stack indefinitely (about 20 retries per second,
each adding frames and a WARN with a stack trace).

Retry through the client's scheduled executor with the client's configured
reconnection backoff (initialBackoffInterval doubling up to
maxBackoffInterval), reset the backoff after a successful read, and stop
retrying when the client's executor has been shut down. Both the classic
and the mapped table view share this path.
try {
((ScheduledExecutorService) client.getScheduledExecutorProvider().getExecutor())
.schedule(() -> readTailMessages(reader), delayMillis, TimeUnit.MILLISECONDS);
} catch (RejectedExecutionException e) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we link this scheduled retry to the TableView lifecycle? When backoff starts, there is no in-flight readNextAsync(), allowing closeAsync() to complete while this queued task remains. A pending refreshAsync() would then not encounter AlreadyClosedException until the delay ends—which could be as long as maxBackoffInterval. If schedule() is rejected, the current catch block stops the retry loop, leaving pending refresh futures unresolved indefinitely.

We should cancel or prevent the pending retry when closing, and immediately fail pendingRefreshRequests upon close or scheduling rejection. Adding tests for closing the TableView with a pending retry, and for rejected scheduling, would be beneficial as well.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in ec3af86: closeAsync() now cancels a retry that is still waiting for its delay and fails the pending refreshAsync() futures right away with AlreadyClosedException, instead of leaving them until the closed reader is found once the delay has elapsed. A retry that the scheduler rejects because the client is shutting down fails them the same way. testCloseCancelsPendingTailReadRetryAndFailsPendingRefreshes and testRejectedTailReadRetryFailsPendingRefreshes cover the two cases.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up in dcc60f0: a refresh that was still fetching the last message ids when the table view closed registers afterwards and could never complete; it now fails right away too. These failures keep the AlreadyClosedException wrapped the way the closed reader's failed read delivered it, since TableViewTest#testRefreshTaskCanBeCompletedWhenReaderClosed (and any callback looking at getCause()) relies on that.

private void readTailMessages(Reader<T> reader) {
reader.readNextAsync()
.thenAccept(msg -> {
tailReadBackoff.reset();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we move tailReadBackoff.reset() after handleMessage(msg) succeeds? This exceptionally also sees exceptions thrown from handleMessage, not only failures from readNextAsync().

Mapping errors are isolated now, but other processing code such as a TopicCompactionStrategy can still throw. With the reset before handleMessage, consecutive processing failures always retry with the initial delay, so the exponential backoff never really grows for that path.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Moved the reset after handleMessage(msg) in ec3af86: a failure thrown while handling a message now backs off like a failed read instead of retrying with the initial delay every time. testProcessingFailureBacksOffLikeAReadFailure drives it with a TopicCompactionStrategy that throws, as in your example.

…dling failures

A retry waiting for its backoff delay outlived closeAsync(): the closed
reader was only noticed once the delay had elapsed, up to
maxBackoffInterval, and a refresh pending at that time waited just as
long. Closing now cancels the pending retry and fails the pending
refreshes right away, and so does a retry that the client's executor
rejects because it is shutting down.

The backoff is reset only after a message has been handled, so a failure
thrown while handling it backs off like a failed read instead of
retrying with the initial delay every time.
…e before

closeAsync() fails the refresh requests it finds, but a refresh that is
still fetching the last message ids at that moment registers afterwards,
when the cancelled retry has left no read that could ever complete it.
Fail it right away instead.

A refresh that fails because the table view or the client closed now
carries the AlreadyClosedException the way the closed reader's failed
read used to deliver it, wrapped in a CompletionException, so callbacks
that look at getCause() keep working (TableViewTest relies on that).
try {
ScheduledFuture<?> retry = ((ScheduledExecutorService) client.getScheduledExecutorProvider().getExecutor())
.schedule(() -> readTailMessages(reader), delayMillis, TimeUnit.MILLISECONDS);
pendingTailReadRetry.set(retry);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There still seems to be a race between scheduling the retry and publishing its ScheduledFuture.

schedule() makes the task eligible to run before it returns the future. If that retry runs immediately, fails again, and schedules another retry, the nested call can store F2 in pendingTailReadRetry, then the original schedule() returns and overwrites it with the already-running/completed F1. A later closeAsync() would cancel F1 while F2 is still queued.

This is easier to hit than it may look because startingBackoffInterval supports sub-millisecond values, while backoff.next().toMillis() turns those delays into 0.

Could we avoid relying on a ScheduledFuture that is only published after scheduling? At minimum the scheduled runnable should check closed before starting another read; ideally the retry ownership/token should be published before the task can run and updated with CAS-like semantics.

A test where schedule() invokes the runnable before returning its ScheduledFuture should reproduce this ordering.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Right, I missed the 0 ms case. Dropped the future in ea520b2: close sets a stop flag first and the retry checks it before reading. Cancel didn't help anyway, the scheduler keeps a cancelled task queued until its delay is up. Added a test that runs the retry inside schedule(), it failed before the fix.

Comment on lines +552 to +556
} catch (RejectedExecutionException e) {
// The client is shutting down; the reader will be closed with it and nothing will be read any more.
log.info().attr("reader", reader.getTopic())
.log("Client is closed, giving up retrying tail messages.");
failPendingRefreshRequests(alreadyClosed("Client already closed"));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One smaller lifecycle case still looks incomplete here. Once schedule() is rejected, this tail-read loop has permanently stopped, but we only fail the refresh requests that are already registered and leave closed == false.

A refreshAsync() that starts after this catch can still fetch the last message ids, register itself, pass the closed check, and then wait forever because there is no tail read left to advance it.

Since the scheduled executor uses the default rejection policy, rejection here is effectively a terminal condition. Would it make sense to mark the TableView/tail loop terminal here as well (either closed = true or a separate stopped flag), so future refreshes fail immediately too? A test that calls refreshAsync() after the scheduler rejection would cover this case.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in ea520b2. Rejection now sets the same stop flag as close, and so does a read failing because the reader was closed under us. A refresh registering after that fails right away. Two tests added.

…s future

The ScheduledFuture of a retry was only published after schedule() had
returned, but with a zero delay (a sub-millisecond initialBackoffInterval
rounds down to 0 ms) the retry can run, fail again and queue its
successor before that, so the stored future could be stale and
closeAsync() cancelled the wrong one. Cancelling gained nothing anyway:
the client's scheduler keeps a cancelled task queued until its delay
elapses. Publish a stop cause before failing the pending refreshes and
let the scheduled task check it instead.

Every end of the tail-read loop sets that stop cause: closeAsync(), a
read that fails because the reader was closed under the table view, and
a schedule() the client's scheduler rejects while shutting down. A
refresh that registers after any of them fails right away rather than
waiting for a read that will never come.
@namest504

namest504 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor Author

Two things I noticed while working on this. I left them out on purpose, to keep this PR limited to how the loop waits:

  1. When handling a message throws (in practice, a TopicCompactionStrategy raising), that message is already consumed, so it is skipped with only a log, and the view stays out of sync for that key until a later update or a rebuild. That predates this PR. I'd like to open a follow-up issue to surface these failures, along the lines of onMappingError or a metric, rather than widen this change.

  2. With startingBackoffInterval set to 0, the retry delay stays 0 forever, so an immediately failing read would spin. The client's reconnection backoff behaves the same with that setting, so I left it alone, but I can add a minimum delay here if you prefer.

@Override
public CompletableFuture<Void> refreshAsync() {
CompletableFuture<Void> completableFuture = new CompletableFuture<>();
reader.thenCompose(reader -> getLastMessageIdOfNonEmptyTopics(reader).thenAccept(lastMessageIds -> {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think there is still one gap in the refresh lifecycle here.

stopCause is only checked after getLastMessageIdOfNonEmptyTopics() has completed (and even after the empty-topic fast path). So a refresh started after the tail loop has already stopped can still wait on getLastMessageIdsAsync() indefinitely/for its operation timeout, even though we already know no tail read will ever run again. If the lookup returns an empty result, the refresh even completes successfully before checking stopCause.

There is a similar window for a refresh that was already fetching last message ids when closeAsync() calls stopTailReads(): it has not been added to pendingRefreshRequests yet, so failPendingRefreshRequests() cannot see it. testRefreshInFlightWhileClosingFailsRightAway currently demonstrates this by asserting the refresh is still pending after close and only making it fail after manually completing lastMessageIds.

Could we make the stop signal cover the whole refresh operation, including the last-message-id lookup phase? At minimum we should check stopCause before starting the lookup and before the empty-result fast path; ideally a refresh already in that lookup phase should also be completed when stopTailReads() runs.

A useful regression test would leave getLastMessageIdsAsync() incomplete, call closeAsync() (or stop via scheduler rejection), and verify that the refresh future is completed exceptionally without completing the lookup future.

@namest504 namest504 Oct 2, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You're right, the stop only reached refreshes that were already registered. aa1483f tracks a refresh from the moment refreshAsync() is called, so a stop fails it while it is still looking up the last message ids, and a refresh started after the stop fails before any lookup, empty topic included. Your test is in for close, a closed reader and a rejected read; the lookup future stays incomplete in all three.

While testing this against a broker I found one more hole. A retry already queued on the client's scheduler is dropped silently when that scheduler shuts down; only a schedule() call made afterwards is rejected. Nothing is reading while a retry waits, so the loop ended without a stop cause and a waiting refresh stayed pending. I reproduced it by closing the client during a 6 s retry delay.

d6c30f3 moves the delay to CompletableFuture.delayedExecutor, like ScalableTopicProducer does, so the retry still runs and finds the closed reader. Its read now starts on a common-pool thread. The rejected-schedule branch is gone; two rules replace it: a read rejected because the reader's executor was shut down stops the loop, and a failure is not retried once the client is closed.

…ast message ids

stopTailReads() only reached the refreshes registered in
pendingRefreshRequests, so a refresh that was still fetching the last
message ids was left alone until the lookup answered, and a refresh
started after the stop went through the lookup first and even
succeeded on an empty topic. Track a refresh from the moment
refreshAsync() is called: it is added to the tracked set before it
reads the stop cause, and the stop sets the cause before it goes
through the set. The stop cause is checked before the lookup and
before the empty-topic shortcut. A refresh leaves the tracked set and
pendingRefreshRequests when it completes, whoever completes it.

The stop check moves into readTailMessages(), so the success path and
start() do not issue a tail read after a stop either, and the first
stop cause is kept.
A retry waiting on the client's scheduler is discarded without notice
when that scheduler is shut down, with the client or, when it is
shared, under it: shutdownNow() returns the queued tasks without
running or cancelling them, and schedule() only throws for a task
submitted afterwards. No read is in flight while a retry waits, so the
tail reads ended there with no stop cause, and a refresh waiting for a
message stayed pending until the table view was closed. On master the
failed read was retried after sleeping 50 ms on the same thread, so a
closed reader was found.

Run the retry through CompletableFuture.delayedExecutor() instead,
which does not depend on the client. The retry then finds the closed
reader, and the rejected-schedule branch goes away with the scheduler.

Two cases the retry meets now that it outlives the client's executors:

- The executors were shut down without the reader being closed
  (PulsarClientImpl#shutdown() does not close the consumers, and
  shared resources can be closed under the client). The reader then
  throws RejectedExecutionException from readNextAsync(). The read is
  invoked through FutureUtil.supplySafely() and a rejected read stops
  the tail reads like a closed reader. Any other failure thrown by the
  read is retried like a failed read.
- The client was closed but the reader was left failed rather than
  closed, and answers NotConnectedException for good. A failure is not
  retried once the client is closed.

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thanks for working on this fix and for the careful handling of the close races; the stop cause plus the active-refresh set makes every refresh settle on shutdown. One non-blocking suggestion on releasing the delayed retry when the table view is closed.

.exception(ex)
.log("Reader was interrupted while reading tail messages. Retrying..");
// readTailMessages() checks the stop cause itself: a retry still waiting at a stop reads nothing.
runAfterDelay(delayMillis, () -> readTailMessages(reader));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[SUGGESTION] A pending retry keeps a closed table view reachable until its delay elapses

The lambda queued by runAfterDelay at AbstractTableViewImpl.java:618 captures this and reader. closeAsync() sets the stop cause, but the delayed task stays in the JDK delay queue until it fires, so a closed table view (data map, listeners, reader, client) stays reachable for up to the current backoff delay, about a minute at the default maximum. The retry itself is harmless, since readTailMessages returns on the stop cause. Not blocking; if it is easy, keep the Future of the delayed task (for example via CompletableFuture.runAsync(retry, CompletableFuture.delayedExecutor(...))) and cancel it in stopTailReads, or document that bounded retention is accepted.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for catching this, you're right that the closed view stays reachable from the queued task. I tried your suggestion first, keeping the future from runAsync() and cancelling it in stopTailReads(), but as far as I can tell it doesn't release the view: the task stays in the JDK delayer's queue until the delay is up and still references the retry, so the view is held for the same time. I checked with a weak reference to a closed view, and it's only collected once the delay has passed. Please let me know if I missed something there.

One thing that did work in my test is keeping the retry in a reference that stopTailReads() clears, so the queued task only holds that reference and not the view. I didn't want to push more changes to an approved PR without asking, so I put it on a branch in my fork for now: namest504/pulsar@d31fbf0909 (diff: namest504/pulsar@fix/tableview-tail-retry...fix/tableview-tail-retry-holder). It also stops a failure that arrives after the close, from a read that was still in flight, from queueing a new retry.

If you'd like it in this PR, I'll push it. If you'd rather keep the PR as it is, the retention is already noted in the description. Whatever you think is best.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug] TableView tail-read retry sleeps on the client callback thread and can recurse without bound

3 participants