Most Pulsar explanations stop at the architecture: stateless brokers, BookKeeper storage, subscriptions and cursors. That is the right place to start, and it is covered in Apache Pulsar architecture. But most production incidents with Pulsar happen in the client library: a producer that buffers until the send timeout fires, a consumer that prefetches a thousand large messages into memory, a retry policy that reorders keyed data, or a delayed message that is delivered immediately because the subscription type ignores delays.

This page is about the client side. It walks through what the Java client does with a message on each side of the broker, the defaults it ships with, and how to choose settings deliberately. Defaults quoted here come from the Java client's configuration classes and the Pulsar messaging documentation; other language clients are similar but not identical, so check yours.

Advertisement

The producer pipeline

When the application calls sendAsync, the message does not go to the network immediately. It joins a pending queue, is grouped into a batch, optionally compressed, framed and written to the broker connection. The broker persists it to BookKeeper and only then returns a receipt with the message ID. The future completes on that receipt, or fails when the send timeout expires.

Inside a Pulsar Java client: where messages waitProducersendAsync()app threadPending queuememory limitedBatch container1 ms / 1000 / 128 KiBCompress + frameNONE by defaultBrokerpersist, then receiptConsumerReceiver queue1000 permitsreceive()app threadAck groupingflushed every 100 mssenddispatch up to permitsackacksflow: more permitsProducer latency = queue wait + batch delay + broker persist. Consumer memory = receiver queue x message size.Every default in this picture is a trade-off you can change per producer or consumer.
Messages wait in three places in a client: the producer's pending queue and batch, and the consumer's receiver queue.
Producer settingJava defaultEffect
enableBatchingtrueGroups messages into one entry; fewer, larger writes.
batchingMaxPublishDelay1 msLongest a message waits for its batch to fill.
batchingMaxMessages1000Batch closes at this many messages.
batchingMaxBytes128 KiBBatch closes at this size.
sendTimeout30 sFuture fails if no receipt arrives in time; 0 disables it.
blockIfQueueFullfalseWhen the queue is full, throw instead of blocking the caller.
compressionTypeNONELZ4, ZLIB, ZSTD and SNAPPY are available; compression applies per batch.

Batching is why Pulsar producers achieve high throughput, and the delay setting is the latency you pay for it. For a low-rate topic, the 1 ms default adds almost nothing. For a high-rate topic, a few milliseconds of delay can multiply batch size and cut broker load substantially. Compression works on whole batches, so it is only worthwhile once batches are large.

The send timeout is a correctness setting, not just a performance one. A timed-out send may still have been persisted; the client simply did not hear back. Retrying blindly creates duplicates unless deduplication is on, which is the next topic.

Routing on partitioned topics

A partitioned topic is a set of ordinary topics, one per partition, and the producer chooses the partition. Messages with a key are hashed to a partition, so every message for one key lands in the same partition and keeps its order. Messages without a key are spread by the routing mode; in the Java client the default for keyless messages is round-robin, and with batching enabled the client keeps sending to one partition for the duration of a batch rather than switching every message, so batches stay full.

Pick keys by the unit whose order matters, such as a customer or an account, and watch for skew: one huge customer makes one hot partition. Partition counts can be increased but not decreased, and increasing them changes which partition a key hashes to for new messages, so per-key order can break across the change. The same trade-offs in Kafka's model are discussed in Kafka consumers.

Advertisement

Idempotent producers and deduplication

With deduplication enabled on the broker or namespace, each producer has a name, and every message carries a sequence ID that increases within that producer for a topic or partition. The broker remembers the highest sequence ID it has persisted per producer and drops anything at or below it. Only one producer with a given name can publish to a topic at a time.

This turns retries into safe operations, but only if the producer's identity survives a restart. Give the producer a stable name, as in the example below, and when it restarts it resumes from the last sequence ID the broker acknowledged. If you supply sequence IDs yourself, derive them from something durable, such as an offset in the source you are copying from, so a replay produces the same IDs. A random producer name per process start defeats the whole mechanism.

PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://pulsar.internal:6650")
    .build();

Producer<Order> producer = client.newProducer(Schema.AVRO(Order.class))
    .topic("persistent://shop/prod/orders")
    .producerName("orders-writer-1")             // stable name: needed for dedup
    .enableBatching(true)
    .batchingMaxPublishDelay(5, TimeUnit.MILLISECONDS)
    .batcherBuilder(BatcherBuilder.KEY_BASED)    // safe for Key_Shared consumers
    .compressionType(CompressionType.LZ4)
    .sendTimeout(30, TimeUnit.SECONDS)
    .blockIfQueueFull(true)                      // backpressure instead of exceptions
    .create();

