Apache Pulsar Geo-Replication for Industrial IoT Telemetry

Apache Pulsar Geo-Replication for Industrial IoT Telemetry

Apache Pulsar Geo-Replication for Industrial IoT Telemetry

Last Updated: September 2026

A factory does not care that your message bus is elegant. It cares that a vibration alarm in Monterrey is visible to the reliability team in Stuttgart within seconds, that the same raw waveform never leaves Mexico if the plant’s data policy says so, and that a submarine cable fault does not stop the line historian from writing. Apache Pulsar geo-replication is one of the few open-source mechanisms that lets you express all three requirements in configuration rather than in glue code. But it is also easy to misuse: teams replicate everything everywhere, discover their WAN bill and replication backlog in the same week, and conclude the technology is at fault.

This rewrite treats geo-replication as a budgeting problem first and a topology problem second. You will see how replication actually works underneath, how to split telemetry into namespaces with different replication scope, how to size the WAN and the backlog with worked arithmetic, and what changed in the Pulsar ecosystem through September 2026.

What this covers: what changed since the original post, the mechanics of asynchronous replication, a reference architecture for a five-region manufacturer, namespace design for sovereignty, replicated subscriptions and client failover, edge pre-processing, cost arithmetic, comparison with Kafka-family alternatives, and the operational mistakes that hurt in production.

What changed for September 2026

This article replaces an April 2026 version. Several claims in that version were wrong or have aged, and the ecosystem has moved. Here is the short list, with the detail in the sections that follow.

  • Versions. Apache Pulsar 4.0 is the current long-term support (LTS) line, released in October 2024; per the project’s support table its active support window ends on 21 October 2026, with security-only support afterwards. The 4.2 feature line shipped in April 2026, and a 5.0 milestone (5.0.0-M2, explicitly not for production) appeared in September 2026. If you are planning a new deployment this quarter, plan for 4.x now and a deliberate 5.x evaluation later.
  • Oxia is not a “3.x default”. The old post implied that Oxia had replaced ZooKeeper for new clusters. The current migration documentation describes an online ZooKeeper-to-Oxia migration framework that requires Pulsar 5.0 or later on all brokers, bookies, and auto-recovery daemons. In 4.x, ZooKeeper (or another supported store) remains the safe production choice, and 4.2 removed the Etcd backend.
  • Geo-replication hardening in 4.2. The 4.2.0 release notes include PIP-433 (ensure the topic exists before geo-replication starts), stricter validation when enabling namespace-level replication, and fixes to replicator behavior around snapshots and disconnections. These are exactly the edge cases that bite IIoT deployments on flaky links.
  • Kafka-family alternatives now have replication features too. The old post said Redpanda had no native replication story. Redpanda now documents an enterprise Shadowing feature for asynchronous cross-region disaster recovery, and StreamNative positions its Ursa engine and Universal Linking as a lakehouse-native, Kafka-compatible option. The comparison later in this article is rewritten accordingly.
  • Corrections to mechanics. Ledger rollover defaults, bookie quorum arithmetic, the claimed behavior of Pulsar Functions, and several regulatory statements were overstated or wrong. Each is fixed in place and called out where it matters.

If you only read one section for the delta, read this one; if you are designing a system, read on.

Context and Background

Industrial telemetry has three properties that ordinary web event streams do not share. It is continuous and machine-paced: a bearing sensor does not go quiet at night, so capacity planning uses sustained rates, not diurnal averages. It is jurisdictionally entangled: a single plant can produce data that is operationally shared, contractually restricted, and legally localized, sometimes on the same topic hierarchy. And it is loss-intolerant at the historian, latency-tolerant at the analytics tier: control loops stay on the plant network, but quality and maintenance analytics can accept seconds of delay in exchange for completeness.

Those properties explain why a single global broker cluster is rarely the right answer. Cross-ocean round trips add tens to hundreds of milliseconds to every synchronous acknowledgment, a single regional outage would take down every plant, and policy boundaries cannot be enforced inside one flat namespace. The pragmatic pattern is one cluster per region or per large site, with asynchronous replication of selected data between them. That is precisely the model Apache Pulsar’s geo-replication implements.

Pulsar separates serving from storage. Stateless brokers handle producers, consumers, and dispatch. Apache BookKeeper stores each topic as a sequence of append-only ledgers. A metadata store coordinates ownership and configuration. Geo-replication sits on top: brokers in one cluster run internal replicators that read from local topics and publish to the same topic name in remote clusters. The official documentation describes the semantic bluntly: messages are first persisted to the local cluster and then replicated asynchronously by brokers, with lag governed by network round-trip time and the health of the link.

