apache / apache/pulsar

[Bug] Consumer subscription failed. If some partitions on a topic are on a dead broker

Open
#18,743 3 comments 0 reactions 0 assignees View on GitHub
Stale type/bug
Dominant language
Java
Stars
15.3k
Forks
3.8k
Avg merge
1d 14h
Merged PRs (30d)
160

Description

### Search before asking

- [X] I searched in the [issues](https://github.com/apache/pulsar/issues) and found nothing similar.

### Version

2.8.1

### Minimal reproduce step

1、A topic has multiple partitions, and the pulsar io thread is in the sleep or block state on a broker.
2、Create consumers on topic.

```
PulsarClient client = PulsarClient.builder().serviceUrl(localClusterUrl).build();
ConsumerBuilder consumerBuilder = client.newConsumer(Schema.STRING);
List topics = new ArrayList<>();
topics.add("persistent://public/default/test-string1");
consumerBuilder.topics(topics)
.subscriptionName("test")
.subscriptionType(SubscriptionType.Shared)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
```

### What did you expect to see?

I hope that consumers can be successfully created and messages can be received on other normal partitions except the partition on the problematic broker that cannot create consumers.

### What did you see instead?

```
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-1][test] Subscribing to topic on cnx [id: 0xa363196f, L:/172.32.147.245:4291 - R:/172.32.149.121:16650], consumerId 1
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConnectionPool - [[id: 0x029dd152, L:/172.32.147.245:4292 - R:/172.32.149.122:16650]] Connected to server
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConnectionPool - [[id: 0xcea16deb, L:/172.32.147.245:4293 - R:/172.32.149.123:16650]] Connected to server
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-0][test] Subscribing to topic on cnx [id: 0x029dd152, L:/172.32.147.245:4292 - R:/172.32.149.122:16650], consumerId 0
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-2][test] Subscribing to topic on cnx [id: 0xcea16deb, L:/172.32.147.245:4293 - R:/172.32.149.123:16650], consumerId 2
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-1][test] Subscribed to topic on /172.32.149.121:16650 -- consumer: 1
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-0][test] Subscribed to topic on /172.32.149.122:16650 -- consumer: 0
[pulsar-client-io-1-1] WARN org.apache.pulsar.common.protocol.PulsarHandler - [[id: 0xcea16deb, L:/172.32.147.245:4293 - R:/172.32.149.123:16650]] Forcing connection to close after keep-alive timeout
[pulsar-client-io-1-1] WARN org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-2][test] Failed to subscribe to topic on /172.32.149.123:16650
[pulsar-client-io-1-1] WARN org.apache.pulsar.client.impl.MultiTopicsConsumerImpl - [persistent://public/default/test-string1] Failed to subscribe for topic [persistent://public/default/test-string1] in topics consumer org.apache.pulsar.client.api.PulsarClientException$TimeoutException: 6 request timedout after ms 30000
[pulsar-client-io-1-1] WARN org.apache.pulsar.client.impl.ClientCnx - [id: 0xcea16deb, L:/172.32.147.245:4293 ! R:/172.32.149.123:16650] 6 request timedout after ms 30000
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ClientCnx - [id: 0xcea16deb, L:/172.32.147.245:4293 ! R:/172.32.149.123:16650] Disconnected
[pulsar-external-listener-3-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-2] [test] Closed Consumer (not connected)
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-0] [test] Closed consumer
[pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerImpl - [persistent://public/default/test-string1-partition-1] [test] Closed consumer
[pulsar-client-internal-4-1] WARN org.apache.pulsar.client.impl.MultiTopicsConsumerImpl - [persistent://public/default/test-string1] Failed to subscribe for topic [persistent://public/default/test-string1] in topics consumer, subscribe error: org.apache.pulsar.client.api.PulsarClientException$TimeoutException: 6 request timedout after ms 30000
[pulsar-client-internal-4-1] WARN org.apache.pulsar.client.impl.MultiTopicsConsumerImpl - Failed subscription for createPartitionedConsumer: persistent://public/default/test-string1 3, e:{}
java.util.concurrent.CompletionException: org.apache.pulsar.client.api.PulsarClientException$TimeoutException: 6 request timedout after ms 30000
at java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:292)
at java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:308)
at java.util.concurrent.CompletableFuture.uniRun(CompletableFuture.java:714)
at java.util.concurrent.CompletableFuture$UniRun.tryFire(CompletableFuture.java:701)
at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)
at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1990)
at org.apache.pulsar.client.impl.ClientCnx.checkRequestTimeout(ClientCnx.java:1145)
at org.apache.pulsar.client.impl.ClientCnx.lambda$channelActive$0(ClientCnx.java:211)
at org.apache.pulsar.shade.io.netty.util.concurrent.PromiseTask.runTask(PromiseTask.java:98)
at org.apache.pulsar.shade.io.netty.util.concurrent.ScheduledFutureTask.run(ScheduledFutureTask.java:176)
at org.apache.pulsar.shade.io.netty.util.concurrent.AbstractEventExecutor.safeExecute(AbstractEventExecutor.java:164)
at org.apache.pulsar.shade.io.netty.util.concurrent.SingleThreadEventExecutor.runAllTasks(SingleThreadEventExecutor.java:469)
at org.apache.pulsar.shade.io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:500)
at org.apache.pulsar.shade.io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986)
at org.apache.pulsar.shade.io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
at org.apache.pulsar.shade.io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.pulsar.client.api.PulsarClientException$TimeoutException: 6 request timedout after ms 30000
... 11 more
[pulsar-client-io-6-1] INFO org.apache.pulsar.client.impl.ConnectionPool - [[id: 0xcdc97132, L:/172.32.147.245:4359 - R:/172.32.149.121:16650]] Connected to server
consumer is ? :null
create consumer fail: Failed to subscribe persistent://public/default/test-string1 with 3 partitions
6 request timedout after ms 30000
```

### Anything else?

_No response_

### Are you willing to submit a PR?

- [ ] I'm willing to submit a PR!

Contributor guide

Open the contributing guide

Research direction

Start with the minimal Java client reproduction and trace partitioned subscription handling through MultiTopicsConsumerImpl and ConsumerImpl, using the ClientCnx timeout and ConnectionPool logs as context. Verify behavior when one partition's broker is unresponsive: healthy partitions should remain consumable while the failed partition is handled without failing the whole subscription.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.