producer.newMessage()
    .key(order.customerId())                     // same key -> same partition
    .value(order)
    .sendAsync()
    .whenComplete((msgId, err) -> {
        if (err != null) metrics.sendFailed(err); // timeout or rejection: decide, don't drop
    });

Large messages: chunking or a claim check

Brokers enforce a maximum message size. Chunking lets the producer split a large payload into chunks that the consumer reassembles. It has firm constraints: it cannot be combined with batching, so you must disable batching on that producer, and it works only on persistent topics. The consumer must buffer partial chunks until the whole message arrives, which costs memory when several producers interleave large messages.

Often the better answer is a claim check: store the payload in object storage and publish a small message with its location and checksum. The topic stays fast, the payload is fetched only by consumers that need it, and you avoid tuning chunk buffers. Use chunking when the payload must travel with the event and size is occasionally, not routinely, large.

Consumer flow control

Pulsar consumers pull with permits. A consumer tells the broker how many messages it can take, the broker pushes up to that many, and the client stores them in a receiver queue until the application calls receive. When the queue drains to half, the client grants more permits. The Java default receiver queue is 1000 messages per consumer, with a cap of 50,000 across all partitions of a partitioned topic.

Prefetching is good for throughput and dangerous for memory and fairness. A thousand 1 MB messages is a gigabyte in one consumer's heap. On a Shared subscription, a consumer that prefetched 1000 messages holds them even while another consumer sits idle. For slow handlers or large messages, shrink the queue to tens or hundreds; for fast handlers of small messages, the default or larger helps. The Java client also supports a zero-size queue, which fetches one message at a time, at a large throughput cost.

Acknowledgements flow back the other way and are grouped, flushed every 100 ms by default, so a crash can redeliver up to that window of already-processed messages. Handlers must be idempotent regardless.

Choosing a retry mechanism

Pulsar offers four ways to handle a message you could not process, and mixing them up causes both duplicates and stalls.

MechanismDefaultUse it forWatch out for
Ack timeoutOffDetecting handlers that hang without failingFires on slow-but-healthy handlers and duplicates work
Negative ackRedelivered after 1 minuteUnknown, possibly transient failuresCan reorder messages on ordered subscription types
Retry letter topicOff; enableRetry(true)Known slow-to-clear causes with a chosen delayRetries are new messages on <topic>-<subscription>-RETRY
Dead letter topicOff; set maxRedeliverCountPoison messages that never succeedSomeone must read <topic>-<subscription>-DLQ
Consumer<Order> consumer = client.newConsumer(Schema.AVRO(Order.class))
    .topic("persistent://shop/prod/orders")
    .subscriptionName("fulfilment")
    .subscriptionType(SubscriptionType.Shared)
    .receiverQueueSize(100)                      // large messages, slow handler
    .enableRetry(true)                           // retry letter topic
    .deadLetterPolicy(DeadLetterPolicy.builder()
        .maxRedeliverCount(5)
        .build())
    .negativeAckRedeliveryDelay(10, TimeUnit.SECONDS)
    .subscribe();

while (running) {
    Message<Order> msg = consumer.receive();
    try {
        fulfil(msg.getValue());                  // must be idempotent
        consumer.acknowledge(msg);
    } catch (WarehouseBusyException e) {
        consumer.reconsumeLater(msg, 2, TimeUnit.MINUTES); // known, slow-to-clear cause
    } catch (Exception e) {
        consumer.negativeAcknowledge(msg);       // unknown failure: quick redelivery
    }
}

The pattern above uses a negative ack for unknown failures, reconsumeLater with a delay for a known transient cause, and a dead letter policy so a poison message stops after five attempts instead of cycling forever. Alert on the dead letter topic's backlog, and write a small tool to inspect and republish from it, because a dead letter topic nobody reads is just slower data loss.

Delayed delivery

A producer can ask for a message to be delivered later with deliverAfter or at a time with deliverAt. The message is stored immediately; the broker's delayed delivery tracker keeps an index from delivery time to message ID and holds the message back from dispatch until it is due. Delayed delivery is enabled by default on the broker.

Only Shared and Key_Shared subscriptions honour delays. An Exclusive or Failover subscription on the same topic does not delay those messages, so one topic can feed a scheduler on one subscription and an auditor that sees everything immediately on another. The tracker's index lives in broker memory, so millions of messages scheduled days ahead cost broker resources and slow recovery after a broker restart. For long horizons, a scheduler service or a database with a due-time index is usually a better tool, with Pulsar handling the short delays.