For readers coming from the MQTT side of the plant, the ingest edge is covered in our guides to MQTT-to-Kafka bridging with Kafka Connect and the MQTT and Sparkplug B reference architecture; the payload conventions there carry straight over to Pulsar topics. If you want the generic, non-industrial treatment of Pulsar replication internals, see the sibling article Apache Pulsar geo-replication architecture in 2026. This article deliberately stays on the telemetry problem: rates, payload shapes, sovereignty tiers, and plant-to-hub economics. The authoritative reference for everything below is the Apache Pulsar geo-replication administration guide.

A note on scope and honesty. Numbers in the worked examples are illustrative arithmetic on stated assumptions, not benchmarks. Where Pulsar defaults are quoted, they come from the project’s broker.conf on the main branch at the time of writing; always check the file for your exact version.

Reference Architecture: Region Clusters, a Hub, and Tiered Replication Scope

Apache Pulsar geo-replication copies each message asynchronously from the cluster where it was produced to the clusters listed in the namespace or topic replication policy, using broker-resident replicators that read locally and publish remotely. Producers acknowledge after local durability only, so the WAN never sits in the write path. The design consequence for IIoT is that you choose, per namespace, how far each stream travels.

Apache Pulsar geo-replication reference architecture for industrial IoT telemetry

Figure 1: Plant gateways feed a regional Pulsar cluster; Pulsar Functions reduce raw streams to summaries, and only summaries and events cross the WAN to the hub and peers.

Figure 1 shows the pattern this article defends. Each region runs its own cluster with its own metadata store. Plant gateways translate OPC UA and MQTT into Pulsar topics inside the region. Functions running next to the regional cluster reduce raw streams into summaries. A per-namespace policy then decides whether a stream stays local, goes to the hub, or goes everywhere. The lakehouse consumes from the hub. Cold ledgers offload to object storage in the same jurisdiction as the cluster that wrote them.

How a replicated message actually moves

When a namespace lists more than one cluster in its replication set, each broker that owns a topic starts one replicator per remote cluster. A replicator is a durable cursor on the local topic plus a producer connected to the remote cluster. It reads entries the cursor has not yet acknowledged, publishes them to the remote topic of the same name, and advances the cursor when the remote broker acknowledges. Because the cursor is durable, a broker restart or a link outage simply leaves a backlog; when the link returns, the replicator resumes from where it stopped.

Three properties fall out of this and matter for telemetry. First, ordering is preserved per topic partition per source cluster, because a single replicator drains a cursor in order. Second, there is no global order across sources: if two regions both produce to a replicated topic, each cluster ends up with an interleaving of the two streams that can differ from cluster to cluster. For sensor data keyed by asset this is harmless; for anything that needs a total order it is a design error. Third, a message carries the identity of its origin cluster, which is how replicators avoid echoing a message back to where it came from in a mesh.

The ledger detail in the previous version of this post was partly wrong, so here is the correct picture. A Pulsar topic (or each partition of a partitioned topic) is stored as a chain of ledgers. A ledger rolls over based on count, size, and time thresholds, not the “50 MB or 5 minutes” the old text claimed. In the current broker.conf, managedLedgerMaxEntriesPerLedger defaults to 50,000 entries, managedLedgerMinLedgerRolloverTimeMinutes to 10, and managedLedgerMaxLedgerRolloverTimeMinutes to 240. Replication does not ship sealed ledger files across regions; it reads entries through the managed-ledger cursor like any other consumer and republishes them. That distinction matters: replication cost scales with message count and bytes, not with ledger boundaries.

The storage tier: bookies, quorums, and what “2 of 3” really means

BookKeeper writes each entry to an ensemble of bookies according to three numbers: ensemble size, write quorum, and ack quorum. The broker acknowledges a produce once ack-quorum bookies have persisted the entry. The stock defaults in broker.conf are managedLedgerDefaultEnsembleSize=2, managedLedgerDefaultWriteQuorum=2, and managedLedgerDefaultAckQuorum=2, which tolerate no bookie loss on the write path. The commonly recommended production shape is ensemble 3, write quorum 3, ack quorum 2, which lets one bookie be slow or down without stalling producers while still storing three copies.

The old post claimed that with a write quorum of 2, “2 bookies can be down without data loss.” That is wrong. Ack quorum 2 of write quorum 3 tolerates one unavailable bookie per ensemble for writes; losing two at once removes the ability to make progress until the ensemble can change. Size regional clusters so that an ensemble change to a fresh bookie is always possible: five bookies is a sensible minimum for a region that must survive one failure during a rolling upgrade.

