Using Pulsar Schema Registry to Prevent Silent Data Corruption
Learn how Apache Pulsar’s built‑in Schema Registry centralizes schema validation, enforces compatibility, and reduces integration bugs in event‑driven pipelines.
15 Sept 2026, 22:13 UTC

The problem: silent schema drift
In many event‑driven systems, producers and consumers evolve their message formats independently. When a producer adds, removes, or changes a field without coordinating with consumers, the broker happily stores the bytes. Later, a consumer attempts to deserialize the payload with an outdated schema and fails silently or throws an exception that surfaces only after processing lag builds up. The result is data loss, replay storms, or costly downtime that is hard to trace back to a schema mismatch.
Thesis: let Pulsar’s Schema Registry be the single source of truth
Apache Pulsar includes a built‑in Schema Registry that lives alongside the broker and stores schemas in BookKeeper. By registering a schema once per topic, the broker can validate every incoming message against that schema before persisting it. Consumers configured to use the registry automatically fetch the matching schema by ID, guaranteeing that deserialization uses the exact contract the producer intended.
How the registry works in practice
- Registration – A developer posts an Avro, JSON, or Protobuf schema via
pulsar-admin schemasor the REST API. The registry assigns a monotonically increasing version‑ID and stores the schema under the topic name. - Producing – The producer client (e.g., the Pulsar Java client with the schema‑registry serializer) looks up the latest schema ID, embeds that ID in the message metadata, and sends the payload. The broker checks that the ID matches a known schema and validates the payload.
- Consuming – The consumer client uses the schema‑registry deserializer. It reads the schema ID from the message, fetches the schema from the registry (usually via a local cache), and deserializes the payload. If the broker had rejected the message earlier, it never reaches the consumer.
Because validation happens at the broker, badly formed messages are kept out of the topic entirely, eliminating the "poison pill" scenario.
Worked example: registering an Avro schema and enforcing compatibility
Where to run: on a host with a Pulsar standalone installation (the bin/pulsar standalone command starts broker, BookKeeper, and the schema registry by default). You need OS‑level permission to execute pulsar-admin and to run a Java client.
- Start Pulsar
bin/pulsar standalone - Define a simple Avro schema (save as
order-v1.avsc):{ "type": "record", "name": "Order", "namespace": "com.example", "fields": [ {"name": "orderId", "type": "string"}, {"name": "amount", "type": "double"}, {"name": "items", "type": {"type": "array", "items": "string"}} ] } - Register the schema for topic
public/default/orderspulsar-admin schemas create public/default/orders \ --schema-file order-v1.avsc \ --schema-type AVROThe command returns a schema ID (e.g.,
0). No special privileges beyond being able to executepulsar-adminare required. - Java producer with schema‑registry serializer
Add the dependency (Maven):
<dependency> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-client-avro</artifactId> <version>3.0.0</version> </dependency>Producer code (replace placeholders):
import org.apache.pulsar.client.api.*; import org.apache.pulsar.client.api.schema.*; import java.util.*; public class OrderProducer { public static void main(String[] args) throws Exception { PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Schema avroSchema = Schema.AVRO(GenericRecord.class, new SchemaDefinition("order-v1.avsc")); Producer producer = client.newProducer(avroSchema) .topic("public/default/orders") .create(); // Valid message GenericRecord valid = new GenericData.Record(avroSchema.getSchemaInfo().getSchema()); valid.put("orderId", "ord-123"); valid.put("amount", 42.5); valid.put("items", Arrays.asList("apple", "banana")); producer.send(valid); System.out.println("Sent valid order"); // Invalid message – missing required field "amount" GenericRecord invalid = new GenericData.Record(avroSchema.getSchemaInfo().getSchema()); invalid.put("orderId", "ord-999"); // amount omitted intentionally invalid.put("items", Arrays.asList("orange")); try { producer.send(invalid); } catch (PulsarClientException e) { System.out.println("Broker rejected invalid message: " + e.getMessage()); } producer.close(); client.close(); } }When the producer attempts to send the invalid record, the broker returns a
SchemaValidationErrorand does not persist the message. You can verify this by checking the topic’s backlog withpulsar-admin topics stats public/default/orders– the msgRateIn should not increase for the rejected send. - Consumer that uses the registry deserializer
Consumer consumer = client.newConsumer(avroSchema) .topic("public/default/orders") .subscriptionName("my-sub") .subscribe(); while (true) { Message msg = consumer.receive(); System.out.println("Received: " + msg.getValue()); consumer.acknowledge(msg); }Only the valid message will be delivered; the invalid one never reaches this loop.
Trade‑offs and limitations
- Network overhead – Each producer/consumer performs a schema lookup (typically a RPC to the BookKeeper‑backed registry). In high‑throughput scenarios this latency can become noticeable unless you enable client‑side caching (the Java client caches schemas by ID for a configurable TTL).
- Operational complexity** – The registry adds another moving part: BookKeeper must be healthy and sufficiently provisioned. If BookKeeper is unavailable, schema fetches fail and producers may be blocked until the registry recovers. Monitoring BookKeeper ledger lag and read/write latency is therefore essential.
- Compatibility mode choice** – The registry lets you enforce BACKWARD, FORWARD, or FULL compatibility. A overly strict mode (e.g., FULL when you need to drop a field) will block legitimate evolution, while a lax mode (NONE) defeats the purpose. Choose the mode that matches your product’s versioning policy and test it in a staging environment before promoting to production.
Actionable checklist
- Deploy a Pulsar standalone cluster (or use your existing cluster) and verify that the schema registry is enabled (
bin/pulsar statusshows the registry component). - Register a schema for each critical topic using
pulsar-admin schemas create. - Configure producer and consumer clients to use the schema‑registry serializer/deserializer for the appropriate schema type (Avro, JSON, Protobuf).
- Set a compatibility mode that reflects your evolution strategy (
pulsar-admin schemas set-compatibility --schema-compatibility BACKWARD). - Enable client‑side caching (e.g.,
schemaCacheSizeandschemaCacheExpiryin the Java client) to reduce lookup latency. - Monitor BookKeeper health (ledger size, write latency) and the registry’s request/response metrics via Pulsar’s Prometheus endpoint.
- In staging, deliberately publish a message that violates the registered schema and confirm that the broker returns a schema‑validation error and that the message does not appear in consumer logs.
By treating the schema as a first‑class artifact stored in Pulsar’s Schema Registry, you shift schema‑related bugs from runtime surprises to pre‑commit checks, dramatically reducing the chance of silent data corruption in production.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.