Search before reporting
Read release policy
User environment
Pulsar master (776248a7fe), Java 21. The code has been the same since #22064 (Feb 2024, released in 3.3.0).
Issue Description
TableViewImpl.readTailMessages() retries a failed tail read by calling Thread.sleep(50) inside the CompletableFuture.exceptionally callback and then calling itself again:
https://github.com/apache/pulsar/blob/776248a7fe/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TableViewImpl.java#L420-L431
That callback runs on whichever thread completed the failed read, and that is never a thread it is safe to block:
-
Shared client thread. The read future is completed on the consumer's internalPinnedExecutor (ConsumerImpl.java:1757), one of the client's numIoThreads internal executor threads shared by every consumer and producer of that client. Each retry would stall that thread for 50 ms, delaying whatever else is pinned to it. CODING.md explicitly forbids this: "Never block on event-loop / async-execution threads — no Thread.sleep...". This is the only Thread.sleep in the client's main code. On master the way to reach this branch is an exception thrown while applying a message (for example a schema decode failure in Message.getValue()), which the handler then treats as a reader failure.
-
Caller's own stack, unbounded recursion. If readNextAsync() ever returns an already-failed future, exceptionally runs synchronously on the caller's thread and readTailMessages calls itself on the same stack: sleep 50 ms, log a WARN with a stack trace, recurse, never return. ConsumerBase.receiveAsync() does return such a future when verifyConsumerState() throws (Failed/Uninitialized), although I could not find a way for a reader that has already been created to get into those states in the current code, so this is a latent path rather than something I have observed. It is reproducible in a unit test: with readNextAsync() failing immediately, the calling thread spins at ~20 retries per second and the number of readTailMessages frames in the logged stack trace grows linearly (2 on the first retry, 597 after 214 retries in 10 s).
-
No backoff, no upper bound. Every exception other than AlreadyClosedException is retried after a fixed 50 ms, with a WARN plus stack trace each time.
To be clear about impact: I have not seen this cause a production incident. The report is that a documented rule is violated on a path that is reachable, with a failure mode (unbounded synchronous recursion) that only needs one immediately-failing read to appear.
Error messages
Logged 20 times per second while the retry loop is spinning (from the unit test described below; the stack trace grows with every retry):
WARN o.a.p.c.i.TableViewImpl - Reader was interrupted while reading tail messages. Retrying..
{reader=null, topic=non-persistent://tenant/ns/tail-retry}
java.util.concurrent.CompletionException: org.apache.pulsar.client.api.PulsarClientException$NotConnectedException: Not connected to broker
at java.base/java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:332)
at java.base/java.util.concurrent.CompletableFuture.uniAcceptNow(CompletableFuture.java:747)
at java.base/java.util.concurrent.CompletableFuture.uniAcceptStage(CompletableFuture.java:735)
at java.base/java.util.concurrent.CompletableFuture.thenAccept(CompletableFuture.java:2214)
at org.apache.pulsar.client.impl.TableViewImpl.readTailMessages(TableViewImpl.java:408)
at org.apache.pulsar.client.impl.TableViewImpl.lambda$start$1(TableViewImpl.java:122)
...
at org.apache.pulsar.client.impl.TableViewImpl.readTailMessages(TableViewImpl.java:430)
at org.apache.pulsar.client.impl.TableViewImpl.lambda$readTailMessages$14(TableViewImpl.java:430)
...
Reproducing the issue
Unit test against TableViewImpl with a mocked Reader whose readNextAsync() returns FutureUtil.failedFuture(new PulsarClientException.NotConnectedException()) on every call, on a non-persistent topic so that start() invokes readTailMessages on the test thread. start() never returns; with a 10 s TestNG timeout the run ends with ThreadTimeoutException after 214 retries.
Additional information
Suggested fix: schedule the retry on client.getScheduledExecutorProvider() with the client's configured reconnection backoff (initialBackoffInterval doubling up to maxBackoffInterval), reset the backoff after a successful read, and stop when the client's executor is shut down. That is the pattern ConsumerImpl already uses for delayed retries (for example in seekAsyncInternal).
Are you willing to submit a PR?
Search before reporting
Read release policy
User environment
Pulsar master (
776248a7fe), Java 21. The code has been the same since #22064 (Feb 2024, released in 3.3.0).Issue Description
TableViewImpl.readTailMessages()retries a failed tail read by callingThread.sleep(50)inside theCompletableFuture.exceptionallycallback and then calling itself again:https://github.com/apache/pulsar/blob/776248a7fe/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TableViewImpl.java#L420-L431
That callback runs on whichever thread completed the failed read, and that is never a thread it is safe to block:
Shared client thread. The read future is completed on the consumer's
internalPinnedExecutor(ConsumerImpl.java:1757), one of the client'snumIoThreadsinternal executor threads shared by every consumer and producer of that client. Each retry would stall that thread for 50 ms, delaying whatever else is pinned to it.CODING.mdexplicitly forbids this: "Never block on event-loop / async-execution threads — noThread.sleep...". This is the onlyThread.sleepin the client's main code. On master the way to reach this branch is an exception thrown while applying a message (for example a schema decode failure inMessage.getValue()), which the handler then treats as a reader failure.Caller's own stack, unbounded recursion. If
readNextAsync()ever returns an already-failed future,exceptionallyruns synchronously on the caller's thread andreadTailMessagescalls itself on the same stack: sleep 50 ms, log a WARN with a stack trace, recurse, never return.ConsumerBase.receiveAsync()does return such a future whenverifyConsumerState()throws (Failed/Uninitialized), although I could not find a way for a reader that has already been created to get into those states in the current code, so this is a latent path rather than something I have observed. It is reproducible in a unit test: withreadNextAsync()failing immediately, the calling thread spins at ~20 retries per second and the number ofreadTailMessagesframes in the logged stack trace grows linearly (2 on the first retry, 597 after 214 retries in 10 s).No backoff, no upper bound. Every exception other than
AlreadyClosedExceptionis retried after a fixed 50 ms, with a WARN plus stack trace each time.To be clear about impact: I have not seen this cause a production incident. The report is that a documented rule is violated on a path that is reachable, with a failure mode (unbounded synchronous recursion) that only needs one immediately-failing read to appear.
Error messages
Logged 20 times per second while the retry loop is spinning (from the unit test described below; the stack trace grows with every retry):
Reproducing the issue
Unit test against
TableViewImplwith a mockedReaderwhosereadNextAsync()returnsFutureUtil.failedFuture(new PulsarClientException.NotConnectedException())on every call, on a non-persistent topic so thatstart()invokesreadTailMessageson the test thread.start()never returns; with a 10 s TestNG timeout the run ends withThreadTimeoutExceptionafter 214 retries.Additional information
Suggested fix: schedule the retry on
client.getScheduledExecutorProvider()with the client's configured reconnection backoff (initialBackoffIntervaldoubling up tomaxBackoffInterval), reset the backoff after a successful read, and stop when the client's executor is shut down. That is the patternConsumerImplalready uses for delayed retries (for example inseekAsyncInternal).Are you willing to submit a PR?