On disks, the old text described the ledger disk as holding a “columnar index.” It does not. BookKeeper bookies write to a journal (a sequential write-ahead log that determines write latency, best on fast NVMe or a dedicated device) and flush to ledger storage with an entry-location index, which serves reads. Keeping the journal on a separate device from ledger storage is still the single most useful bookie tuning step, but the reason is the isolation of sequential synchronous writes from read and compaction traffic, not a columnar layout.

Configuring the pieces

Cluster registration, tenant permissions, and namespace policy are three separate steps, and confusing them is the most common way replication silently fails to start. The following commands follow the official guide and assume three clusters named eu-fra, na-hub, and apac-pune.

# 1. Register every peer cluster in each cluster's metadata (run on each side)
bin/pulsar-admin clusters create \
  --broker-url pulsar+ssl://broker.na-hub.example.com:6651 \
  --url https://admin.na-hub.example.com:8443 \
  na-hub

# 2. Create a tenant that is allowed to exist on these clusters
bin/pulsar-admin tenants create plant-eu \
  --admin-roles plant-eu-admin \
  --allowed-clusters eu-fra,na-hub

# 3. Create the namespace and set its replication set
bin/pulsar-admin namespaces create plant-eu/summary
bin/pulsar-admin namespaces set-clusters plant-eu/summary \
  --clusters eu-fra,na-hub

The important field is --allowed-clusters on the tenant. It is the guardrail: a producer can narrow the replication set per message, but it cannot escape the tenant’s allowed list. The documentation is explicit that clusters at the namespace level is the default target set, while allowed-clusters is what is permitted at all. For sovereignty, that difference is the enforcement mechanism, and we return to it in the namespace section.

You also need a decision about configuration metadata. By default each cluster keeps independent configuration, so tenants and namespaces must be created on every participating cluster. Alternatives are a shared configuration store across clusters or the metadata-sync event topic (configurationMetadataSyncEventTopic) that the docs describe for propagating configuration. Shared stores are convenient but they couple regions’ control planes; for jurisdictions that must be able to operate on their own, independent configuration managed by infrastructure-as-code is the more defensible choice. Topic policies, note, are replicated via geo-replication whenever the namespace has replication enabled, whichever configuration approach you pick.

Namespace Design: Four Tiers of Telemetry

The core design move in this article is to stop thinking of “the telemetry topic” and to define four tiers with different replication scope, retention, and failure behavior. Each tier is a namespace, and each namespace inherits its policy at creation time so that nobody has to remember it later.

Namespace tiers mapping telemetry classes to Pulsar geo-replication scope

Figure 2: A tenant per region with allowed-clusters set, and four namespaces whose replication scope shrinks as data gets closer to the control loop.

Figure 2 lays out the tiers. They are worth defining precisely because every later decision follows from them.

Tier 1: raw telemetry, local only

Raw is everything the gateways produce at native rate: vibration waveforms, high-frequency current samples, camera metadata, PLC scan snapshots. It has the highest volume, the least cross-site value, and the strongest sovereignty sensitivity. The namespace’s replication set is the local cluster only, retention is measured in days, and older ledgers offload to in-region object storage. Nothing about the plant’s raw data crosses a border unless someone deliberately changes the tenant’s allowed-clusters, which is a reviewable change.

Tier 2: summaries, region to hub

Summaries are the product of edge reduction: one-second RMS values, min-max-mean windows, state-change events per asset, counters. They are one to two orders of magnitude smaller than raw, and they are what global reliability and quality teams actually query. This namespace replicates one-way from each region to the hub. Because summaries are derived data, you can scrub or pseudonymize operator identifiers and exact facility coordinates in the reduction step, before anything is eligible to leave the region. That is a better privacy control than trying to filter at the replication layer.

Tier 3: operational events, full mesh

Alarms, equipment state changes, shift handovers, and maintenance-start events are low in volume and high in value to every region. Replicating them in a full mesh costs little and enables follow-the-sun operations: the night shift in Pune sees what the Stuttgart day shift changed. Because volume is small, this is also the only tier where replicated subscriptions and active-active consumers are worth their complexity.

Tier 4: control and setpoints, never replicated

Setpoints, interlocks, and command acknowledgments belong in the plant’s own control system and network segment. If you use Pulsar for command distribution at all, keep it in a namespace with a replication set of exactly one cluster, no summary subscription, and network policy that mirrors your IEC 62443 zone design. The article on IEC 62443 zones and conduits covers why bidirectional bridges between zones deserve extra scrutiny. Nothing here is a substitute for a safety-rated control path.

