Apache Pulsar is a distributed log and message queue in one system. Kafka keeps each partition's data on the broker that serves it. Pulsar splits the two jobs: brokers accept and dispatch messages but keep no durable data of their own, and a separate cluster of Apache BookKeeper storage nodes, called bookies, holds the log. Almost everything distinctive about operating Pulsar follows from that split.
This article builds the architecture from first principles, then works through sizing, failure modes and an operations checklist. For the Kafka model it is compared against, see Apache Kafka streaming architecture.
The two layers and why they are separate
A Pulsar cluster has three kinds of process. Brokers speak the client protocol, own topics, cache recent entries and dispatch messages to consumers. Bookies store entries on disk and serve reads of older data. A metadata store records which broker owns which topic, which ledgers make up each topic and where cursors are. Pulsar made the metadata layer pluggable (PIP-45); ZooKeeper remains fully supported, and the current development docs recommend Oxia for new production clusters. For the coordination primitives involved, see ZooKeeper, in depth.
Because a broker holds no data it cannot lose, a topic can move from one broker to another in seconds: the new owner reads the topic's ledger list from metadata and carries on. There is no partition copy to rebuild, and serving and storage scale independently.
Naming, tenancy and bundles
Every topic has a full name of the form persistent://tenant/namespace/topic. Tenants carry authentication and allowed clusters; namespaces carry almost every policy: persistence quorums, retention, TTL, backlog quotas, geo-replication and rate limits.
A partitioned topic is a set of ordinary topics named orders-partition-0, orders-partition-1 and so on, with a client-side router choosing one per message. Ordering exists only within a partition, so a key that must stay ordered has to be routed consistently with partition_key.
Brokers do not own topics one by one. Each namespace's topic-name hash space is divided into bundles, and the load manager assigns whole bundles to brokers. Hot bundles can be split, and bundles are unloaded and reassigned to rebalance; clients reconnect automatically and see a short pause rather than an error.
Storage: ledgers, entries and the managed ledger
BookKeeper's unit of storage is the ledger: an append-only sequence of entries with exactly one writer. Once a ledger is closed it is immutable. A Pulsar topic is therefore not one ledger but a managed ledger, an ordered list of ledgers kept in metadata. The broker writes to the newest open ledger and rolls over on size or time limits or ownership changes. Deleting old data means dropping whole ledgers.
Each ledger has three numbers. The ensemble size E is how many bookies the ledger is spread over. The write quorum Qw is how many of them receive each entry. The ack quorum Qa is how many must confirm before the write counts. They must satisfy E >= Qw >= Qa. With E greater than Qw, entries are striped round-robin across the ensemble, which spreads load and disk use; with E equal to Qw every bookie holds everything.
Inside a bookie, an entry is appended to the journal, a write-ahead log that is fsynced before the bookie acknowledges, and also placed in a write cache that is later flushed to entry log files with an index (RocksDB in the default DbLedgerStorage). The journal is written sequentially and read only during recovery, so putting it on its own fast disk keeps write latency independent of read traffic: a consumer replaying a week of history does not slow down producers.
The life of a write
- The producer sends a batch to the broker that owns the topic's bundle. If the client talks to the wrong broker, a lookup redirects it.
- The broker appends the entry to the open ledger by sending it to Qw bookies chosen from the ensemble.
- Each bookie writes the journal, fsyncs, and acknowledges.
- When Qa acknowledgements arrive, the broker advances the ledger's last-add-confirmed position, acknowledges the producer with a message ID of the form (ledger id, entry id, partition, batch index), and makes the entry readable to consumers.
- Consumers that are caught up are served from the broker's cache without touching bookies; lagging ones read from bookies.
With Qa smaller than Qw, one slow bookie does not stall writes: the broker is acknowledged by the fastest two of three. The third copy still arrives, or autorecovery rebuilds it. Qa equal to Qw makes tail latency track your slowest disk.
import pulsar
client = pulsar.Client("pulsar://broker.internal:6650")
producer = client.create_producer(
"persistent://payments/prod/orders", # tenant/namespace/topic
send_timeout_millis=30000,
block_if_queue_full=True, # apply backpressure instead of raising
)
def publish(order):
# partition_key routes every event of one order to the same partition,
# and Key_Shared consumers later keep per-key order on that key.
producer.send(order.to_json().encode(), partition_key=order.order_id)block_if_queue_full turns a full client queue into backpressure. Deduplication, enabled per namespace or broker, drops retried sends already stored, keyed by producer name and sequence ID; without it a timeout plus retry can store a message twice.
Failover without split brain: fencing
Suppose broker A owns the orders topic, stalls in a long garbage-collection pause, and loses its metadata session. Broker B takes ownership. A is not dead; when it wakes up it will keep trying to append to the ledger it had open. Two writers on one ledger would corrupt the log.
BookKeeper prevents this with fencing. B opens A's last ledger in recovery mode. That operation tells the bookies of the ledger's ensemble to fence it, after which they reject any further add from the old writer. B then reads the tail to find the last entry that reached Qa, closes the ledger at that point, and starts a new ledger for its own writes. When A's delayed writes arrive, they fail, A learns it lost ownership, and its producers reconnect to B. Anything A acknowledged had reached Qa bookies and survives; anything unacknowledged is retried by the producer.
Subscriptions, cursors and acknowledgements
Consumers never delete messages. Each subscription has a durable cursor that records the mark-delete position, below which everything is acknowledged, plus ranges of individually acknowledged messages above it. Data can be removed only once every subscription's cursor has passed it and the retention policy no longer requires it.
| Type | Consumers served | Ordering | Use it for |
|---|---|---|---|
| Exclusive | One; others are rejected | Full per partition | Single-reader pipelines |
| Failover | One active, others on standby | Full per partition | Ordered processing with a hot spare |
| Shared | Many, round-robin per message | None | Work queues, delayed delivery |
| Key_Shared | Many, split by key hash | Per key | Scaling ordered, keyed workloads |
Shared subscriptions are what make Pulsar a queue: individual acknowledgements, negative acknowledgements and a dead-letter topic make it natural to retry one poison message without blocking the rest. Cumulative acknowledgement, which acknowledges everything up to a message, is only available on Exclusive and Failover subscriptions. For how consumer groups rebalance in the Kafka model, compare consumer rebalancing.
import pulsar
client = pulsar.Client("pulsar://broker.internal:6650")
consumer = client.subscribe(
"persistent://payments/prod/orders",
subscription_name="billing",
consumer_type=pulsar.ConsumerType.KeyShared, # scale out, keep per-key order
negative_ack_redelivery_delay_ms=10000,
dead_letter_policy=pulsar.ConsumerDeadLetterPolicy(
max_redeliver_count=5,
dead_letter_topic="persistent://payments/prod/orders-billing-DLQ",
),
)
while True:
msg = consumer.receive()
try:
bill(msg.data()) # must be idempotent: delivery is at-least-once
consumer.acknowledge(msg) # moves this subscription's cursor
except Exception:
consumer.negative_acknowledge(msg) # redeliver; DLQ after 5 attemptsDelivery is at least once: a crash after processing but before the acknowledgement arrives redelivers the message, so the handler must be idempotent. Pulsar also has transactions that atomically acknowledge input and publish output, the same pattern described in exactly-once stream processing; use them when duplicates are genuinely unacceptable, and measure the latency they add.
Retention, TTL and backlog quota are three different things
This is the most misread part of Pulsar, and it decides both your storage bill and whether producers stall.
- Retention keeps data that every subscription has already acknowledged, for a time or size limit, so new readers can replay it. Without retention, acknowledged data is eligible for deletion immediately.
- Message TTL acts on unacknowledged data: messages older than the TTL are automatically acknowledged for every subscription, so a dead consumer cannot pin storage forever. It silently skips data, which is the point, and also the risk.
- Backlog quota caps how much unacknowledged data a topic may accumulate, and chooses what happens at the limit:
producer_request_holdmakes producers wait,producer_exceptionfails their sends, andconsumer_backlog_evictiondiscards the oldest backlog.
# Durability for the namespace: stripe over 3 bookies, 3 copies, ack after 2
pulsar-admin namespaces set-persistence \
--bookkeeper-ensemble 3 --bookkeeper-write-quorum 3 --bookkeeper-ack-quorum 2 \
--ml-mark-delete-max-rate 0 payments/prod
# Keep acknowledged data for 3 days or 500G, whichever is hit first
pulsar-admin namespaces set-retention --size 500G --time 3d payments/prod
# Cap unacknowledged backlog per topic; hold producers when exceeded
pulsar-admin namespaces set-backlog-quota \
--limit 200G --policy producer_request_hold payments/prod
# Auto-acknowledge anything older than 7 days that nobody consumed
pulsar-admin namespaces set-message-ttl --messageTTL 604800 payments/prod
# Where is the backlog?
pulsar-admin topics stats persistent://payments/prod/ordersA forgotten subscription is the classic incident: someone created it for a test, never consumed, and it now holds every ledger since. With no TTL and a hold policy, producers eventually block. topics stats shows each subscription's msgBacklog; alert on it per subscription, not only per topic. The general pattern of letting a slow consumer push back on a fast producer is covered in backpressure in streaming systems.
Tiered storage and geo-replication
Because a topic is a list of immutable ledgers, closed ledgers can be copied to object storage and deleted from bookies while the ledger list simply points at the new location. Reads of old data become slower and cheaper; the bookie tier only has to hold the recent working set. This is the same economic trade-off as tiered storage in Kafka, arrived at more naturally because segment boundaries already exist.
Geo-replication is configured per namespace. Each cluster's broker runs a replicator, effectively an internal subscription, that reads local topics and publishes them to remote clusters asynchronously. Replication lag is a backlog like any other, and a message acknowledged in one region is lost if that region fails before it ships. For synchronous durability, stretch one BookKeeper cluster across regions instead and pay the round trip on every write.
Worked example: sizing a payments namespace
A payments team expects 50 MB/s of compressed producer traffic, wants three copies, and needs three days of replay. They choose E=3, Qw=3, Qa=2 across three availability zones with a rack-aware placement policy so that each copy lands in a different zone.
ingress_mb_s = 50 # producer bytes per second, after compression
qw, qa = 3, 2 # write quorum, ack quorum
retention_days = 3
bookie_net_in = ingress_mb_s * qw # 150 MB/s across the bookie tier
bookie_disk = bookie_net_in * 2 # journal + entry log: ~300 MB/s
stored_tb = ingress_mb_s * 86400 * retention_days * qw / 1e6 # ~38.9 TB replicated
print(bookie_net_in, bookie_disk, round(stored_tb, 1))The bookie tier receives 150 MB/s of network traffic and writes roughly 300 MB/s to disk, because every byte goes to the journal and then the entry log. Three days of retention holds about 39 TB of replicated data. The team spreads that over six bookies, not three, so that autorecovery has somewhere to re-replicate after a loss, and plans to raise E to five later to stripe each ledger over more machines. Offloading ledgers older than twelve hours to object storage would shrink the bookie footprint to about a sixth of that, at the cost of slower replays.
Failure modes
| Symptom | Likely cause | What to do |
|---|---|---|
| Producers block or time out | Backlog quota reached by one stuck subscription | Find it in topic stats; fix, delete it or set a TTL |
| Storage grows without bound | Abandoned subscription, or retention set far longer than needed | Audit subscriptions per namespace; alert on msgBacklog |
| Latency spikes during failover | Bundle unloads and ledger fencing after broker GC pauses | Tune GC, raise session timeouts carefully, split hot bundles |
| Write latency tracks the slowest disk | Qa equal to Qw, or journal sharing a disk with the entry log | Qa below Qw; dedicated journal device |
| Duplicates after producer retries | Deduplication disabled | Enable deduplication and set stable producer names |
| Out-of-order processing | Shared subscription used for ordered data | Key_Shared or Failover |
| Under-replicated ledgers after a bookie loss | Autorecovery not running or no spare capacity | Run autorecovery; keep headroom beyond Qw bookies |
Trade-offs against a broker-stores-data design
Pulsar's separation gives fast rebalancing, independent scaling of compute and storage, very large topic counts and natural tiered storage. It costs an extra hop and three services to monitor instead of one. For a small team with a few high-throughput topics, a single-layer log is usually simpler; for a multi-tenant platform mixing queues and streams, Pulsar's model pays for itself.
What to do next
- Draw your tenant and namespace layout and decide which policies belong to each namespace.
- Choose E, Qw and Qa per namespace with Qa below Qw, and set a rack- or zone-aware placement policy.
- Put each bookie's journal on its own low-latency device and verify fsync latency under load.
- Set retention, TTL and backlog quota explicitly for every namespace, and write down why each value was chosen.
- Pick the subscription type per consumer group from the ordering requirement, and make every handler idempotent.
- Enable deduplication where producers retry, and configure a dead-letter topic for Shared and Key_Shared subscriptions.
- Alert on per-subscription backlog, bundle unload rate and under-replicated ledgers.
- Run a failover drill: pause a broker process, confirm producers reconnect and no acknowledged message is lost.