Reader unusable: Exclusive consumer is already connected
- Dominant language
- Java
- Stars
- 15.3k
- Forks
- 3.8k
- Avg merge
- 1d 14h
- Merged PRs (30d)
- 160
Description
**Describe the bug**
We have a topic with a more or less complex AVRO schema and quite a few messages on it (38345) . We read the topic completely using Reader interface.
This has been **tested on 2.7.4**.
The first time we create the Reader everything goes fine. The second time we manage to read all messages but some error messages are seen on Pulsar log (tested on standalone, but the same situation happens in a AKS cluster). At this time we discover that the Reader subscription is never deleted (stays there). Next times we try to repeat the operation client fails to return the new Reader instance with an error ("Exclusive consumer is already connected"). A new subscription is left open for each try.
This is the error trace:
```
java.util.concurrent.CompletionException: org.apache.pulsar.client.api.PulsarClientException$ConsumerBusyException: Exclusive consumer is already connected
at java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:292) [na:1.8.0_111]
at java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:308) [na:1.8.0_111]
at java.util.concurrent.CompletableFuture.uniRun(CompletableFuture.java:700) [na:1.8.0_111]
at java.util.concurrent.CompletableFuture$UniRun.tryFire$$$capture(CompletableFuture.java:687) ~[na:1.8.0_111]
at java.util.concurrent.CompletableFuture$UniRun.tryFire(CompletableFuture.java) ~[na:1.8.0_111]
at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) [na:1.8.0_111]
at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977) [na:1.8.0_111]
at org.apache.pulsar.client.impl.ClientCnx.handleError(ClientCnx.java:655) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.common.protocol.PulsarDecoder.channelRead(PulsarDecoder.java:197) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:324) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:296) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:719) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:655) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:581) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:493) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) ~[pulsar-client-2.7.4.jar:2.7.4]
at org.apache.pulsar.shade.io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) ~[pulsar-client-2.7.4.jar:2.7.4]
at java.lang.Thread.run(Thread.java:745) ~[na:1.8.0_111]
Caused by: org.apache.pulsar.client.api.PulsarClientException$ConsumerBusyException: Exclusive consumer is already connected
at org.apache.pulsar.client.impl.ClientCnx.getPulsarClientException(ClientCnx.java:1027) ~[pulsar-client-2.7.4.jar:2.7.4]
... 23 common frames omitted
```
These are the **logs** on standalone Docker container on **first (successful operation)**:
```
standalone_1 | 13:19:54.093 [pulsar-web-67-8] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/namespaces/dbus HTTP/1.1" 200 51 "-" "Pulsar-Java-v2.7.4" 21
standalone_1 | 13:19:54.101 [pulsar-web-67-13] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo/partitioned HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 6
standalone_1 | 13:19:54.119 [pulsar-web-67-9] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/partitioned HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 15
standalone_1 | 13:19:54.132 [pulsar-web-67-1] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo HTTP/1.1" 200 66 "-" "Pulsar-Java-v2.7.4" 8
standalone_1 | 13:19:54.145 [pulsar-web-67-12] INFO org.apache.pulsar.broker.PulsarService - created admin with url http://localhost:8080
standalone_1 | 13:19:54.210 [ForkJoinPool.commonPool-worker-5] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x00000000_0x40000000
standalone_1 | 13:19:54.211 [ForkJoinPool.commonPool-worker-3] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x40000000_0x80000000
standalone_1 | 13:19:54.211 [ForkJoinPool.commonPool-worker-6] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0xc0000000_0xffffffff
standalone_1 | 13:19:54.213 [ForkJoinPool.commonPool-worker-4] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x80000000_0xc0000000
standalone_1 | 13:19:54.216 [ForkJoinPool.commonPool-worker-3] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x40000000_0x80000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 10
standalone_1 | 13:19:54.220 [ForkJoinPool.commonPool-worker-5] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x00000000_0x40000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 14
standalone_1 | 13:19:54.220 [ForkJoinPool.commonPool-worker-6] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0xc0000000_0xffffffff HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 13
standalone_1 | 13:19:54.221 [ForkJoinPool.commonPool-worker-4] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x80000000_0xc0000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 14
standalone_1 | 13:19:54.227 [AsyncHttpClient-99-1] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 104
standalone_1 | 13:19:54.235 [pulsar-30-8] INFO org.apache.pulsar.broker.namespace.OwnershipCache - Trying to acquire ownership of dbus/traffic-events-demo/0xc0000000_0xffffffff
standalone_1 | 13:19:54.238 [pulsar-ordered-OrderedExecutor-0-0-EventThread] INFO org.apache.pulsar.broker.namespace.OwnershipCache - Successfully acquired ownership of /namespace/dbus/traffic-events-demo/0xc0000000_0xffffffff
standalone_1 | 13:19:54.256 [bookie-io-1-3] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.262 [bookie-io-1-3] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x4c4e415d, L:/127.0.0.1:3181 - R:/127.0.0.1:42866]
standalone_1 | 13:19:54.262 [bookkeeper-io-63-1] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x91eee5d3, L:/127.0.0.1:42866 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.262 [bookkeeper-io-63-1] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x91eee5d3, L:/127.0.0.1:42866 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.262 [bookie-io-1-4] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.262 [bookie-io-1-4] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x2894edfa, L:/127.0.0.1:3181 - R:/127.0.0.1:42868]
standalone_1 | 13:19:54.265 [bookie-io-1-5] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.265 [bookkeeper-io-63-3] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0xfd45b68c, L:/127.0.0.1:42868 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.265 [bookie-io-1-5] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x87b2d164, L:/127.0.0.1:3181 - R:/127.0.0.1:42870]
standalone_1 | 13:19:54.266 [bookkeeper-io-63-3] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0xfd45b68c, L:/127.0.0.1:42868 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.265 [bookkeeper-io-63-2] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x3d3b7e0e, L:/127.0.0.1:42870 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.266 [bookkeeper-io-63-2] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x3d3b7e0e, L:/127.0.0.1:42870 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.267 [bookie-io-1-6] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.267 [bookie-io-1-6] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x5681c9c5, L:/127.0.0.1:3181 - R:/127.0.0.1:42874]
standalone_1 | 13:19:54.267 [bookie-io-1-7] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.267 [bookie-io-1-7] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x47fbc2ba, L:/127.0.0.1:3181 - R:/127.0.0.1:42872]
standalone_1 | 13:19:54.268 [bookie-io-1-8] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.268 [bookie-io-1-8] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x8417212a, L:/127.0.0.1:3181 - R:/127.0.0.1:42876]
standalone_1 | 13:19:54.268 [bookie-io-1-9] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.271 [bookie-io-1-9] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x1b76672d, L:/127.0.0.1:3181 - R:/127.0.0.1:42878]
standalone_1 | 13:19:54.271 [bookkeeper-io-63-4] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x6ca18a24, L:/127.0.0.1:42872 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.272 [bookkeeper-io-63-4] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x6ca18a24, L:/127.0.0.1:42872 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.272 [bookie-io-1-10] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.272 [bookie-io-1-10] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x7819cb70, L:/127.0.0.1:3181 - R:/127.0.0.1:42880]
standalone_1 | 13:19:54.272 [bookie-io-1-11] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.272 [bookie-io-1-11] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0xa8e654fb, L:/127.0.0.1:3181 - R:/127.0.0.1:42882]
standalone_1 | 13:19:54.274 [bookie-io-1-12] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.274 [bookie-io-1-12] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x09d7284d, L:/127.0.0.1:3181 - R:/127.0.0.1:42884]
standalone_1 | 13:19:54.277 [bookkeeper-io-63-6] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0xfe00b9b6, L:/127.0.0.1:42876 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.277 [bookkeeper-io-63-5] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x2c04e458, L:/127.0.0.1:42874 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.279 [bookie-io-1-16] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.279 [bookie-io-1-16] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x1780420c, L:/127.0.0.1:3181 - R:/127.0.0.1:42892]
standalone_1 | 13:19:54.279 [bookie-io-1-15] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.279 [bookie-io-1-14] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.280 [bookie-io-1-15] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0xde831a1d, L:/127.0.0.1:3181 - R:/127.0.0.1:42890]
standalone_1 | 13:19:54.280 [bookie-io-1-14] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x8c25e46f, L:/127.0.0.1:3181 - R:/127.0.0.1:42888]
standalone_1 | 13:19:54.280 [bookkeeper-io-63-8] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x68dab97c, L:/127.0.0.1:42882 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.279 [bookie-io-1-13] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.280 [bookkeeper-io-63-8] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x68dab97c, L:/127.0.0.1:42882 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.281 [bookie-io-1-13] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x00861e2a, L:/127.0.0.1:3181 - R:/127.0.0.1:42886]
standalone_1 | 13:19:54.282 [bookkeeper-io-63-5] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x2c04e458, L:/127.0.0.1:42874 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.277 [bookkeeper-io-63-7] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0xc350b6a1, L:/127.0.0.1:42878 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.282 [bookkeeper-io-63-7] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0xc350b6a1, L:/127.0.0.1:42878 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.282 [bookie-io-1-1] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.282 [bookie-io-1-1] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0xa7352d8d, L:/127.0.0.1:3181 - R:/127.0.0.1:42894]
standalone_1 | 13:19:54.279 [bookkeeper-io-63-6] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0xfe00b9b6, L:/127.0.0.1:42876 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.288 [bookkeeper-io-63-10] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0xbb5cc097, L:/127.0.0.1:42884 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.288 [bookkeeper-io-63-9] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x03c5bc62, L:/127.0.0.1:42880 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.288 [bookkeeper-io-63-10] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0xbb5cc097, L:/127.0.0.1:42884 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.288 [bookkeeper-io-63-9] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x03c5bc62, L:/127.0.0.1:42880 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.288 [bookkeeper-io-63-13] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x4c063f7f, L:/127.0.0.1:42890 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.288 [bookkeeper-io-63-13] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x4c063f7f, L:/127.0.0.1:42890 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.288 [bookkeeper-io-63-11] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0xd021a631, L:/127.0.0.1:42888 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.288 [bookkeeper-io-63-14] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x35dea621, L:/127.0.0.1:42892 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.289 [bookie-io-1-2] INFO org.apache.bookkeeper.proto.AuthHandler - Authentication success on server side
standalone_1 | 13:19:54.289 [bookkeeper-io-63-11] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0xd021a631, L:/127.0.0.1:42888 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.289 [bookie-io-1-2] INFO org.apache.bookkeeper.proto.BookieRequestHandler - Channel connected [id: 0x83b4eca7, L:/127.0.0.1:3181 - R:/127.0.0.1:42896]
standalone_1 | 13:19:54.289 [bookkeeper-io-63-14] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x35dea621, L:/127.0.0.1:42892 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.289 [bookkeeper-io-63-12] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x02d3527c, L:/127.0.0.1:42886 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.289 [bookkeeper-io-63-12] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x02d3527c, L:/127.0.0.1:42886 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.293 [bookkeeper-io-63-15] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x5a026f1a, L:/127.0.0.1:42894 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.293 [bookkeeper-io-63-16] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - Successfully connected to bookie: 127.0.0.1:3181 [id: 0x129a942f, L:/127.0.0.1:42896 - R:127.0.0.1/127.0.0.1:3181]
standalone_1 | 13:19:54.293 [bookkeeper-io-63-15] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x5a026f1a, L:/127.0.0.1:42894 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.293 [bookkeeper-io-63-16] INFO org.apache.bookkeeper.proto.PerChannelBookieClient - connection [id: 0x129a942f, L:/127.0.0.1:42896 - R:127.0.0.1/127.0.0.1:3181] authenticated as BookKeeperPrincipal{ANONYMOUS}
standalone_1 | 13:19:54.317 [ForkJoinPool.commonPool-worker-4] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/schemas/dbus/traffic-events-demo/traffic-situation-events/schema HTTP/1.1" 200 40459 "-" "Pulsar-Java-v2.7.4" 87
standalone_1 | 13:19:54.327 [pulsar-ordered-OrderedExecutor-0-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - Opening managed ledger dbus/traffic-events-demo/persistent/traffic-situation-events
standalone_1 | 13:19:54.330 [bookkeeper-ml-workers-OrderedExecutor-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Creating ledger, metadata: {component=[109, 97, 110, 97, 103, 101, 100, 45, 108, 101, 100, 103, 101, 114], pulsar/managed-ledger=[100, 98, 117, 115, 47, 116, 114, 97, 102, 102, 105, 99, 45, 101, 118, 101, 110, 116, 115, 45, 100, 101, 109, 111, 47, 112, 101, 114, 115, 105, 115, 116, 101, 110, 116, 47, 116, 114, 97, 102, 102, 105, 99, 45, 115, 105, 116, 117, 97, 116, 105, 111, 110, 45, 101, 118, 101, 110, 116, 115], application=[112, 117, 108, 115, 97, 114]} - metadata ops timeout : 60 seconds
standalone_1 | 13:19:54.337 [main-EventThread] INFO org.apache.bookkeeper.client.LedgerCreateOp - Ensemble: [127.0.0.1:3181] for ledger: 10155
standalone_1 | 13:19:54.337 [bookkeeper-ml-workers-OrderedExecutor-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Created ledger 10155
standalone_1 | 13:19:54.342 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Loading cursor __compaction
standalone_1 | 13:19:54.342 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Recovering from bookkeeper ledger cursor: __compaction
standalone_1 | 13:19:54.343 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Cursor __compaction recovered to position 9783:1813
standalone_1 | 13:19:54.343 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Recovery for cursor __compaction completed. pos=9783:1813 -- todo=0
standalone_1 | 13:19:54.343 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Successfully initialize managed ledger
standalone_1 | 13:19:54.343 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.pulsar.broker.service.AbstractTopic - Disabling publish throttling for persistent://dbus/traffic-events-demo/traffic-situation-events
standalone_1 | 13:19:54.344 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events] There are no replicated subscriptions on the topic
standalone_1 | 13:19:54.344 [bookkeeper-ml-scheduler-OrderedScheduler-1-0] INFO org.apache.pulsar.broker.service.BrokerService - Created topic persistent://dbus/traffic-events-demo/traffic-situation-events - dedup is disabled
standalone_1 | 13:19:54.354 [bookkeeper-ml-workers-OrderedExecutor-1-0] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo/traffic-situation-events/lastMessageId HTTP/1.1" 200 207 "-" "Pulsar-Java-v2.7.4" 31
standalone_1 | 13:19:54.861 [pulsar-web-67-14] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:19:54.868 [pulsar-web-67-13] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:19:54 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:20:04.847 [SyncThread-7-1] INFO org.apache.bookkeeper.bookie.Journal - garbage collected journal 17fe4c0a8f1.txn
```
Looking at the subscriptions only the "__compaction" one is alive (we have compaction enabled and try to read from the compacted topic always):
```
$ docker-compose exec standalone /pulsar/bin/pulsarctl subscriptions list persistent://dbus/traffic-events-demo/traffic-situation-events
+-------------------+
| SUBSCRIPTIONS |
+-------------------+
| __compaction |
+-------------------+
```
On **second invocation** we sucessfully read all messages but these are the logs from the Docker container:
```
standalone_1 | 13:21:19.676 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.ServerCnx - New connection from /172.16.2.1:55140
standalone_1 | 13:21:19.757 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55140] Subscribing on topic persistent://dbus/traffic-events-demo/traffic-situation-events / reader-a5525623e6
standalone_1 | 13:21:19.772 [ForkJoinPool.commonPool-worker-1] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events][reader-a5525623e6] Creating non-durable subscription at msg id -1:-1:-1:-1
standalone_1 | 13:21:19.772 [ForkJoinPool.commonPool-worker-1] INFO org.apache.bookkeeper.mledger.impl.NonDurableCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Created non-durable cursor read-position=9783:0 mark-delete-position=9783:-1
standalone_1 | 13:21:19.772 [ForkJoinPool.commonPool-worker-1] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Opened new cursor: NonDurableCursorImpl{ledger=dbus/traffic-events-demo/persistent/traffic-situation-events, ackPos=9783:-1, readPos=9783:0}
standalone_1 | 13:21:19.773 [ForkJoinPool.commonPool-worker-1] INFO org.apache.bookkeeper.mledger.impl.ManagedCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events-reader-a5525623e6] Rewind from 9783:0 to 9783:0
standalone_1 | 13:21:19.775 [ForkJoinPool.commonPool-worker-1] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events] There are no replicated subscriptions on the topic
standalone_1 | 13:21:19.776 [ForkJoinPool.commonPool-worker-1] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events][reader-a5525623e6] Created new subscription for 10
standalone_1 | 13:21:19.776 [ForkJoinPool.commonPool-worker-1] INFO org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55140] Created subscription on topic persistent://dbus/traffic-events-demo/traffic-situation-events / reader-a5525623e6
standalone_1 | 13:21:24.863 [pulsar-web-67-16] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:21:24 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 5
standalone_1 | 13:21:24.872 [pulsar-web-67-3] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:21:24 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 5
standalone_1 | 13:21:24.972 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55140] Closing consumer: consumerId=10
standalone_1 | 13:21:24.973 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer - Removing consumer Consumer{subscription=PersistentSubscription{topic=persistent://dbus/traffic-events-demo/traffic-situation-events, name=reader-a5525623e6}, consumerId=10, consumerName=0e946, address=/172.16.2.1:55140}
standalone_1 | 13:21:24.976 [pulsar-io-50-11] WARN org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55140] Got exception java.lang.NullPointerException
standalone_1 | at org.apache.bookkeeper.mledger.impl.ManagedCursorContainer.removeCursor(ManagedCursorContainer.java:128)
standalone_1 | at org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.deactivateCursor(ManagedLedgerImpl.java:3122)
standalone_1 | at org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.setInactive(ManagedCursorImpl.java:956)
standalone_1 | at org.apache.pulsar.broker.service.persistent.PersistentSubscription.deactivateCursor(PersistentSubscription.java:302)
standalone_1 | at org.apache.pulsar.broker.service.persistent.PersistentSubscription.removeConsumer(PersistentSubscription.java:255)
standalone_1 | at org.apache.pulsar.broker.service.Consumer.close(Consumer.java:300)
standalone_1 | at org.apache.pulsar.broker.service.Consumer.close(Consumer.java:296)
standalone_1 | at org.apache.pulsar.broker.service.ServerCnx.handleCloseConsumer(ServerCnx.java:1504)
standalone_1 | at org.apache.pulsar.common.protocol.PulsarDecoder.channelRead(PulsarDecoder.java:172)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
standalone_1 | at io.netty.handler.flow.FlowControlHandler.dequeue(FlowControlHandler.java:200)
standalone_1 | at io.netty.handler.flow.FlowControlHandler.channelRead(FlowControlHandler.java:162)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
standalone_1 | at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:324)
standalone_1 | at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:296)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
standalone_1 | at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
standalone_1 | at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
standalone_1 | at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919)
standalone_1 | at io.netty.channel.epoll.AbstractEpollStreamChannel$EpollStreamUnsafe.epollInReady(AbstractEpollStreamChannel.java:795)
standalone_1 | at io.netty.channel.epoll.EpollEventLoop.processReady(EpollEventLoop.java:480)
standalone_1 | at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:378)
standalone_1 | at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986)
standalone_1 | at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
standalone_1 | at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
standalone_1 | at java.lang.Thread.run(Thread.java:748)
standalone_1 |
standalone_1 | 13:21:24.977 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.ServerCnx - Closed connection from /172.16.2.1:55140
standalone_1 | 13:21:24.989 [pulsar-io-50-11] INFO org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer - Removing consumer Consumer{subscription=PersistentSubscription{topic=persistent://dbus/traffic-events-demo/traffic-situation-events, name=reader-a5525623e6}, consumerId=10, consumerName=0e946, address=/172.16.2.1:55140}
standalone_1 | 13:21:24.989 [pulsar-io-50-11] WARN org.apache.pulsar.broker.service.ServerCnx - Consumer Consumer{subscription=PersistentSubscription{topic=persistent://dbus/traffic-events-demo/traffic-situation-events, name=reader-a5525623e6}, consumerId=10, consumerName=0e946, address=/172.16.2.1:55140} was already closed: org.apache.pulsar.broker.service.BrokerServiceException$ServerMetadataException: Consumer was not connected
```
By looking now at the subscriptions on the topic we see the one related to the reader:
```
$ docker-compose exec standalone /pulsar/bin/pulsarctl subscriptions list persistent://dbus/traffic-events-demo/traffic-situation-events
+-------------------+
| SUBSCRIPTIONS |
+-------------------+
| reader-a5525623e6 |
| __compaction |
+-------------------+
```
On third an sucessive operations Pulsar client always fails with the reported error above and subscriptions keep on growing on each iteration. These are Docker container logs:
```
standalone_1 | 13:22:50.195 [pulsar-web-67-7] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/tenants HTTP/1.1" 200 26 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:22:50.200 [pulsar-web-67-11] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/namespaces/dbus HTTP/1.1" 200 51 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:22:50.205 [pulsar-web-67-14] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo/partitioned HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:22:50.213 [pulsar-web-67-5] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/partitioned HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:22:50.219 [pulsar-web-67-1] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo HTTP/1.1" 200 66 "-" "Pulsar-Java-v2.7.4" 4
standalone_1 | 13:22:50.230 [pulsar-web-67-6] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0xc0000000_0xffffffff HTTP/1.1" 200 66 "-" "Pulsar-Java-v2.7.4" 2
standalone_1 | 13:22:50.233 [ForkJoinPool.commonPool-worker-3] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x00000000_0x40000000
standalone_1 | 13:22:50.233 [ForkJoinPool.commonPool-worker-5] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x80000000_0xc0000000
standalone_1 | 13:22:50.233 [ForkJoinPool.commonPool-worker-4] INFO org.apache.pulsar.broker.admin.v2.NonPersistentTopics - [null] Namespace bundle is not owned by any broker dbus/traffic-events-demo/0x40000000_0x80000000
standalone_1 | 13:22:50.239 [ForkJoinPool.commonPool-worker-3] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x00000000_0x40000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 15
standalone_1 | 13:22:50.239 [ForkJoinPool.commonPool-worker-5] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x80000000_0xc0000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 11
standalone_1 | 13:22:50.239 [ForkJoinPool.commonPool-worker-4] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo/0x40000000_0x80000000 HTTP/1.1" 204 0 "-" "Pulsar-Java-v2.7.4" 15
standalone_1 | 13:22:50.242 [AsyncHttpClient-99-1] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/non-persistent/dbus/traffic-events-demo HTTP/1.1" 200 2 "-" "Pulsar-Java-v2.7.4" 25
standalone_1 | 13:22:50.255 [ForkJoinPool.commonPool-worker-5] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/schemas/dbus/traffic-events-demo/traffic-situation-events/schema HTTP/1.1" 200 40459 "-" "Pulsar-Java-v2.7.4" 10
standalone_1 | 13:22:50.264 [bookkeeper-ml-workers-OrderedExecutor-1-0] INFO org.eclipse.jetty.server.RequestLog - 172.16.2.1 - - [07/Apr/2022:13:22:50 +0200] "GET /admin/v2/persistent/dbus/traffic-events-demo/traffic-situation-events/lastMessageId HTTP/1.1" 200 207 "-" "Pulsar-Java-v2.7.4" 5
standalone_1 | 13:22:50.269 [pulsar-io-50-12] INFO org.apache.pulsar.broker.service.ServerCnx - New connection from /172.16.2.1:55168
standalone_1 | 13:22:50.297 [pulsar-io-50-12] INFO org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55168] Subscribing on topic persistent://dbus/traffic-events-demo/traffic-situation-events / reader-e490b7f71b
standalone_1 | 13:22:50.307 [ForkJoinPool.commonPool-worker-5] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events][reader-e490b7f71b] Creating non-durable subscription at msg id -1:-1:-1:-1
standalone_1 | 13:22:50.308 [ForkJoinPool.commonPool-worker-5] INFO org.apache.bookkeeper.mledger.impl.NonDurableCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Created non-durable cursor read-position=9783:0 mark-delete-position=9783:-1
standalone_1 | 13:22:50.308 [ForkJoinPool.commonPool-worker-5] INFO org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events] Opened new cursor: NonDurableCursorImpl{ledger=dbus/traffic-events-demo/persistent/traffic-situation-events, ackPos=9783:-1, readPos=9783:0}
standalone_1 | 13:22:50.308 [ForkJoinPool.commonPool-worker-5] INFO org.apache.bookkeeper.mledger.impl.ManagedCursorImpl - [dbus/traffic-events-demo/persistent/traffic-situation-events-reader-e490b7f71b] Rewind from 9783:0 to 9783:0
standalone_1 | 13:22:50.309 [ForkJoinPool.commonPool-worker-5] ERROR org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events] Failed to create subscription: reader-e490b7f71b error: java.util.concurrent.CompletionException: java.lang.NullPointerException
standalone_1 | 13:22:50.309 [ForkJoinPool.commonPool-worker-5] WARN org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55168][persistent://dbus/traffic-events-demo/traffic-situation-events][reader-e490b7f71b] Failed to create consumer: consumerId=11, java.util.concurrent.CompletionException: java.lang.NullPointerException
standalone_1 | 13:22:50.418 [pulsar-io-50-12] INFO org.apache.pulsar.broker.service.ServerCnx - [/172.16.2.1:55168] Subscribing on topic persistent://dbus/traffic-events-demo/traffic-situation-events / reader-e490b7f71b
standalone_1 | 13:22:50.429 [ForkJoinPool.commonPool-worker-5] INFO org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events][reader-e490b7f71b] Creating non-durable subscription at msg id -1:-1:-1:-1
standalone_1 | 13:22:50.429 [ForkJoinPool.commonPool-worker-5] WARN org.apache.pulsar.broker.service.persistent.PersistentTopic - [persistent://dbus/traffic-events-demo/traffic-situation-events][reader-e490b7f71b] Consumer 11 81d7a already connected
standalone_1 | 13:22:54.860 [pulsar-web-67-14] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:54 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 3
standalone_1 | 13:22:54.867 [pulsar-web-67-4] INFO org.eclipse.jetty.server.RequestLog - 127.0.0.1 - - [07/Apr/2022:13:22:54 +0200] "GET /admin/v2/persistent/public/functions/coordinate/stats?getPreciseBacklog=false&subscriptionBacklogSize=false HTTP/1.1" 200 1682 "-" "Pulsar-Java-v2.7.4" 4
```
An the subscriptions:
```
$ docker-compose exec standalone /pulsar/bin/pulsarctl subscriptions list persistent://dbus/traffic-events-demo/traffic-situation-events
+-------------------+
| SUBSCRIPTIONS |
+-------------------+
| reader-e490b7f71b |
| reader-a5525623e6 |
| __compaction |
+-------------------+
```
**To Reproduce**
Steps to reproduce the behavior:
Described above.
**Expected behavior**
Not failing and not leaving reader subscriptions behind.
**Screenshots**
**Desktop (please complete the following information):**
- WSL2 on Windows (standalone) and AKS
**Additional context**
Tested on 2.7.4
Contributor guide
Research direction
Start by reproducing the repeated Reader creation and subscription cleanup behavior described in the issue, then trace the client path through the mentioned ClientCnx.java and PulsarDecoder.java stack frames. Done means repeated complete reads no longer leave subscriptions behind or fail with ConsumerBusyException; the issue does not name a test to run.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- backend, distributed-systems
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 35/100