Regulation: what the law does and does not say

The old post said GDPR “mandates” that personal data stay in-region, that China’s law requires sensitive industrial data to stay in China, and that India’s DPDP Act “imposes similar localization rules.” Those statements overreach. GDPR does not require EU-only storage; it restricts transfers of personal data outside the EU unless a transfer mechanism such as an adequacy decision or standard contractual clauses applies. China’s data laws do impose localization and export-assessment requirements for certain categories of important and personal data. India’s Digital Personal Data Protection Act permits cross-border transfer except to jurisdictions the government restricts. Industrial sensor readings are often not personal data at all, but operator badges, shift rosters, and camera-derived metadata can be.

The engineering conclusion does not depend on the legal detail: make the data class explicit, keep the highest-risk class local by construction, and let counsel decide what may travel. Pulsar’s allowed-clusters and per-namespace scope are how you make “by construction” auditable. This is engineering guidance, not legal advice.

Sizing the WAN and the Backlog: Worked Arithmetic

Geo-replication fails in production for arithmetic reasons far more often than for software reasons. The rule is simple: the sustained replication rate to each destination must stay below the usable bandwidth of the path to that destination, with headroom for catch-up after an outage. Everything else in this section is a way of making that inequality checkable before go-live.

Sequence of a message from gateway to hub showing local ack and asynchronous replicator

Figure 3: The producer is acknowledged after the local bookie quorum; the replicator drains its own cursor to the hub afterwards, so WAN latency shows up as replication delay, not as produce latency.

Figure 3 is the mental model to keep. The producer acknowledgment (steps 1 to 4) is entirely local. The replicator (steps 5 to 8) is a separate consumer of the topic. When the WAN degrades, produce latency is unaffected and the replicator cursor falls behind, which surfaces as backlog. That is the correct failure mode for telemetry, and it is also why you must monitor the backlog: nothing else will tell you.

A five-region example with stated assumptions

Assume five regions with ten plants each, and assume each plant produces 5,000 messages per second at an average of 250 bytes per message including keys and properties. These are illustrative figures chosen to make the arithmetic legible, not measurements from a particular factory.

Quantity Calculation Result
Per plant raw rate 5,000 msg/s x 250 B 1.25 MB/s (about 10 Mbit/s)
Per region raw rate 10 plants x 1.25 MB/s 12.5 MB/s (about 100 Mbit/s)
Fleet raw rate 5 regions x 12.5 MB/s 62.5 MB/s (about 500 Mbit/s), roughly 5.4 TB per day
Naive full mesh of raw 12.5 MB/s x 4 peers, per region 50 MB/s (about 400 Mbit/s) egress per region
Summaries at 1:50 reduction 12.5 MB/s / 50 0.25 MB/s (about 2 Mbit/s) per region
Hub ingest of summaries 4 regions x 0.25 MB/s 1 MB/s (about 8 Mbit/s)
7-day local raw retention 12.5 MB/s x 604,800 s about 7.6 TB per region before replication
Stored on bookies at write quorum 3 7.6 TB x 3 about 22.7 TB per region

The table makes the argument for tiers. Replicating raw in a full mesh would need roughly 400 Mbit/s of sustained egress per region, before protocol overhead, TLS, and catch-up capacity. Replicating one-second summaries at a 1:50 reduction needs about 2 Mbit/s. The 1:50 ratio is an assumption; a plant that sends slow process variables gets little reduction, while one dominated by vibration waveforms gets far more. Measure your own ratio from a week of real traffic before you commit to a link size.

The last two rows are a reminder that local storage, not WAN, is usually the larger bill. Retention on bookies is multiplied by the write quorum, and long raw retention is what tiered storage is for. Offloading is disabled by default (managedLedgerOffloadThresholdInBytes=-1), so it does nothing until you configure a threshold or an automatic namespace policy.

Catch-up headroom

A replicator that has been offline for an hour must replay an hour of data on top of live traffic. With summaries at 0.25 MB/s, one hour is 900 MB, trivial. With raw at 12.5 MB/s, one hour is 45 GB, and on a 100 Mbit/s path that takes over an hour of fully saturated link with no capacity left for live data. The practical rule: size the link so that the catch-up window for your longest tolerated outage is a fraction of that outage, and rate-limit replicators so they cannot starve other traffic. Pulsar exposes a replicator dispatch rate policy for exactly this (pulsar-admin namespaces set-replicator-dispatch-rate), and replicationConnectionsPerBroker (default 16) and replicationProducerQueueSize (default 1,000) control how much parallelism and in-flight data each broker devotes to a remote link.