Key_Shared rules

Key_Shared lets many consumers share a subscription while keeping per-key order. Its default mode, AUTO_SPLIT, divides the key hash range among consumers automatically; STICKY lets you assign ranges yourself. Batching interacts with it: a batch is dispatched as a unit, so a batch mixing keys would break per-key assignment. The documentation requires producers to either disable batching or use key-based batching, BatcherBuilder.KEY_BASED in Java, as in the producer above.

Consumers joining or leaving move hash ranges, and the broker must avoid delivering a key's newer messages to its new owner while older ones are still unacknowledged at the old one. Slow consumers can therefore stall keys during scale events. Negative acks can also deliver a failed message after later messages for the same key, so if strict per-key order matters, retry inside the handler rather than nacking.

Schemas and transactions

Typed producers and consumers register a schema with the broker, which stores versions per topic and checks each new version against the namespace's compatibility strategy before allowing a producer to connect. That turns a breaking change into a failed producer start rather than a consumer crash at 3 a.m. Look up the strategy your namespace actually enforces with pulsar-admin rather than assuming one, and evolve Avro or Protobuf schemas with optional fields and defaults.

For read-process-write pipelines, transactions make the output messages and the input acknowledgement commit atomically. They require the transaction coordinator on the broker and transactions enabled on the client. They add latency and broker work, and consumers of the output see transactional messages only after commit. The trade-offs against idempotent consumers are covered in exactly-once semantics.

// broker: transactionCoordinatorEnabled=true ; client: .enableTransaction(true)
Transaction txn = client.newTransaction()
    .withTransactionTimeout(1, TimeUnit.MINUTES)
    .build().get();

Message<Order> in = inputConsumer.receive();
Invoice invoice = price(in.getValue());

invoiceProducer.newMessage(txn).value(invoice).sendAsync();
inputConsumer.acknowledgeAsync(in.getMessageId(), txn);
txn.commit().get();   // output visible and input acknowledged together, or neither

Worked example: an order pipeline

An order topic receives 2,000 messages per second at peaks, average 2 KB, keyed by customer, consumed by a fulfilment service whose handler takes about 20 ms because it calls a warehouse API. Each consumer therefore handles about 50 messages per second, so 40 consumers on a Shared subscription are needed at peak, with headroom say 50.

With the default receiver queue of 1000, each consumer could hold 20 seconds of work in memory and, worse, starve its peers at the end of a burst. A queue of 100 holds two seconds of work, which keeps every consumer busy without hoarding. The producer raises its batching delay to 5 ms: at 2,000 per second that is about 10 messages per batch instead of 2. Warehouse throttling responses go to reconsumeLater with a two-minute delay, unknown errors are nacked, and five failures route to the dead letter topic. Because the handler is idempotent on order ID, the occasional redelivery from grouped acks is harmless.

Failure modes

  • Send timeout duplicates: retries after timeouts without deduplication produce copies.
  • Random producer names: deduplication silently stops working across restarts.
  • Prefetch hoarding: large receiver queues cause memory pressure and idle peers on Shared subscriptions.
  • Nack reordering: keyed data processed out of order after a negative ack.
  • Ignored delays: delayed messages read immediately by an Exclusive or Failover subscription.
  • Unread dead letter topics: poison messages accumulate unnoticed.
  • Key_Shared with plain batching: producers that did not switch to key-based batching.

What to do next

  1. List every producer and record its name, batching delay, send timeout and whether deduplication is enabled on its namespace.
  2. Give every producer that retries a stable name and turn on deduplication where duplicates matter.
  3. Size each consumer's receiver queue from handler time and message size, not the default.
  4. Choose one retry mechanism per failure class, and add a dead letter policy with an alert and a replay tool.
  5. Switch Key_Shared producers to key-based batching and avoid nacks where per-key order matters.
  6. Check the schema compatibility strategy on each namespace and test schema changes in CI.
  7. Read about log compaction in log compaction if consumers only need the latest value per key.
Key takeaway: Most Pulsar behaviour that surprises teams lives in the client. Producers batch for 1 ms by default and fail sends after 30 seconds, and a timed-out send may still have been persisted, so give retrying producers stable names and enable deduplication. Consumers prefetch 1000 messages by default, so size receiver queues from message size and handler time. Choose nacks, the retry letter topic and the dead letter topic deliberately, use key-based batching with Key_Shared, remember that delays apply only to Shared and Key_Shared subscriptions, and use transactions only where duplicates are truly unacceptable.