Stopping Schema Drift in Streaming Pipelines with Pulsar's Built-In Schema Registry
Schema drift corrupts streaming data silently. Pulsar's built-in schema registry enforces compatibility broker-side — here's how it works, an Avro example, and the trade-offs.
12 Jul 2025, 17:30 UTC

The silent failure mode nobody pages you for
A producer team adds a field. A consumer team, three services away, still deserializes with the old definition. Nothing crashes — the consumer just starts writing nulls into a column, and you find out two weeks later from a finance report. This is schema drift, and in byte-oriented streaming setups it is one of the most common ways pipelines corrupt data without raising a single error.
Apache Pulsar addresses this directly: it ships with a built-in schema registry that lives on the broker side, enforces compatibility rules at produce time, and attaches a schema version to every message. The useful takeaway is that you get schema safety without deploying a separate registry service — but only if you pick a compatibility strategy deliberately and keep it enforced.
How the registry actually works
When a producer connects to a typed topic, it registers its schema with the broker. Pulsar stores schema definitions in its metadata store (ZooKeeper in classic deployments, or the configured metadata store in newer ones) and assigns each version an ID. Every message carries that ID, so a consumer can fetch the exact schema the producer used and deserialize correctly — even when producer and consumer are on different schema versions, as long as the versions are compatible.
The broker enforces a schema compatibility check before accepting a new schema version. The main strategies are:
- BACKWARD — new consumers can read data written with the previous schema. Typical when consumers upgrade first.
- FORWARD — old consumers can read data written with the new schema. Typical when producers upgrade first.
- FULL — both directions hold. The safest default for most teams.
- ALWAYS_COMPATIBLE / NONE — no enforcement. Avoid in production; it defeats the point.
Because the check happens broker-side, a bad schema change is rejected at the producer — before a single incompatible message lands on the topic. That is the core win over client-side-only validation.
A worked example with Avro
Assume Pulsar 2.10+ or 3.x, with a Java client. Define a POJO and produce with an Avro schema:
// Producer side — run in your application, needs produce permission on the topic
Producer<User> producer = client.newProducer(Schema.AVRO(User.class))
.topic("persistent://team/events/users")
.create();
User u = new User();
u.id = 1;
u.name = "Alice";
producer.send(u); // schema {id, name} registered automaticallyLater, the producer team adds an email field with a default value:
public class User {
public int id;
public String name;
public String email = ""; // default keeps it BACKWARD compatible
}With the topic's compatibility set to BACKWARD (or FULL), the broker accepts this new version because old data can still be read by the new schema — the missing email resolves to the default. Consumers on the old schema keep working; consumers that upgrade see the new field. Had the field been added without a default, the broker would reject the schema registration and the producer would fail fast with an incompatible-schema error.
To set the strategy, run this with the Pulsar admin CLI (requires tenant/namespace admin permissions):
pulsar-admin schemas set-schema-compatibility-strategy \
--compatibility BACKWARD \
persistent://team/events/usersVerify with pulsar-admin schemas get persistent://team/events/users and confirm the returned schema and version match what you expect. You can also check the enforced strategy via pulsar-admin schemas get-schema-compatibility-strategy on the topic.
Trade-offs worth knowing upfront
The registry is not free. Each new producer registration involves a metadata-store lookup, so schema availability is coupled to your ZooKeeper (or metadata store) ensemble — if that layer is unhealthy, schema operations degrade too. In practice, schema info is cached broker-side, so steady-state message flow is unaffected, but cold producers and schema updates depend on it.
Second, compatibility enforcement is only as good as your strategy choice. Setting NONE to "unblock" a release is a common shortcut that quietly reintroduces the exact drift problem the registry exists to prevent. Treat compatibility config like a migration policy: review changes to it the way you review database migrations.
Third, schema evolution rules differ by format. Avro's default-value rules are forgiving; Protobuf and JSON Schema have their own constraints. Test an evolution in a staging cluster before assuming it passes.
What to do next
Three concrete steps: (1) pick a compatibility strategy per namespace — FULL is a reasonable default — and set it before the first production message, since retrofitting is harder; (2) add a staging check that produces with the proposed new schema and confirms the broker accepts it, so incompatibilities surface in CI rather than at deploy time; (3) monitor schema-related metrics and topic state in Pulsar Manager or via pulsar-admin schemas get in your alerting loop, so an unexpected schema version bump from any team is visible immediately. Schema drift is a process problem as much as a technical one — the registry gives you the enforcement point, and these habits make it stick.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.