Long-fat networks deserve a special note. Replication throughput per connection is bounded by TCP window size divided by round-trip time. A 250 ms intercontinental path with a default window can leave a large pipe mostly idle. Increase connections, tune socket buffers, and test with representative message sizes rather than trusting link capacity numbers from the carrier.

Monitoring the backlog like a safety metric

Watch three signals per replicator: the backlog (messages or bytes behind), the outbound replication rate, and the replication delay. Pulsar exports replication metrics to Prometheus, and topic stats from bin/pulsar-admin topics stats include a per-remote-cluster replication section. The exact metric names have shifted across versions, so verify them against the metrics reference for your release rather than copying names from older blog posts, including the old version of this one, which named two metrics without qualification.

Alert on trend, not on threshold alone: a backlog that grows for fifteen minutes while the source rate is flat means the path is degraded or the destination is slow. Also set a backlog quota policy deliberately. A quota with an eviction or producer-blocking policy interacts with replicator cursors, and choosing the wrong policy can either drop unreplicated messages or stall local producers during a WAN outage. Test that behavior in staging by cutting the link, because the correct choice differs between a historian that must not lose data and a dashboard feed that should.

Edge Pre-Processing with Pulsar Functions

The cheapest byte to replicate is the one you never produce. Pulsar Functions give you a lightweight consume-apply-publish runtime that fits reduction jobs: filter heartbeats, compute windowed aggregates, attach asset metadata, and route to the right output topic. Functions can be written in Java, Python, or Go.

The old post described Functions as “embedded in the broker” with “sub-millisecond latency.” Both statements need correcting. Functions run in function workers, which can run alongside the brokers or as a separate deployment; the documentation notes that separate workers keep Functions from consuming broker resources. Instances run as threads (Java only), as processes, or as Kubernetes StatefulSets. Latency depends on the runtime, batching, and topic hops, and is typically milliseconds rather than sub-millisecond. Delivery is at-least-once by default, with at-most-once and effectively-once (deduplicated output) as options.

# rms_summary.py - illustrative 1-second RMS window per asset
import math, json
from pulsar import Function

class RmsSummary(Function):
    def process(self, item, context):
        msg = json.loads(item)
        asset = msg["asset"]
        key = f"win:{asset}:{int(msg['ts'] // 1000)}"
        state = context.get_state(key)
        sq, n = (0.0, 0) if state is None else tuple(json.loads(state))
        v = float(msg["v"])
        sq, n = sq + v * v, n + 1
        context.put_state(key, json.dumps([sq, n]))
        # Emission of closed windows is done by a companion timer topic
        # or by a downstream Flink job; keep this function stateless where you can.
        return None

The code makes one honest point: windowed logic with late data, out-of-order events, and closed-window emission is where Functions stop being the right tool. For anything beyond stateless filtering, simple counters, and routing, use a stream processor. Our comparison of Flink, Spark Streaming, and Kafka Streams covers when to reach for Flink, which also has a mature Pulsar connector.

Deployment is a single command once packaged:

bin/pulsar-admin functions create \
  --tenant plant-eu --namespace raw --name rms-summary \
  --py rms_summary.py --classname rms_summary.RmsSummary \
  --inputs persistent://plant-eu/raw/vibration \
  --output persistent://plant-eu/summary/vibration-1s \
  --processing-guarantees EFFECTIVELY_ONCE

Two design rules help. First, keep the output topic in the replicated namespace and the input in the local one; this makes the reduction step the only door through which data reaches the WAN. Second, treat the reduction function as the privacy boundary: strip or hash fields there. That keeps sovereignty a property of data flow rather than of downstream filtering.

Delivery semantics you can actually rely on

The old post’s promise of cross-region exactly-once was overstated, and the FAQ answer contradicted the body. The accurate statement: within a cluster, Pulsar offers producer-side deduplication (disabled by default, brokerDeduplicationEnabled=false) and transactions in supported configurations; replication is asynchronous at-least-once at the boundary. For telemetry, the robust design is to make sinks idempotent using a natural key such as asset ID plus source timestamp plus sequence, and to treat duplicates as normal. Iceberg upserts by key, or a time-series store that overwrites identical points, achieves the end-to-end effect without depending on a cross-cluster protocol. See the Iceberg, Delta, and Hudi decision record for sink trade-offs.

