Repository navigation
Conversation
eb6bb2f to
fa9ddf5
Compare
…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.
fa9ddf5 to
02383af
Compare
| try { | ||
| ((ScheduledExecutorService) client.getScheduledExecutorProvider().getExecutor()) | ||
| .schedule(() -> readTailMessages(reader), delayMillis, TimeUnit.MILLISECONDS); | ||
| } catch (RejectedExecutionException e) { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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(); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| } 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")); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
Two things I noticed while working on this. I left them out on purpose, to keep this PR limited to how the loop waits:
|
| @Override | ||
| public CompletableFuture<Void> refreshAsync() { | ||
| CompletableFuture<Void> completableFuture = new CompletableFuture<>(); | ||
| reader.thenCompose(reader -> getLastMessageIdOfNonEmptyTopics(reader).thenAccept(lastMessageIds -> { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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)); |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
Fixes #26671
Main Issue: #26671
Motivation
When a tail read fails with anything other than
AlreadyClosedException,AbstractTableViewImpl.readTailMessages()callsThread.sleep(50)inside theexceptionallycallback 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 whenreadNextAsync()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
Backoffdelay, built from the client'sstartingBackoffInterval/maxBackoffInterval, throughCompletableFuture.delayedExecutorinstead 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.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.ScalableTopicProducertimes its retries the same way.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.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 anAlreadyClosedExceptionwrapped in aCompletionException, the shape the failed read of a closed reader delivers, so callbacks that look atgetCause()keep working.FutureUtil.supplySafely, so a reader that throws instead of returning a failed future is handled like a failed read.Behaviour changes:
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.createAsync(): it is handled like any other failed read.What this change does not cover:
maxBackoffInterval. The retry itself does nothing then. The same held for a task queued on the client's scheduler before this change.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.Verifying this change
TableViewImplTestcovers the retry timing (testTailReadFailureSchedulesRetryInsteadOfRecursinghangs 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,MappedTableViewImplTestand, inpulsar-broker,TableViewTest,ServiceUnitStateChannelTest,LoadDataStoreTestandServiceUnitStateCompactionTest, 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
Tail-read retries of a
TableViewno 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-neededMatching PR in forked repository
PR in forked repository: namest504#2