Apache Pulsar Explanation¶
How Apache Pulsar works and why it is built that way. Pulsar has three tiers: stateless brokers, Apache BookKeeper bookies for durable storage, and a metadata store (ZooKeeper today, with Oxia the recommended store from 5.0). The sections below cover the storage model, subscriptions, geo-replication, transactions, the 5.0 scalable-topic design, and the security model. For exact defaults, ports, and version facts, see Reference. For tasks, see How-to Guides.
Component Overview¶
Clients talk only to brokers. Brokers own topics through namespace bundles and write entries to BookKeeper ledgers. Brokers and bookies coordinate through the metadata store. Older ledgers can be offloaded to object storage.
flowchart TB
Client["Producer / Consumer<br/>(pulsar://, pulsar+ssl://)"]
Proxy["Pulsar Proxy<br/>(optional, stateless)"]
subgraph Brokers["Pulsar brokers (stateless)"]
Broker["Broker<br/>ManagedLedger + ManagedCursor<br/>bundle ownership, dispatchers<br/>Transaction Coordinator (opt.)"]
LB["Load manager<br/>(Modular or Extensible)"]
end
subgraph BK["Apache BookKeeper"]
Bookie["Bookies<br/>journal + DbLedgerStorage"]
AutoRecovery["AutoRecovery<br/>(Auditor + ReplicationWorker)"]
end
subgraph Meta["Metadata layer"]
MS["Local metadata store<br/>ZooKeeper or Oxia"]
CS["Configuration store<br/>(optional, shared across clusters)"]
end
subgraph Tiered["Tiered storage"]
Offloader["Ledger offloader"]
Cloud["S3 / GCS / Azure Blob / filesystem"]
end
FW["Functions worker<br/>(in-broker or separate)"]
Client --> Proxy
Client --> Broker
Proxy --> Broker
Broker -->|"addEntry / readEntry"| Bookie
Broker --> MS
LB --> MS
Bookie --> MS
AutoRecovery --> MS
AutoRecovery --> Bookie
Broker --> Offloader
Offloader --> Cloud
Broker -.->|"tenants, namespaces, policies"| CS
FW --> Broker
Layer responsibilities¶
| Layer | Holds | State? |
|---|---|---|
| Broker | Topic ownership (through namespace bundles), subscriptions, ManagedLedger / ManagedCursor handles, entry cache, optional Functions worker and transaction coordinator | Stateless with respect to messages |
| BookKeeper bookies | Persistent message data as ledgers (sequences of entries) | Stateful |
| Metadata store (ZooKeeper / Oxia) | Cluster topology, ledger metadata, broker-bundle ownership, load reports, schemas | Stateful (small) |
| Configuration store | Clusters, tenants, namespaces, policies, partitioned-topic metadata | Stateful (small). By default it is the same store as the local metadata store |
| Tiered storage | Offloaded ledgers | Stateful (large) |
Because brokers hold no message data, a broker failure only means its bundles are reassigned. The new owner reopens the ManagedLedger from BookKeeper and carries on. No data copy or partition rebalancing is needed. This is the core reason for the compute/storage split.
Metadata Store: ZooKeeper, Oxia, and Migration¶
Pulsar reaches its metadata through a pluggable MetadataStore API (introduced by PIP-45). The backends are:
- ZooKeeper — the historical default. It is bundled with Pulsar and remains fully supported in 5.0.
- Oxia — a sharded metadata and coordination store built for Pulsar-scale workloads. It is Apache-2.0 licensed and a CNCF Sandbox project. Pulsar gained an Oxia plugin in 3.3 (PIP-335). Oxia becomes the recommended store for new clusters in 5.0, and it is the required backend for scalable-topic lookups.
- etcd — removed in 5.0 (PIP-462). It saw little production use, and the
jetcdclient is poorly maintained. - RocksDB / memory — standalone and test use only.
ZooKeeper keeps its whole dataset on every ensemble member, so metadata volume and watch count are bounded by a single node. Oxia shards its keyspace across nodes. Very large topic counts and the scalable-topic design depend on that sharding.
5.0 adds a live migration framework (PIP-454). A DualMetadataStore wrapper lets brokers and bookies move from ZooKeeper to Oxia without taking the data plane down.
The migration moves through these phases. It falls back to ZooKeeper automatically if any step fails.
stateDiagram-v2
[*] --> NOT_STARTED
NOT_STARTED --> PREPARATION: pulsar-admin metadata-migration start
PREPARATION --> COPYING: all brokers and bookies recreated ephemeral nodes in Oxia
COPYING --> COMPLETED: persistent nodes copied with versions preserved
PREPARATION --> FAILED: error
COPYING --> FAILED: error
FAILED --> NOT_STARTED: automatic revert to ZooKeeper, retry
COMPLETED --> [*]
During PREPARATION and COPYING, metadata writes are blocked. Publish and consume keep working, but topic and subscription creation and load-balancing moves wait. The docs say these two phases "typically complete in under 30 seconds". After COMPLETED, operators point metadataStoreUrl and the bookies at Oxia. See How-to Guides.
Tenant / Namespace / Topic Hierarchy¶
my-tenant/
├── ns-prod/
│ ├── persistent://my-tenant/ns-prod/orders (partitioned 12 ways)
│ └── persistent://my-tenant/ns-prod/audit-log
└── ns-staging/
└── non-persistent://my-tenant/ns-staging/debug
- Tenant — security and administrative boundary. It lists the clusters it may use (
allowed-clusters). - Namespace — policy boundary: retention, backlog quota, replication clusters, schema rules, dispatch rates, offload policy. It is split into bundles (hash ranges) that the load manager assigns to brokers.
- Topic — the message stream.
persistent://topics are durable.non-persistent://topics are best-effort and in memory only. 5.0 addstopic://for scalable topics.
Topic Storage Model — Ledgers and Entries¶
Each topic partition is a ManagedLedger: an ordered chain of BookKeeper ledgers. Only the last ledger is open for writes. A ledger rolls over by entry count or age (managedLedgerMaxEntriesPerLedger=50000, 10 to 240 minutes by default). Closed ledgers are immutable, so they can be offloaded or deleted as whole units.
flowchart LR
subgraph Topic["ManagedLedger for persistent://my-tenant/ns/orders-partition-0"]
L1["Ledger 101<br/>(closed, offloaded)"]
L2["Ledger 102<br/>(closed)"]
L3["Ledger 103<br/>(open, writable)"]
end
subgraph Bookies["BookKeeper bookies (ensemble per ledger)"]
BkA["Bookie A"]
BkB["Bookie B"]
BkC["Bookie C"]
BkD["Bookie D"]
end
OffStore["Object storage (S3/GCS/Azure)"]
L2 --> BkA
L2 --> BkB
L2 --> BkC
L3 --> BkB
L3 --> BkC
L3 --> BkD
L1 -.->|"offload"| OffStore
Each ledger picks its own ensemble. After a bookie fails, new ledgers simply go to other bookies, and AutoRecovery re-replicates under-replicated fragments in the background. No partition "reassignment" step is needed, unlike Kafka.
Write quorum / Ack quorum / Ensemble size¶
BookKeeper writes use three numbers: E, Qw, and Qa.
| Param | Meaning |
|---|---|
| Ensemble size (E) | Number of bookies that hold this ledger. Entries are striped across them. |
| Write quorum (Qw) | Number of bookies each entry is written to. |
| Ack quorum (Qa) | Number of bookie acks needed before the broker treats the write as durable. |
The shipped broker.conf default is E=2, Qw=2, Qa=2. Production clusters commonly use E=3, Qw=3, Qa=2 for three copies with one slow bookie tolerated on the ack path. Namespaces can override this with set-persistence.
Write path¶
The publish path below shows why durability does not depend on the broker's local disk.
sequenceDiagram
autonumber
participant P as Producer
participant B as Owner broker
participant ML as ManagedLedger
participant BK1 as Bookie 1
participant BK2 as Bookie 2
participant BK3 as Bookie 3
P->>B: lookup topic (HTTP or binary)
B-->>P: owner broker address
P->>B: CommandSend (batch)
B->>ML: asyncAddEntry
par write quorum Qw=3
ML->>BK1: addEntry
ML->>BK2: addEntry
ML->>BK3: addEntry
end
BK1-->>ML: ack (journal fsync)
BK2-->>ML: ack (journal fsync)
Note over ML: Qa=2 acks reached
ML-->>B: entry persisted (ledgerId, entryId)
B-->>P: CommandSendReceipt (MessageId)
B->>B: dispatch to subscriptions from cache
Bookies ack only after the journal write is synced (journalSyncData=true). A lost broker loses nothing that was acknowledged.
Subscription Types¶
A subscription is a named, durable cursor on a topic. Several subscriptions with different types can read the same topic at once. That is how Pulsar covers queueing and pub-sub on the same topic.
flowchart LR
Topic["persistent://t/ns/orders"]
Sub1["S1: Exclusive"]
Sub2["S2: Failover"]
Sub3["S3: Shared"]
Sub4["S4: Key_Shared"]
Topic --> Sub1
Topic --> Sub2
Topic --> Sub3
Topic --> Sub4
C1["Consumer A"]
C2["Consumer B"]
C3["Consumer C"]
Sub1 --> C1
Sub2 --> C1
Sub2 -.->|"standby"| C2
Sub3 --> C1
Sub3 --> C2
Sub3 --> C3
Sub4 -->|"key=k1"| C1
Sub4 -->|"key=k2"| C2
Sub4 -->|"key=k3"| C3
The behavior of each type is listed in Reference. Pulsar 4.0 reworked Key_Shared (PIP-379, "draining hashes"). When consumers join or leave, a hash range moves only after the old consumer has drained its in-flight messages for that range. Earlier versions used a "recently joined consumers" rule that could stall dispatch.
Acknowledgement is per message (individual ack) or cumulative. The broker tracks the mark-delete position plus "holes" of individually acked messages. That is why very sparse acks create large cursor state (PIP-381 addressed large PositionInfo).
Geo-Replication¶
Geo-replication is asynchronous and configured per namespace (or per topic). For each remote cluster, the owner broker creates a replicator. The replicator is an internal cursor plus producer that reads locally persisted messages and republishes them to the same topic on the remote cluster.
flowchart LR
subgraph US["Cluster us-east"]
PUS["Producer P1"]
BrokerUS["Broker (owner)"]
BookieUS["Bookies"]
RepUS["Replicators<br/>to eu-west, ap-southeast"]
end
subgraph EU["Cluster eu-west"]
BrokerEU["Broker (owner)"]
BookieEU["Bookies"]
end
subgraph AP["Cluster ap-southeast"]
BrokerAP["Broker (owner)"]
BookieAP["Bookies"]
end
CS["Configuration store<br/>(optional shared, or per cluster<br/>plus configurationMetadataSyncEventTopic)"]
PUS --> BrokerUS
BrokerUS --> BookieUS
BookieUS --> RepUS
RepUS -->|"async republish"| BrokerEU
RepUS -->|"async republish"| BrokerAP
BrokerEU --> BookieEU
BrokerAP --> BookieAP
BrokerUS -.-> CS
BrokerEU -.-> CS
BrokerAP -.-> CS
Key semantics (from the geo-replication docs):
- Local first. A message is persisted in the local cluster before it is replicated. Producers are never blocked by WAN outages. The backlog is replicated when connectivity returns.
- Full copies. Each cluster keeps its own copy. Subscriptions are local, so consumers in each cluster see every message published in any cluster of the replication set.
- Ordering is guaranteed only per producer. There is no cross-cluster conflict resolution, because messages are appended rather than overwritten. Active-active designs have to accept interleaving from different regions.
- Replicated subscriptions (
enableReplicatedSubscriptions=true, on by default) sync only the mark-delete position, so a consumer can fail over to another region. They are designed for failover, not active-active consumption, and need 2-way replication. - A shared configuration store is not required. Clusters can keep independent configuration stores (the default), share one, or sync configuration through
configurationMetadataSyncEventTopic. - Controls: tenant and namespace
allowed-clustershard-limit where data can go. Namespaceclusterssets the default. Topic-level policies can be local or global. Producers can narrow targets per message (replicationClusters,disableReplication) but cannot bypassallowed-clusters. - Cascading deletes: with shared or synced configuration, removing a cluster from a namespace's
clustersdeletes that namespace's topics on the removed cluster.
Transactions¶
Pulsar transactions (PIP-31) make a consume-process-produce loop atomic across multiple topics and partitions. They are off by default (transactionCoordinatorEnabled=false).
| Component | Role | Persistence |
|---|---|---|
| Transaction coordinator (TC) | Module inside a broker. Allocates TxnIDs, tracks the partitions and subscriptions in each transaction, drives commit or abort, and aborts on timeout | Transaction log (a Pulsar system topic) |
| Transaction buffer (TB) | Per partition. Holds transactional messages as invisible until commit | Messages are written to the real topic. Aborted-transaction state lives in snapshots |
| Pending ack state | Per subscription. Holds acks made inside a transaction, so no other transaction can ack those messages | Pending-ack log (cursor ledger) |
The commit flow below is two-phase: the TC logs the outcome, then tells every participant to materialize or discard.
sequenceDiagram
autonumber
participant C as Client
participant TC as Transaction Coordinator
participant TL as Transaction log
participant BA as Broker (topic out)
participant BI as Broker (topic in)
C->>TC: newTransaction
TC->>TL: append TxnID OPEN
TC-->>C: TxnID
C->>TC: addPartitionToTxn(out-partition-0)
TC->>TL: log partition
C->>BA: send(msg, TxnID) into transaction buffer
C->>TC: addSubscriptionToTxn(in, sub)
TC->>TL: log subscription
C->>BI: ack(msgId, TxnID) into pending ack state
C->>TC: commit(TxnID)
TC->>TL: append COMMITTING
TC->>BA: endTxn COMMIT (messages become visible)
TC->>BI: endTxnOnSubscription COMMIT (acks applied)
TC->>TL: append COMMITTED
TC-->>C: committed
Consumers get read-committed isolation: they never see messages from open or aborted transactions. Transactions give exactly-once processing inside Pulsar. For external sinks you still need idempotent writes or a transactional connector. PIP-439 (4.2) adds transaction support for Pulsar Functions. PIP-473 (5.0) redesigns transactions for scalable topics around the metadata store.
Scalable Topics (Pulsar 5.0)¶
Classic partitioned topics fix a partition count up front. You can increase it but never decrease it, and changing it remaps hash(key) % N, which breaks per-key ordering. Scalable topics (PIP-460, topic://tenant/ns/name) split the key hash space into segments. Each segment is backed by an internal topic. A per-topic controller in the broker splits hot segments and merges cold neighbours. The split and merge history forms a DAG, and keys always move from a parent range to a child range that contains them, so per-key order holds across resizes.
- Scalable topics need the V5 client API (
pulsar-client-v5). Only Java has it in the milestones. Other SDKs are planned before GA. - Segment lookups build on Oxia streaming watch sessions. This is one reason Oxia is the recommended store in 5.0.
- Existing topics can be migrated in place (PIP-475). Classic topics remain fully supported.
- Status: preview in 5.0.0-M1 and M2. Not for production (see Reference).
Pulsar Functions and IO¶
flowchart LR
Input["Input topic(s)"]
Func["Pulsar Function instance<br/>(thread, process, or K8s pod runtime)"]
Output["Output topic"]
State["State storage<br/>(BookKeeper table service)"]
DLQ["Dead-letter topic"]
Input --> Func
Func --> Output
Func -.-> State
Func -.->|"on repeated failure"| DLQ
Functions are lightweight, serverless-style processors (Java, Python, Go). They run in a Functions worker, either embedded in the broker (functionsWorkerEnabled) or as a separate worker cluster (recommended for production isolation). Sources and sinks (Pulsar IO) use the same runtime. In 5.0 the IO connectors leave the core repository (PIP-465) and get their own release cadence.
Schemas¶
The built-in schema registry stores a version history per topic, in the metadata store and BookKeeper. It supports Avro, JSON, Protobuf, primitive types, and KeyValue composites. Compatibility checks run on the broker (schemaCompatibilityStrategy=FULL by default). The registry is not wire-compatible with the Confluent Schema Registry REST API. Kafka-protocol layers such as KoP bundled their own Confluent-compatible endpoint. PIP-420 (4.1) lets Pulsar clients integrate with third-party schema registries.
Protocol Handlers, KoP, and StreamNative Ursa¶
Brokers can load protocol handlers (NAR plugins, messagingProtocols) that speak other wire protocols on top of ManagedLedger and cursors:
- KoP (Kafka on Pulsar) mapped Kafka topics, offsets, and consumer groups onto Pulsar topics and cursors. StreamNative has archived it and points users to its commercial Kafka service.
- MoP (MQTT) and AoP (AMQP 0-9-1) remain StreamNative-maintained plugins. AoP is limited to basic produce and consume.
StreamNative Ursa is a separate, commercial engine, not part of Apache Pulsar. It is a "lakehouse-native", leaderless streaming engine. Brokers write directly to object storage in open table formats (Iceberg and Delta Lake), which avoids cross-AZ replication traffic. It won the VLDB 2025 Best Industry Paper award (StreamNative, PVLDB paper). StreamNative markets Ursa as "lakehouse-native streaming for Kafka and Pulsar". "Ursa For Kafka", a Kafka fork on Ursa storage, entered limited public preview in April 2026. StreamNative has open-sourced the formally verified (TLA+ and Fizzbee) spec of Ursa's leaderless log protocol under Apache-2.0. Stanislav Kozlovski reported in May 2026 that StreamNative committed on 2026-04-10 to open-source Ursa-for-Kafka within nine months (about January 2027). StreamNative's "95% lower cost" headline (product page) is a vendor claim; no independent cost study was found.
Performance Characteristics¶
No reference throughput figure
The rows below describe what bounds each workload, not measured numbers. An earlier per-broker throughput range had no recorded benchmark conditions and was removed (2026-09-28). Measure on your own hardware with OpenMessaging Benchmark before you size anything.
| Workload | Notes |
|---|---|
| Single broker, Qw=3 ledgers, NVMe bookies | Bounded by bookie journal fsync latency, batching and network; varies widely with hardware |
| Cluster-wide aggregate | Scales with broker and bookie counts independently |
| Tiered-storage cold read | Adds object-store latency to the first byte. offloadedReadPriority controls whether bookies or the tier are read first |
| Geo-replication lag | About WAN RTT plus remote broker load under normal conditions |
| Pulsar Functions overhead | Depends on runtime (thread, process, or K8s). No verified figure |
Cost Analysis¶
| Cost | Driver |
|---|---|
| Compute | Brokers plus bookies plus metadata store. Brokers can be small. Bookies need fast journals and memory for caches. |
| Storage | Bookie disks for the hot tier, object storage for the cold tier. |
| Network egress | Geo-replication and tiered-storage uploads. In the cloud, cross-AZ bookie replication is a major line item. That is the problem Ursa-style leaderless designs target. |
| Operations | Three tiers need more on-call expertise than Kafka (KRaft) or NATS. Oxia in 5.0 is meant to simplify the metadata tier. |
| Managed offerings | StreamNative Cloud, DataStax Astra Streaming (IBM), and Tencent Cloud TDMQ take the operational cost off your team. |
Security Model¶
Pulsar splits security into authentication (who you are, mapped to a role), authorization (what that role may do on a tenant, namespace, or topic), and encryption in transit, at rest, and end to end. The setup steps are in How-to Guides. Provider classes and the hardening checklist are in Reference.
Authentication and roles¶
Each connection is authenticated by one of the configured authenticationProviders: JWT, mTLS, OIDC, Athenz, SASL/Kerberos, or Basic. The provider maps the credential to a role string. The Pulsar Proxy forwards the client's role as originalPrincipal, and the broker authorizes the original client, not the proxy.
Authorization¶
The AuthorizationService asks the configured AuthorizationProvider (by default PulsarAuthorizationProvider) whether a role may produce, consume, or manage functions, sources, sinks, or packages. Permissions are granted at namespace or topic level. Super-users (superUserRoles) can do anything. Tenant admins (listed on the tenant) can manage namespaces and grant permissions inside their tenant.
subscriptionAuthMode controls subscription names. With None, any role may create or use any subscription. With Prefix, the subscription name must start with the role name. Prefix isolates cursors between roles that share a topic.
Encryption¶
- In transit: TLS on client-broker, broker-broker (replication), broker-bookie, and broker-metadata-store links. mTLS is supported. Protocols and ciphers are set with
tlsProtocolsandtlsCiphers. - End to end: producers encrypt each payload with a per-message AES data key. That key is encrypted with one or more consumer RSA or ECDSA public keys and carried in the message metadata. Brokers and bookies only see ciphertext, so a compromised broker cannot read payloads. Consumers can register a
decryptFailListener(PIP-436, 4.1). - At rest: Pulsar does not encrypt bookie ledgers itself. Use disk encryption (dm-crypt or cloud volumes). For tiered storage, turn on bucket-default server-side encryption (SSE-S3 or SSE-KMS).
Audit logging¶
Brokers log authentication failures and authorization denials. PIP-402 (4.2) can keep role and originalPrincipal values out of logs. No built-in structured "audit log topic" was found in the 4.x release notes. Earlier versions of this note claimed one. Ship broker logs to a SIEM instead.
Threat Model¶
| Threat | Mitigation |
|---|---|
| Broker compromise reading data | End-to-end encryption for sensitive payloads. |
| Bookie ledger leakage | At-rest disk encryption. Restrict bookie host access. |
| Geo-replication credential theft | Separate replication credentials per cluster, rotated. mTLS for cross-cluster links. Restrict with allowed-clusters. |
| Function code injection | Review Function packages. Restrict the functions/packages permissions. Isolate the Functions worker network (see CVE-2024-27317). |
| Metadata tampering | ZooKeeper ACLs (digest/sasl) and TLS, or Oxia on a private network. No public metadata endpoints. |
| Cross-tenant subscription poisoning | subscriptionAuthMode=Prefix. Tenant-scoped roles. |
| Schema poisoning | Disable schema auto-update. Admins register schemas. |
| Stale-token replay | Short JWT lifetimes, audience pinning (tokenAudience), key rotation. |
| Broker-bookie spoofing | mTLS between brokers and bookies. Bookie authentication through BookKeeper SASL or TLS. |
| Malicious schema in client deserialization | Keep Java clients current (Avro RCE CVE-2024-47561). |
| Tiered-storage bucket misconfiguration | Bucket policies block public access. Cloud audit logging on the bucket. |
| DoS through unbounded backlog or topic listing | Backlog quotas, dispatch and publish rate limits. PIP-442 memory limits for topic-list commands (4.2). |
Comparison Hooks¶
- vs Kafka — Pulsar's compute/storage split and per-ledger ensembles mean no partition reassignment after failures. Kafka's single-tier model is simpler to run. Kafka also has diskless and tiered-storage work that narrows the gap.
- vs Redpanda — Redpanda is a Kafka-API, single-binary engine. Pulsar is multi-tenant and multi-protocol but heavier.
- vs NATS — both have multi-tenant designs. NATS uses accounts, Pulsar uses tenants and namespaces. Pulsar wins on durable-stream scale and NATS on latency and footprint.
- Cross-broker view: Streaming brokers comparison.