Telemetry evolves: sensors get firmware, fields get added. Register schemas per topic, use the same compatibility strategy in every cluster, and test a backward-compatible evolution across the link in staging before it happens on a Friday. Whether you use Avro, Protobuf, or JSON schema matters less than agreeing on one and forbidding untyped producers on replicated namespaces. Sparkplug B payloads are Protobuf-encoded; if you preserve them as opaque bytes on the topic, replication is trivial but consumers lose introspection. Our Sparkplug B and Unified Namespace guide discusses the structuring trade-off.

Failover: Replicated Subscriptions and Client Behavior

Replication moves messages; it does not move consumers. By default, “geo-replication replicates topic data, not subscriptions,” as the concepts documentation puts it. If a consumer group in Frankfurt fails over to the hub, it arrives at a cluster that has the data but has no memory of how far the group had read. Without extra machinery, you choose between reprocessing from the earliest retained message and skipping to the latest, and both are wrong for a maintenance analytics job.

Replicated subscriptions close most of that gap. When a consumer sets replicateSubscriptionState(true), brokers periodically write snapshot markers into the replicated topic. By correlating markers across clusters, the system can advance the mark-delete position of the subscription in the remote cluster to a point known to be safe. The documentation describes the position as kept in sync across clusters within a sub-second timeframe under normal conditions, with the snapshot frequency (replicatedSubscriptionsSnapshotFrequencyMillis, default 1,000 ms) and timeout (30 seconds by default) configurable. It requires two-way geo-replication between the participating clusters.

Consumer<byte[]> consumer = client.newConsumer()
    .topic("persistent://plant-eu/events/alarms")
    .subscriptionName("reliability-alerts")
    .subscriptionType(SubscriptionType.Failover)
    .replicateSubscriptionState(true)
    .subscribe();

The caveat is in the documentation and worth repeating: replicated subscriptions are intended for failover, not active-active consumption. If consumers in both clusters read the same replicated subscription at the same time, most messages are processed twice and some may be processed in either cluster. Because the synchronization is snapshot-based and asynchronous, a failover can also re-deliver a window of already-processed messages. Consumers must be idempotent, which we already required for the sinks.

Client failover decision flow with replicated subscription state and switch-back delay

Figure 4: The client probes the primary, switches after a configured failure window, resumes near the replicated position, tolerates duplicates, and returns only after the primary has been stable for a switch-back delay.

Client-side cluster failover

Figure 4 depends on the client noticing the outage. Pulsar’s Java client offers two mechanisms. AutoClusterFailover probes the primary service URL and switches to a secondary after the primary has been continuously failing longer than failoverDelay; it switches back after the primary has been healthy for switchBackDelay, with a probe checkInterval defaulting to 30 seconds. ControlledClusterFailover instead asks an external URL provider which cluster to use, which is the better fit when operations, not a timer, should decide, for example after a controlled maintenance or a DR declaration.

ServiceUrlProvider failover = AutoClusterFailover.builder()
    .primary("pulsar+ssl://broker.eu-fra.example.com:6651")
    .secondary(List.of("pulsar+ssl://broker.na-hub.example.com:6651"))
    .failoverDelay(60, TimeUnit.SECONDS)
    .switchBackDelay(15, TimeUnit.MINUTES)
    .checkInterval(30, TimeUnit.SECONDS)
    .build();

PulsarClient client = PulsarClient.builder()
    .serviceUrlProvider(failover)
    .build();

Choose the delays with your data in mind. A sixty-second failover delay avoids flapping on a brief network blip, and a long switch-back delay avoids the worse failure of bouncing consumers between clusters during an unstable recovery. For gateway producers inside a plant, failing over to a different region is usually the wrong behavior: it would push data into a cluster whose namespace may not accept that tenant’s data. Producer failover for sovereign data should stay within the region, for instance to a second availability zone or a standby cluster in the same jurisdiction.

Synchronous replication is a different tool

The concepts documentation also describes a stronger, synchronous mode built on BookKeeper, where a write is acknowledged only after a majority of data centers persist it. That gives stronger consistency at the price of cross-datacenter round trips on every produce. For IIoT it fits the rare stream where losing an acknowledged message across a site loss is unacceptable and the sites are close enough for the latency to be tolerable, such as two campuses in one metro area. For continent-scale plant networks, asynchronous replication with idempotent consumers is the sensible default.

Comparing the Alternatives in September 2026

Geo-replication is no longer a differentiator that only Pulsar can claim, so the comparison has to be honest about what each option actually provides.

