apache / apache/pulsar

[Bug] Chunk message is incomplete when resending to a different cluster

Open
#24,153 0 comments 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 asking

- [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 don't get bug fixes. I will attempt to reproduce the issue on a supported version of Pulsar client and Pulsar broker.

### Version

- Pulsar client: 4.0.2
- Java: 17.0.x
- Broker: 4.0.3

### Minimal reproduce step

```java
@Test
public void chunkMessage() throws Exception {
String r1="pulsar://localhost:6650";
String r2="pulsar://localhost:6651";
@Cleanup
PulsarClient pulsarClient = PulsarClient.builder()
.serviceUrl(r1)
.build();
String topic = "test" + System.nanoTime();
// Send a message to the r1 cluster.
@Cleanup
Producer producer = pulsarClient.newProducer()
.enableBatching(false)
.enableChunking(true)
.chunkMaxMessageSize(10)
.topic(topic).create();
// When sending the chunk message, update the service URL to r2.
// Assuming the chunk message [0,1,2,3,4,5], when the url to r2, and then the r1 holds [0,1,2], r2 holds [3,4,5].
new Thread(() -> {
try {
System.out.println("sleep...");
Thread.sleep(500);
System.out.println("update...");
// Update the service URL to r2.
pulsarClient.updateServiceUrl(r2);
} catch (PulsarClientException | InterruptedException e) {
throw new RuntimeException(e);
}
}).start();
try {
System.out.println("send...");
MessageId send = producer.send(new byte[4000000]);
System.out.println(send);
System.out.println("send ok...");
} catch (Exception e) {
throw new RuntimeException(e);
}

// Try to receive the message from r2.
@Cleanup
PulsarClient pulsarClient2 = PulsarClient.builder()
.serviceUrl(r2)
.build();
@Cleanup
Consumer consumer = pulsarClient2.newConsumer()
.topic(topic)
.subscriptionName("test")
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscribe();
Message receive = consumer.receive(30, TimeUnit.SECONDS);
assertNotNull(receive); // Got null.
}
```

### What did you expect to see?

I can receive a chunk message from the second cluster.

### What did you see instead?

I cannot receive a chunk message from the second cluster.

### Anything else?

_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 the provided chunkMessage() reproducer, focusing on PulsarClient.updateServiceUrl(), producer chunking, and the send path. Reproduce the service-url change from r1 to r2, then trace how the chunks are distributed; done means the complete message sent through r1 can be received from r2.

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.