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.

Advertisement

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.

Pulsar: stateless serving layer over a separate storage layerProducerssend(key, payload)Consumerssubscription + cursorBroker Aowns bundle 0x00-0x3fBroker Bowns bundle 0x40-0x7fMetadata storeZooKeeper or Oxiabinary protocoldispatchmanaged ledger: ledger 17 (closed), ledger 18 (open)Bookie 1journal + entry logBookie 2journal + entry logBookie 3journal + entry logBookie 4(in ensemble)add entryownership, ledger listTiered storage (object store)closed ledgers offloadedoffloadE=4 bookies in the ensemble, Qw=3 copies per entry, Qa=2 acks before the producer is acked
Brokers own hash ranges of topics (bundles) and write each entry to a subset of bookies. Metadata records ownership and each topic's ledger list. Closed ledgers can be offloaded to object storage.

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.

Advertisement

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

  1. 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.
  2. The broker appends the entry to the open ledger by sending it to Qw bookies chosen from the ensemble.
  3. Each bookie writes the journal, fsyncs, and acknowledges.
  4. 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.
  5. 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.

TypeConsumers servedOrderingUse it for
ExclusiveOne; others are rejectedFull per partitionSingle-reader pipelines
FailoverOne active, others on standbyFull per partitionOrdered processing with a hot spare
SharedMany, round-robin per messageNoneWork queues, delayed delivery
Key_SharedMany, split by key hashPer keyScaling 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 attempts

Delivery 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_hold makes producers wait, producer_exception fails their sends, and consumer_backlog_eviction discards 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/orders

A 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

SymptomLikely causeWhat to do
Producers block or time outBacklog quota reached by one stuck subscriptionFind it in topic stats; fix, delete it or set a TTL
Storage grows without boundAbandoned subscription, or retention set far longer than neededAudit subscriptions per namespace; alert on msgBacklog
Latency spikes during failoverBundle unloads and ledger fencing after broker GC pausesTune GC, raise session timeouts carefully, split hot bundles
Write latency tracks the slowest diskQa equal to Qw, or journal sharing a disk with the entry logQa below Qw; dedicated journal device
Duplicates after producer retriesDeduplication disabledEnable deduplication and set stable producer names
Out-of-order processingShared subscription used for ordered dataKey_Shared or Failover
Under-replicated ledgers after a bookie lossAutorecovery not running or no spare capacityRun 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

  1. Draw your tenant and namespace layout and decide which policies belong to each namespace.
  2. Choose E, Qw and Qa per namespace with Qa below Qw, and set a rack- or zone-aware placement policy.
  3. Put each bookie's journal on its own low-latency device and verify fsync latency under load.
  4. Set retention, TTL and backlog quota explicitly for every namespace, and write down why each value was chosen.
  5. Pick the subscription type per consumer group from the ordering requirement, and make every handler idempotent.
  6. Enable deduplication where producers retry, and configure a dead-letter topic for Shared and Key_Shared subscriptions.
  7. Alert on per-subscription backlog, bundle unload rate and under-replicated ledgers.
  8. Run a failover drill: pause a broker process, confirm producers reconnect and no acknowledged message is lost.
Key takeaway: Pulsar separates serving from storage: stateless brokers own bundles of topics, and each topic is a managed list of BookKeeper ledgers written to Qw of E bookies and acknowledged after Qa. Fencing makes broker failover safe, cursors make every subscription an independent reader, and the subscription type decides ordering. Operate it by setting quorums, retention, TTL and backlog quotas deliberately per namespace, keeping journals on fast disks and watching per-subscription backlog, because an abandoned cursor is still the most common way to fill a cluster.