Kafka with MirrorMaker 2 (part of Apache Kafka via Kafka Connect) replicates topics, translates consumer group offsets, and syncs some configuration between clusters. It is open source and widely understood. The trade-off is that it is a separate distributed system to run, tune, and monitor, and offset translation is an approximation you must test during DR drills. Our Kafka versus Redpanda versus WarpStream ADR covers the single-cluster side of that ecosystem.

Confluent Cluster Linking is a Confluent feature that mirrors topics with offsets preserved between clusters without a separate Connect deployment. The old post said it “requires Confluent Cloud”; it is also part of Confluent Platform’s commercial offering. It is proprietary and priced with the platform.

Redpanda Shadowing is documented as an enterprise feature for asynchronous cross-region disaster recovery, with shadow topics that become writable on failover, and it supports unidirectional and bidirectional topologies. This corrects the old claim that Redpanda had no comparable feature. Check licensing and whether your version includes it.

StreamNative Ursa and Universal Linking target a different design center. StreamNative describes Ursa as a leaderless, lakehouse-native engine that writes directly to object storage in Iceberg or Delta formats, natively supports both Kafka and Pulsar protocols, and is aimed at latency-relaxed workloads; Universal Linking mirrors Kafka-compatible clusters into it. The vendor’s own cost claims (“10x lower cost” in some scenarios) are marketing figures whose validity depends on workload and topology. Treat Ursa as a commercial platform choice, not as an Apache Pulsar 4.x feature, and note that the vendor pages do not detail geo-replication for Ursa specifically.

NATS JetStream offers mirroring and sourcing between streams and suits lightweight edge-to-core topologies; see NATS JetStream versus Kafka for edge telemetry.

Requirement Pulsar geo-replication Kafka + MirrorMaker 2 Cluster Linking or Shadowing Ursa and Universal Linking
Per-namespace replication scope Native, with tenant allow-list Per-topic patterns you maintain Per-topic link config Vendor platform config
Offset or subscription continuity Replicated subscriptions, failover only Offset translation, approximate Offsets preserved by design Offsets preserved for Kafka sources
Separate service to operate No, in-broker replicators Yes, Connect cluster No, in-broker Managed or private cloud
Open-source core Yes, Apache 2.0 Yes No, commercial No, commercial platform
Best fit Many regions, tiered scope Existing Kafka estate Kafka estate with budget Lakehouse-first, latency-relaxed

The table is a positioning aid, not a benchmark. If your organization runs Kafka everywhere and needs two-site DR, the switching cost to Pulsar will not be repaid by replication alone. If you are designing a multi-region, multi-tenant telemetry backbone from scratch and want scope expressed as policy, Pulsar remains a strong fit.

Trade-offs, Gotchas, and What Goes Wrong

Removing a cluster from a replication set is not a supported rollback. The documentation warns that changing the clusters policy at namespace or topic level can trigger topic deletion on the excluded clusters, and that replicating data to a remote cluster and then removing that cluster from replication is not a supported use case. Treat replication scope changes as migrations with backups, never as toggles. A compliance-driven “stop sending EU data to the hub” ticket can delete the topic on the hub.

Auto-creation races. Enabling replication on a namespace whose topics do not yet exist on the remote side used to produce initialization races; PIP-433 in 4.2.0 addresses the ordering. If you are on an older patch release, pre-create topics and upgrade before relying on automatic creation.

Topic explosion. Per-asset topics look tidy and destroy brokers. Replication multiplies the cost: every replicated topic gets a replicator, a cursor, and a remote producer. Prefer partitioned topics keyed by asset with a modest partition count, and let consumers filter by key or property. Tens of thousands of active replicated topics per cluster call for load testing, not optimism.

Ordering assumptions. Ordering is per source per partition. If two regions both write the same replicated topic, do not build logic that assumes global order; carry event time and source ID in the payload.

Clock and timestamp semantics. Replication preserves the original publish and event times in message metadata, but consumers frequently key windows off arrival time. After an outage, the replicated backlog arrives late; windowing by event time with watermarks avoids double-counting in the hub.

Metadata store dependency. Each cluster’s metadata store is a hard dependency for topic ownership and creation. A tiny three-node ensemble on shared plant hardware is a common single point of failure. Do not attempt a ZooKeeper-to-Oxia migration on 4.x: the current documentation requires 5.0 or later for all components, and 5.0 is only at milestone stage as of September 2026.

Certificates on aging links. Replication uses TLS between clusters; expired certificates stop replication silently unless you alert on the backlog. Automate rotation and test it.

