Scaling Message Throughput with Apache Pulsar Shared Subscriptions
Learn how to implement Apache Pulsar Shared Subscriptions to horizontally scale message consumption across multiple workers and improve system throughput.
18 Oct 2025, 23:19 UTC

The Bottleneck of Single-Consumer Processing
When a single consumer cannot keep up with the volume of messages arriving in a Pulsar topic, the resulting backlog increases latency and risks filling broker storage. While Exclusive and Failover subscription modes ensure strict ordering, they limit processing to one active consumer at a time, creating a performance ceiling.
The solution for high-throughput worker pools is the Shared Subscription. This mode allows multiple consumers to attach to the same subscription, distributing messages in a round-robin fashion. The primary takeaway is that Shared subscriptions trade message ordering for horizontal scalability.
Prerequisites and Constraints
- Pulsar Cluster: An active Apache Pulsar cluster (v2.0+).
- Client Library: The Pulsar Java, Python, or Go client installed on worker nodes.
- Ordering Requirement: Your application must be able to process messages independently. Because messages are distributed across different consumers, the original sequence of production is not preserved during consumption.
Implementing a Shared Consumer Pool
To implement this, you must explicitly set the subscription type to Shared during the consumer creation process. If you do not specify a type, Pulsar defaults to Exclusive, which will cause subsequent consumers attempting to join the same subscription to fail.
Below is a configuration example using the Java Client:
import org.apache.pulsar.client.api.*;
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
Consumer<byte[]> consumer = client.newConsumer()
.topic("persistent://public/default/my-worker-topic")
.subscriptionName("worker-pool-subscription")
.subscriptionType(SubscriptionType.Shared) // Critical for load balancing
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscribe();
Execution Context: Run this code on multiple independent worker instances or pods. Ensure all instances use the exact same subscriptionName ("worker-pool-subscription") to be treated as a single logical group by the broker.
Operational Behavior and Diagnostics
In Shared mode, the broker tracks acknowledgments (ACKs) for every individual message. This differs from Exclusive mode, where the broker typically tracks a cursor position.
| Scenario | Broker Action | Expected Result |
|---|---|---|
| New Consumer Joins | Adds consumer to the round-robin rotation. | Immediate distribution of new messages to the new node. |
| Consumer Crash | Detects connection loss; identifies unacknowledged messages. | Un-ACKed messages are redelivered to other active consumers. |
| Message Processing Failure | Consumer calls negativeAcknowledge(). |
Message is rescheduled for redelivery after a timeout. |
Verification Steps
- Deploy: Start three separate consumer instances using the configuration above.
- Produce: Send a burst of 300 messages to the topic.
- Check Distribution: Monitor the logs of each consumer. You should see roughly 100 messages processed per instance, confirming the round-robin distribution.
- Test Fault Tolerance: Terminate one consumer instance while messages are still being processed. Verify that the messages assigned to the killed instance are eventually picked up and processed by the remaining two consumers.
Limitations and Risk Mitigation
Memory Pressure: Because the broker must track individual ACKs for every message in flight across all shared consumers, a high volume of unacknowledged messages can increase broker memory usage. To mitigate this, keep the receiverQueueSize (the number of messages the client pre-fetches) at a reasonable level relative to your processing speed.
Idempotency: In the event of a consumer crash, a message might be partially processed before being redelivered to another consumer. To prevent duplicate side effects, implement idempotent processing (e.g., using a database unique constraint on a message ID) to ensure that processing the same message twice does not corrupt your data.
Rollback and State Change
Changing a subscription type from Shared back to Exclusive or Failover requires the removal of all existing consumers. If you need to revert the subscription type:
- Stop all consumer instances.
- Use the Pulsar Admin CLI to update the subscription or delete the subscription entirely to reset the state:
bin/pulsar-admin subscriptions delete my-worker-topic --subscription worker-pool-subscription - Restart consumers with the desired
SubscriptionType.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.