apache / apache/pulsar

Send to retry letter topic exception with topic

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

Description

**Describe the bug**
Pulsar v2.6.2

```
org.apache.pulsar.client.api.PulsarClientException$InvalidMessageException: null
at org.apache.pulsar.client.impl.ConsumerBase.reconsumeLaterAsync(ConsumerBase.java:380) ~[pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.ConsumerBase.reconsumeLater(ConsumerBase.java:306) ~[pulsar-client-2.6.2.jar:2.6.2]
at com.newland.boss.nlstreamboss.application.server.housekeeping.bean.PulsarClientFactory.reconsumeLater(PulsarClientFactory.java:102) ~[housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageResource.igniteConsumerMessage(MessageResource.java:188) [housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageResource.dealMessage(MessageResource.java:63) [housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageProvider.run(MessageProvider.java:25) [housekeeping-1.0.0.jar:1.0.0]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_151]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_151]
at java.lang.Thread.run(Thread.java:748) [na:1.8.0_151]

2021-09-12 13:02:23.472 ERROR 107874 --- [ekeeping-pool-7] o.a.pulsar.client.impl.ConsumerImpl : Send to retry letter topic exception with topic: persistent://pulsar/iagw/21_TX_IGN_DLQ, messageId: 8281593:4196:9

org.apache.pulsar.client.api.PulsarClientException$TimeoutException: The producer pulsar_jp-35-43945 can not send message to the topic persistent://pulsar/iagw/21_TX_IGN_RETRY-partition-7 within given timeout
at org.apache.pulsar.client.api.PulsarClientException.unwrap(PulsarClientException.java:856) ~[pulsar-client-api-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.ProducerBase.send(ProducerBase.java:117) ~[pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.TypedMessageBuilderImpl.send(TypedMessageBuilderImpl.java:89) ~[pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.ConsumerImpl.doReconsumeLater(ConsumerImpl.java:671) ~[pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.MultiTopicsConsumerImpl.doReconsumeLater(MultiTopicsConsumerImpl.java:456) [pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.ConsumerBase.reconsumeLaterAsync(ConsumerBase.java:378) [pulsar-client-2.6.2.jar:2.6.2]
at org.apache.pulsar.client.impl.ConsumerBase.reconsumeLater(ConsumerBase.java:306) [pulsar-client-2.6.2.jar:2.6.2]
at com.newland.boss.nlstreamboss.application.server.housekeeping.bean.PulsarClientFactory.reconsumeLater(PulsarClientFactory.java:102) [housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageResource.igniteConsumerMessage(MessageResource.java:188) [housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageResource.dealMessage(MessageResource.java:63) [housekeeping-1.0.0.jar:1.0.0]
at com.newland.boss.nlstreamboss.application.server.housekeeping.threads.MessageProvider.run(MessageProvider.java:25) [housekeeping-1.0.0.jar:1.0.0]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_151]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_151]
at java.lang.Thread.run(Thread.java:748) [na:1.8.0_151]
```
Topic Storage Size Cannot release.
![image](https://user-images.githubusercontent.com/54351417/133034107-d96202f7-e1e6-45fd-8a75-a95178aeca11.png)

**To Reproduce**
Steps to reproduce the behavior:
1.create consumer.
```
String[] topic = topicTmp.split("&");
Consumer consumer = pulsarClient.newConsumer(Schema.JSON(NlMessage.class))
.topic(topic[0])
.enableRetry(true)
.deadLetterPolicy(DeadLetterPolicy.builder()
// 注意这里使用公用的重试和死信队列
.retryLetterTopic(pulsarConfig.getRetryTopic())
.maxRedeliverCount(pulsarConfig.getMaxDeliveryCount())
.deadLetterTopic(pulsarConfig.getDeadLetterTopic())
.build())
.subscriptionName(pulsarConfig.getSubscription())
.subscriptionType(SubscriptionType.Shared)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.ackTimeout(0, TimeUnit.SECONDS)
.subscribe();
```
2.10 threads share a consumer to receive messages.
3.If the message processing fails, call the recommerlater interface.

Contributor guide

Open the contributing guide

Research direction

The relevant entry points are ConsumerImpl.doReconsumeLater, MultiTopicsConsumerImpl.doReconsumeLater, and ConsumerBase.reconsumeLaterAsync; begin by tracing the retry-letter path for the shared consumer configuration shown. Reproduce with retry and dead-letter policies enabled, multiple threads, and reconsumeLater after processing failure. Done means retry publishing no longer produces the reported InvalidMessageException or timeout, and topic storage can be released.

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
32/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.