Security of the bridge. Give each cluster’s replication identity the minimum permissions on the namespaces it replicates. A compromised replicator credential that can publish to any namespace on the hub is a lateral-movement path from an OT site into IT analytics.

Practical Recommendations

Start by classifying your streams into the four tiers, because that classification drives everything else. Write the tier as a namespace policy in infrastructure as code, review changes to allowed-clusters like firewall rule changes, and forbid producers from setting message-level replication overrides on Tier 1 and Tier 4 namespaces.

Then measure before you size. Capture a week of real traffic, compute bytes per second per plant per tier, and measure the reduction ratio your Functions or Flink jobs actually achieve. Size WAN links for the replicated tiers only, with a catch-up window you can defend. Put backlog alerts on every replicator from day one, and rehearse a WAN cut in staging until the numbers match your model.

On versions, deploy the 4.0 LTS or the latest 4.2 patch release depending on your appetite for change, and plan the 4.0-to-later upgrade before its active support window ends in October 2026. Keep metadata on ZooKeeper until a 5.x release is generally available and the migration guide fits your rollback requirements. Pin client and broker versions together, and read the geo-replication notes in each release before upgrading, since replicator fixes have shipped in nearly every feature release.

  • [ ] Streams classified into four tiers, each a namespace with a recorded replication set
  • [ ] Tenants created with allowed-clusters limited to lawful destinations
  • [ ] Bookie ensemble 3, write quorum 3, ack quorum 2 (or justified alternative) and at least five bookies per region
  • [ ] Journal on a separate fast device; offloading configured with an explicit threshold
  • [ ] Reduction functions or Flink jobs strip identifiers before the replicated namespace
  • [ ] Sinks and consumers idempotent by natural key
  • [ ] Replicator rate limits set; backlog and delay alerts live
  • [ ] Failover delays chosen; producers restricted to in-jurisdiction failover
  • [ ] Quarterly drill: cut the WAN, restore, verify backlog drains and no duplicates leak into reports

Frequently Asked Questions

Does Apache Pulsar geo-replication guarantee exactly-once delivery across regions?

No. Replication is asynchronous and effectively at-least-once at the boundary. Pulsar supports producer deduplication (off by default) and transactions within supported configurations, but cross-cluster exactly-once is not something to depend on. Design consumers and sinks to be idempotent using a natural key such as asset ID, source timestamp, and sequence. Then duplicates from replication, failover, or retries are harmless, and you avoid coupling correctness to a protocol that spans a WAN.

Can I replicate some topics but not others inside one cluster?

Yes. Replication scope is a policy on the namespace, with optional topic-level overrides when topic-level policies are enabled, and a producer can narrow the target set per message. The tenant’s allowed-clusters list is the ceiling that none of those can exceed. This is the mechanism for keeping raw telemetry local while summaries go to a hub and operational events go everywhere.

Are replicated subscriptions the same as active-active consumers?

No. Replicated subscriptions keep the subscription’s mark-delete position roughly in sync across clusters so a consumer group can fail over and resume near where it left off. The documentation states they are intended for failover, not simultaneous consumption in several clusters, where messages are largely processed twice. Run one active consumer cluster at a time and expect a small window of redelivery after failover.

Producers keep working, because acknowledgment depends only on the local bookie quorum. Each replicator’s cursor stays put, so a backlog accumulates on the source cluster and drains after the link returns. The risks are disk pressure from retained backlog, backlog quota policies that evict or block, and a catch-up burst that competes with live traffic. Size retention, choose the quota policy deliberately, and rate-limit replicators.

Should I use Pulsar or Kafka MirrorMaker 2 for plant telemetry?

If you already run Kafka across the estate and need simple two-site disaster recovery, MirrorMaker 2 or a commercial linking feature avoids a platform migration. If you are building a multi-region, multi-tenant backbone and want replication scope, tenancy, and tiered storage expressed as broker policy, Pulsar is a strong fit. Decide on operating model and team skills first; replication features alone rarely justify a switch.

Which Pulsar version should I deploy in late 2026?

Use the 4.0 LTS line if you value stability, remembering its active support window ends on 21 October 2026 and security-only support continues for a further year, or use the current 4.2 patch release for newer geo-replication fixes. Do not run the 5.0 milestone in production. Keep ZooKeeper for metadata until the ZooKeeper-to-Oxia migration path, which the docs tie to Pulsar 5.0, is generally available.

Further Reading

By Riju — about

Comments

No comments yet. Why don’t you start the discussion?

Leave a Reply

Your email address will not be published. Required fields are marked *