Overview
Schema-less Kafka topics are a silent failure mode. The producer publishes JSON with a renamed field, the consumer silently deserialises a null into a field it expected to be mandatory, the downstream process writes a zero amount to the general ledger, and the first person to notice is the reconciliation team two days later. There are no exceptions, no error logs, no alerts. The schema mismatch is absorbed by the consumer’s null-tolerance and manifests as corrupted business data.
The contract problem between producer and consumer teams is fundamentally an organisational problem dressed as a technical one. In a microservices architecture, the payments team owns the producer and the fraud team owns the consumer. When the payments team adds a field, the fraud team’s consumer should continue working. When they rename a field, it should not. Without an enforced schema contract, the only mechanism that prevents breakage is inter-team communication — which does not scale, is not auditable, and fails silently when it fails at all.
Confluent Schema Registry is the API contract layer for event streams. It stores versioned schemas, enforces compatibility rules before a new schema version is registered, embeds a schema ID in every message on the wire, and provides a REST API that serialisers and deserialisers call automatically. The producer cannot publish an incompatible schema without the registry refusing the registration. The consumer cannot silently absorb a field change — it fetches the exact schema the producer used and deserialises deterministically.
Confluent’s wire format for schema-registered messages begins with a magic byte (0x00) followed by a 4-byte big-endian schema ID. A raw JSON or Avro message that was not serialised with KafkaAvroSerializer will not start with this magic byte. If your consumer uses KafkaAvroDeserializer, it will throw a SerializationException on the first byte. Always validate that every producer on a registered topic uses the schema-registry serialiser — a single legacy producer publishing raw bytes will break all schema-aware consumers on that topic immediately.
Avro vs Protobuf vs JSON Schema
Schema Registry 7.x supports three serialisation formats. The right choice depends on your team’s constraints, not on benchmarks in isolation. For payment event streams in a regulated environment the differences are operationally significant.
| Criterion | Avro | Protobuf | JSON Schema |
|---|---|---|---|
| Wire size | Smallest — binary, no field names on wire | Small — binary, field numbers only | Largest — text, full field names |
| Compatibility enforcement | Strong — full forward/backward via registry | Strong — field number stability required | Partial — harder to enforce additive-only |
| Code generation | Good (Java/Python/C#) | Excellent — first-class tooling across 12 languages | Poor — inconsistent across generators |
| ISO 20022 mapping ease | Good — unions handle nullable fields | Moderate — optional maps well, nesting verbose | Natural — JSON Schema mirrors MX XML structure |
| Schema registry support | Native — original format | Native from SR 6.0+ | Native from SR 6.1+ |
Recommendation for payment event streams: Avro for internal Kafka topics where throughput and schema governance are priorities; Protobuf for cross-team or cross-organisation event buses where polyglot consumers matter; JSON Schema only for topics that must remain human-readable for audit without tooling. For ISO 20022 message envelopes, Avro with decimal logicalType is the strongest choice — the binary compactness reduces Kafka storage costs on high-volume IPS and SARIE feeds, and the registry enforces the field-level contract that regulated message formats demand.
Schema Registry Architecture
Confluent Schema Registry is a separate JVM process (or set of replicated processes) that sits alongside your Kafka cluster. It stores schemas in an internal Kafka topic (_schemas), exposes a REST API on port 8081, and is addressed by serialisers and deserialisers through a schema.registry.url configuration property.
Three concepts govern how schemas are stored: subjects, schema IDs, and versions. A subject is a named scope for schema evolution, typically following a TopicNameStrategy convention: a topic called payments.pacs008 has subjects payments.pacs008-key and payments.pacs008-value. An alternative, RecordNameStrategy, scopes schemas by their fully-qualified record name rather than by topic — useful when the same schema type is published to multiple topics and you want a single evolution history. A schema ID is a globally unique integer assigned at registration time; it is embedded in every message wire format and used by the consumer to fetch the exact schema for deserialisation. A version is the sequence number within a subject — each new compatible schema registration increments the version.
Compatibility Modes
The compatibility mode on a subject controls what schema changes are allowed when a new version is registered. This is the core governance lever — set it too loosely and the registry provides no meaningful protection; set it too strictly and schema evolution becomes a negotiation that slows delivery.
| Mode | What it allows | What it forbids | Consumer must upgrade? |
|---|---|---|---|
| BACKWARD | Add optional fields (with defaults), remove optional fields | Rename fields, add required fields, change types | No — old consumer reads new messages |
| FORWARD | Add required fields, remove optional fields | Rename fields, remove required fields | Yes — new consumer reads old messages |
| FULL | Add optional fields with defaults only | Remove any field, rename, change types | No in either direction |
| BACKWARD_TRANSITIVE | Same as BACKWARD but checked against all prior versions | Anything not BACKWARD-compatible with every prior version | No — even very old consumers work |
| NONE | Any change | Nothing — no compatibility check performed | Unknown — undefined |
Production guidance for payment event streams: default to FULL. This is the most conservative useful mode — both old consumers reading new messages and new consumers reading old messages work correctly. FULL requires that every new schema version adds only optional fields with default values, which is a discipline that forces explicit migration planning for breaking changes. Use BACKWARD_TRANSITIVE for topics where consumers may be several schema versions behind — common on regulated replay or audit topics where consumers may be recovering from a long outage. Reserve NONE only for development subjects where schema iteration is rapid and no consumer SLA exists.
Do not relax compatibility to FORWARD or BACKWARD on topics that carry payment instructions or transaction events without a documented migration plan reviewed by both the producing and consuming team leads. A FORWARD-only schema allows the producer to add required fields that break old consumers silently if those consumers are not upgraded first — exactly the failure mode Schema Registry exists to prevent. The governance board should review and approve any compatibility relaxation request as a change with a migration plan attached.
Compatibility is enforced in the CI pipeline using the schema-registry-maven-plugin or by calling the registry’s /compatibility endpoint directly. Register schemas in lower environments (dev, test) before merging to main; fail the pipeline if compatibility is rejected. See the CI/CD integration section for configuration details.
ISO 20022 in Avro
ISO 20022 messages — pacs.008, pacs.002, camt.053, pain.001 — are XML-first in the standard, but transporting them as raw MX XML on Kafka topics wastes considerable bandwidth, lacks schema enforcement, and provides no type safety. Mapping ISO 20022 message elements to Avro records gives you binary compactness, registry-enforced compatibility, and generated Java classes that type-check at compile time.
The key mapping rules: ISO 20022’s XML element hierarchy maps to Avro nested records; XML attributes become Avro fields; optional XML elements (minOccurs="0") become Avro union types (["null", actualType]) with a null default; ISO 20022’s ActiveOrHistoricCurrencyAndAmount becomes an Avro record with two fields — an amount and a currency code.
{
"type": "record",
"name": "PaymentEvent",
"namespace": "info.wbadawi.payments.avro",
"doc": "pacs.008.001.09 FIToFICustomerCreditTransfer envelope — internal event schema",
"fields": [
{
"name": "msgId",
"type": "string",
"doc": "GrpHdr/MsgId — ISO 20022 Max35Text"
},
{
"name": "creDtTm",
"type": { "type": "long", "logicalType": "timestamp-millis" },
"doc": "GrpHdr/CreDtTm — creation timestamp, epoch millis UTC"
},
{
"name": "nbOfTxs",
"type": "int",
"doc": "GrpHdr/NbOfTxs"
},
{
"name": "instrAmt",
"type": { "type": "bytes", "logicalType": "decimal", "precision": 18, "scale": 5 },
"doc": "CdtTrfTxInf/IntrBkSttlmAmt — settlement amount, decimal, scale=5"
},
{
"name": "instrAmtCcy",
"type": "string",
"doc": "IntrBkSttlmAmt Ccy attribute — ISO 4217 code e.g. SAR"
},
{
"name": "endToEndId",
"type": "string",
"doc": "CdtTrfTxInf/PmtId/EndToEndId"
},
{
"name": "cdtrIBAN",
"type": ["null", "string"],
"default": null,
"doc": "Cdtr/Acct/Id/IBAN — nullable; not all IPS transfers have IBAN"
},
{
"name": "dbtrIBAN",
"type": ["null", "string"],
"default": null,
"doc": "Dbtr/Acct/Id/IBAN — nullable"
},
{
"name": "pdplClassification",
"type": "string",
"default": "CONFIDENTIAL",
"doc": "PDPL data classification: PUBLIC | INTERNAL | CONFIDENTIAL | RESTRICTED"
},
{
"name": "schemaVersion",
"type": "int",
"default": 1,
"doc": "Internal schema version for migration tracking"
}
]
}
IEEE 754 floating-point cannot represent most decimal fractions exactly. A SAR 100.10 amount stored as a double may round-trip as 100.09999999999999. In a payment system this is a reconciliation defect and potentially a regulatory finding. Always use Avro’s bytes with logicalType: decimal and explicit precision/scale, or store amounts as string if the consuming team cannot handle decimal logicalType. Set scale to 5 to accommodate SAR sub-halala precision for IPS.
Schema Evolution Patterns
Schema evolution is inevitable. Fields get added as the business adds reporting requirements; the naming conventions established in year one look wrong in year three; regulatory changes require new mandatory fields. The question is not whether schemas evolve but whether they evolve safely.
-
Adding optional fields
This is the safe operation. Add a field with a default value and type
["null", T]orTwith a defined default. Old consumers that do not know the field receive its default on deserialisation. New producers that write the field are readable by old consumers. This is the only change allowed under FULL compatibility. -
Renaming fields (never rename — use aliases)
Renaming a field is a breaking change: old consumers using the old name cannot find the value. The correct approach is to use Avro’s
aliasesmechanism: add the new name as a field and mark the old name as an alias. Old consumers reading the old name continue to work; new consumers use the new name. Remove the alias only after all consumers are on the new name — typically after two schema versions. -
Deprecating fields
Mark the field in the
docproperty:"doc": "DEPRECATED: use fooNewField instead. Remove after schema version 12.". Keep the deprecated field in the schema for at least two versions. Do not remove it until all consumers have been updated and verified not to read it. -
Removing fields
Remove the consumer that reads the field before removing the field from the schema. Removing a field that a consumer depends on is a silent data loss — the field deserialises to its default value, which for most types is zero or null. Follow the sequence: deprecate → update consumers → remove field from schema.
-
Changing field types
Avro supports promotions (e.g.,
inttolong,floattodouble) as BACKWARD-compatible changes. Any other type change is incompatible and must be handled as a new field plus a migration period. In practice, type changes on payment fields are a governance event requiring a migration plan, not a schema version bump.
CI/CD Integration
Schema compatibility enforcement belongs in the CI pipeline, not in a production deployment. A failed compatibility check in CI is a build failure that prevents the code from reaching a deployment; a failed compatibility check in production is an outage. The investment in CI integration pays back on the first incompatible change a developer makes to a shared schema.
Two integration points: a pre-commit hook that runs compatibility against the development Schema Registry before the developer pushes; and a CI pipeline step that registers the schema and blocks merge if registration fails. The Maven plugin approach integrates with both:
<!-- Confluent Schema Registry Maven Plugin -->
<plugin>
<groupId>io.confluent</groupId>
<artifactId>kafka-schema-registry-maven-plugin</artifactId>
<version>7.7.0</version>
<configuration>
<schemaRegistryUrls>
<param>${schema.registry.url}</param>
</schemaRegistryUrls>
<subjects>
<!-- subject name maps to schema file -->
<payments.pacs008-value>
src/main/avro/pacs008-payment-event.avsc
</payments.pacs008-value>
</subjects>
<compatibilityLevels>
<!-- enforce FULL on payment topics -->
<payments.pacs008-value>FULL</payments.pacs008-value>
</compatibilityLevels>
</configuration>
<executions>
<!-- test-compatibility: check but don't register (for PR validation) -->
<execution>
<id>test-compatibility</id>
<phase>validate</phase>
<goals><goal>test-compatibility</goal></goals>
</execution>
<!-- register: register on merge to main -->
<execution>
<id>register</id>
<phase>deploy</phase>
<goals><goal>register</goal></goals>
</execution>
</executions>
</plugin>
#!/usr/bin/env bash
# Fail the CI build if the new schema is incompatible with the latest registered version.
# Run this in the PR validation stage against the dev registry.
SCHEMA_FILE="src/main/avro/pacs008-payment-event.avsc"
SUBJECT="payments.pacs008-value"
REGISTRY="https://schema-registry.dev.internal:8081"
# Escape schema JSON for embedding in request body
SCHEMA_JSON=$(cat "${SCHEMA_FILE}" | python3 -c "import sys,json; print(json.dumps(sys.stdin.read()))")
RESPONSE=$(curl -s -o /dev/null -w "%{http_code}" \
-X POST \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data "{\"schema\": ${SCHEMA_JSON}}" \
"${REGISTRY}/compatibility/subjects/${SUBJECT}/versions/latest")
if [ "${RESPONSE}" != "200" ]; then
echo "ERROR: Schema compatibility check failed (HTTP ${RESPONSE}) for subject ${SUBJECT}"
echo "Review compatibility rules before merging. Compatibility mode: FULL"
exit 1
fi
COMPAT=$(curl -s \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data "{\"schema\": ${SCHEMA_JSON}}" \
"${REGISTRY}/compatibility/subjects/${SUBJECT}/versions/latest" \
| python3 -c "import sys,json; d=json.load(sys.stdin); print(d.get('is_compatible','false'))")
if [ "${COMPAT}" != "true" ]; then
echo "ERROR: Schema is not FULL-compatible with latest version of ${SUBJECT}"
exit 1
fi
echo "Schema compatibility check passed for ${SUBJECT}"
Promote schemas through environments using separate subject namespaces: dev.payments.pacs008-value, test.payments.pacs008-value, prod.payments.pacs008-value — or separate Schema Registry instances per environment with the same subject names. The latter is cleaner for access control (production registry access restricted to CI/CD pipelines only; no developer can register directly in prod) but requires an additional registry instance.
SAMA Regulatory Reporting Schemas
SAMA’s Technology Risk Management framework and the PDPL impose specific requirements on event schemas used for regulatory reporting and audit trails. Three requirements shape schema design for regulatory topics: immutability, classification, and archival.
Immutable event schemas for audit trails. A SAMA-reportable event — a transaction that was blocked, a customer data access event, a system security alert — must be queryable with exactly the schema it was written with. Schema Registry’s versioned schema archive is the mechanism: each schema version is immutable once registered; the schema ID embedded in each message points permanently to the schema that was used. Even if the subject is deleted from the registry, a schema archive backed by the _schemas topic (with compaction disabled and infinite retention) preserves all versions.
PDPL data classification embedded as schema metadata. Add a pdplClassification field to every schema that carries personal data. The field should be a required string enum with values PUBLIC, INTERNAL, CONFIDENTIAL, or RESTRICTED. This makes data classification machine-readable at the event level, enabling downstream consumers (data governance pipelines, DLP scanners) to enforce access control without out-of-band configuration. The example Avro schema in the ISO 20022 section demonstrates this pattern.
Schema tamper evidence via registry audit log. Enable Schema Registry’s audit log integration with your SIEM. Every schema registration, compatibility mode change, and subject deletion is an auditable event. A compatibility mode change on a regulated subject from FULL to NONE should trigger an alert — it means someone has bypassed the compatibility enforcement and a breaking change may follow. Configure the registry’s REST API behind an API gateway that enforces mutual TLS and logs all requests to the SIEM.
For topics subject to SAMA replay obligations, configure the _schemas topic with retention.ms=-1 (infinite retention) and cleanup.policy=compact disabled or set to delete. This preserves every schema version permanently. When a regulatory examination requires replaying events from 18 months ago, the schema ID embedded in each message resolves to the exact schema in use at that time — there is no ambiguity about field types or nullability that existed in that version.
Pitfalls
The most destructive anti-pattern: a developer temporarily sets the subject compatibility to NONE, registers a breaking schema change to unblock a deployment, and either forgets to restore the original compatibility mode or leaves it “for now.” NONE mode means the registry performs no validation; any schema passes. If this happens in production, every consumer on the subject is now unprotected from future breaking changes. Set up an alert (Prometheus or SIEM) that fires whenever a subject compatibility mode is changed to NONE in production. Treat it as a P1 incident until the mode is restored.
KafkaAvroSerializer caches the schema ID it receives from the registry after the first registration call. If the Schema Registry is restarted or fails over and the subject is recreated with a different schema history (e.g., after a disaster recovery restore), the cached schema ID in the producer may refer to a schema that no longer exists or now maps to a different schema version. Producers that embed a stale schema ID will produce messages that consumers cannot deserialise. Configure producers with schema.registry.cache.capacity=0 only in disaster recovery scenarios; in normal operation the cache is correct. Monitor registry restart events and verify producer serialisation on all payment topics immediately after a registry failover.
Schema IDs are globally unique within a single registry instance. If dev and prod share the same registry (a common cost-saving decision), a developer registering a schema in dev may produce a schema ID that is higher than the latest prod ID — but the schema contents differ between environments. A message produced in dev accidentally published to a prod topic carries a schema ID that resolves to the wrong schema in the prod registry. The deserialised result is silently wrong. Always run separate Schema Registry instances per environment tier, or use strict namespace prefixes per environment with access control that prevents cross-environment schema reads.
Production Checklist
- Schema Registry deployed with at least 3 replicas;
_schemastopic replication factor 3, min.insync.replicas 2. - Separate Schema Registry instances (or strict namespace isolation) for dev, test, and prod environments.
- All payment topics configured with
FULLorFULL_TRANSITIVEcompatibility; NONE mode alerting enabled. - All producers using
KafkaAvroSerializerwithauto.register.schemas=falsein production. - All consumers using
KafkaAvroDeserializerwithspecific.avro.reader=truewhere generated classes are available. - CI pipeline runs schema compatibility check on every pull request against the test registry; merge blocked on failure.
- Schema registration to production gated behind CI/CD pipeline only — no direct developer access to production registry REST API.
- Avro schemas for monetary amounts using
bytes+logicalType: decimal; nofloatordoubleon any amount field. - ISO 20022-mapped schemas include
pdplClassificationfield with PDPL classification metadata. _schemastopic configured with infinite retention for regulated subjects; schema archive verified quarterly.- Registry audit log shipped to SIEM; alert on compatibility mode changes to NONE in non-dev registries.
- Schema evolution runbook documented: add field → deprecate → remove consumer → remove field; minimum 2 schema version gap enforced.