apache / apache/pulsar

Exclusive subscription mode for partitioned topics not seeing consistent order among different consumers

Open
#11,883 6 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**

In the official documentation of Apache Pulsar it says that an exclusive subscription sees a consistent order for a single consumer. Suppose we have multiple consumers each of them with its own exclusive subscription and they are reading from a partitioned topic. As presented in the docs:

"Decisions about routing and subscription modes **can be made separately in most cases**. In general, throughput concerns should guide partitioning/routing decisions while subscription decisions should be guided by application semantics.

**There is no difference between partitioned topics and normal topics in terms of how subscription modes work**, as partitioning only determines what happens between when a message is published by a producer and processed and acknowledged by a consumer."

Those statements lead us to infer that the readers would get a consistent global ordering among topic partitions. But that is not the case in my tests so far:

The code I have used to test it (Scala):

**Producer:**

package pulsar

import org.apache.pulsar.client.api.PulsarClient
import java.util.UUID

object Producer {

def main(args: Array[String]): Unit = {

val client = PulsarClient.builder()
.serviceUrl(SERVICE_URL)
.allowTlsInsecureConnection(true)
.build()

val producer = client.newProducer()
.topic(TOPIC)
.enableBatching(true)
//.accessMode(ProducerAccessMode.Exclusive)
.create()

for(i<-0 until 100){
val key = UUID.randomUUID().toString.getBytes()
//val key = s"Hello-${i}".getBytes()
producer.newMessage().orderingKey("k0".getBytes()).value(key).send()

println(s"produced msg: ${i.toString}")
}

producer.flush()

producer.close()
client.close()
}

}

**Consumer:**

```
import org.apache.pulsar.client.api.{PulsarClient, SubscriptionInitialPosition, SubscriptionType}

object Consumer {

def main(args: Array[String]): Unit = {

val client = PulsarClient.builder()
.serviceUrl(SERVICE_URL)
.allowTlsInsecureConnection(true)
.build()

var l1 = Seq.empty[String]
var l2 = Seq.empty[String]

val c1 = client.newConsumer()
.topic(TOPIC)
.subscriptionType(SubscriptionType.Exclusive)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscriptionName("c1")
.subscribe()

// To stop the consuption I put a limit (100) - this limit is known
while(l1.length < 100){
val msg = c1.receive()
val str = new String(msg.getData)

println(s"${Console.MAGENTA_B}$str${Console.RESET}")

l1 = l1 :+ str

//c1.acknowledge(msg.getMessageId)
}

val c2 = client.newConsumer()
.topic(TOPIC)
.subscriptionType(SubscriptionType.Exclusive)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscriptionName("c2")
.subscribe()

while(l2.length < 100){
val msg = c2.receive()
val str = new String(msg.getData)

println(s"${Console.GREEN_B}$str${Console.RESET}")

l2 = l2 :+ str

//c2.acknowledge(msg.getMessageId)
}

println()
println(l1)
println()

println()
println(l2)
println()

try {
assert(l1 == l2)
} finally {
c1.close()
c2.close()
client.close()
}
}

}
```

Am I wrong about it or Pulsar does not support the described behavior I expect?

**OBS.: I've tried every configuration for that. I've set Retention policies for the namespace as infinite both in size and number of messages. It does not work :(
I also tried:

- SinglePartitionRouting for producer (it does not matter tho)
- Setting an ordering key for the messages**

Contributor guide

Open the contributing guide

Research direction

Start with the official documentation on exclusive subscriptions, partitioned topics, routing, and ordering, then run the Scala Producer and Consumer reproduction described in the issue. Compare the observed consumer sequences with the documented guarantees and determine whether the behavior is a bug or whether the documentation must clarify the ordering semantics; done means the supported behavior is established and consistently documented or tested.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, scala
Domain
distributed-systems
Issue type
Bug
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.