apache / apache/pulsar

[Bug] Message expiration affects connected consumer

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

Description

### Search before reporting

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

### Read release policy

- [x] I understand that [unsupported versions](https://pulsar.apache.org/contribute/release-policy/#supported-versions) don't get bug fixes. I will attempt to reproduce the issue on a supported version of Pulsar client and Pulsar broker.

### User environment

- Pulsar version: 3.0.x

### Issue Description

When a consumer is connected and actively receiving messages, invoking `Subscription.expireMessages(int expireTimeInSeconds)` should not interfere with its message delivery pipeline. This method is intended to expire messages for disconnected consumers, not to remove data visible to a live consumer.

However, the expiring messages feature causes expired messages to become inaccessible to the connected consumer, making them appear lost, even though they are still present in the topic.

### Error messages

```text

```

### Reproducing the issue

You can add and run this test in the `org.apache.pulsar.client.api.SimpleProducerConsumerTest`.

```
@Test
public void testConsumeWhenMessageExpired() throws Exception {
String topicName = "persistent://public/default/testConsumeWhenExpired";
String subName = "sub";
int receiverQueueSize = 1;

@Cleanup
Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subName)
.receiverQueueSize(receiverQueueSize).subscribe();
@Cleanup
Producer producer = pulsarClient.newProducer().topic(topicName).enableBatching(false).create();

List messageIds = new ArrayList<>();
messageIds.add(producer.send(("my-message-0-ok").getBytes()));
messageIds.add(producer.send(("my-message-1-ok").getBytes()));
Thread.sleep(3 * 1000);
messageIds.add(producer.send(("my-message-2-ok").getBytes()));
messageIds.add(producer.send(("my-message-3-ok").getBytes()));
Thread.sleep(3 * 1000);
// Reset current position to 5 seconds ago.
consumer.seek(System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(5));
messageIds.add(producer.send(("my-message-4-expired").getBytes()));
messageIds.add(producer.send(("my-message-5-expired").getBytes()));
// Sleep 5 seconds to wait the messages to be expired.
Thread.sleep(5 * 1000);
messageIds.add(producer.send(("my-message-6-ok").getBytes()));
Optional topicReference = pulsar.getBrokerService().getTopicReference(topicName);
assertThat(topicReference).isPresent();
// Expire messages.
topicReference.ifPresent(topic -> {
Subscription subscription = topic.getSubscription(subName);
subscription.expireMessages(5);
});

List receivedMessages = new ArrayList<>();
while (true) {
Message receive = consumer.receive(2, TimeUnit.SECONDS);
if (receive == null) {
break;
}
String msg = new String(receive.getData(), StandardCharsets.UTF_8);
receivedMessages.add(msg);
}

List expectedMessages = Lists.newArrayList("my-message-2-ok", "my-message-3-ok", "my-message-4-expired", "my-message-5-expired", "my-message-6-ok");
assertThat(receivedMessages)
.isEqualTo(expectedMessages);

// Debug
System.out.println("All messageIds:");
messageIds.forEach(System.out::println);
}
```

### Additional information

_No response_

### Are you willing to submit a PR?

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

Contributor guide

Open the contributing guide

Research direction

Start with org.apache.pulsar.client.api.SimpleProducerConsumerTest and run the provided test scenario on Pulsar 3.0.x. Trace Subscription.expireMessages(5) and compare delivery for the connected consumer. Done means the received messages match the expected list, including the messages marked expired.

Written by the indexing model from the issue text.

Assessment

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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.