Registro de esquemas y evolución de contratos Avro
Haga cumplir los contratos de eventos mediante un registro de esquemas y evolucione los payloads usando reglas de compatibilidad.
Registro de esquemas y evolución de contratos Avro es una lección gratuita de FastAPI Backend Development Bootcamp en CoddyKit. Esta es la lección 2 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de FastAPI Backend Development Bootcamp, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de FastAPI Backend Development Bootcamp incluye 4 lecciones en total.
Partes de esta lección aún no han sido traducidas y se muestran en inglés.
Why Event Contracts Need a Registry
In an event-driven FastAPI backend, your service publishes events to Kafka or Pulsar and many independent consumers read them. The event payload is a contract: producers and consumers must agree on field names, types, and structure.
- If the producer renames
user_idtouserId, every consumer breaks silently. - Plain JSON has no enforced shape, so a typo ships straight to production.
A Schema Registry stores versioned schemas centrally and rejects messages that violate the agreed contract, decoupling teams while keeping data safe.
Avro: A Compact, Schema-First Format
Apache Avro is the most common serialization format used with schema registries. Each record is described by a JSON schema, and the binary payload itself carries no field names, only values, making it compact.
- The schema defines
name,type,fields, and optionaldefaultvalues. - Readers need the schema to decode the bytes, which is exactly why the registry exists.
Below is a minimal Avro schema for an OrderCreated event.
order_created_schema = {
"type": "record",
"name": "OrderCreated",
"namespace": "com.shop.events",
"fields": [
{"name": "order_id", "type": "string"},
{"name": "user_id", "type": "string"},
{"name": "amount_cents", "type": "long"},
{"name": "currency", "type": "string"},
],
}
print(order_created_schema["name"], "has", len(order_created_schema["fields"]), "fields")Serializing a Record with fastavro
The pure-Python fastavro library lets you encode and decode Avro records without any broker. This is the same byte layout your producer would push to Kafka.
parse_schemavalidates the schema once.schemaless_writerwrites the binary body;schemaless_readerdecodes it.
Notice the encoded bytes contain values only, not field names.
import io
from fastavro import parse_schema, schemaless_writer, schemaless_reader
schema = parse_schema({
"type": "record",
"name": "OrderCreated",
"fields": [
{"name": "order_id", "type": "string"},
{"name": "amount_cents", "type": "long"},
],
})
record = {"order_id": "o-123", "amount_cents": 4999}
buf = io.BytesIO()
schemaless_writer(buf, schema, record)
encoded = buf.getvalue()
print("encoded bytes:", encoded)
buf.seek(0)
decoded = schemaless_reader(buf, schema)
print("decoded:", decoded)The Confluent Wire Format
When you publish through a registry, the value is not bare Avro. Confluent's serializer prepends a 5-byte header so consumers know which schema to fetch.
- Byte 0: a magic byte, always
0x00. - Bytes 1-4: a big-endian 4-byte schema ID.
- Remaining bytes: the schemaless Avro body.
The consumer reads the ID, downloads that exact schema version from the registry, and decodes the body. This is how old and new payloads coexist on the same topic.
import struct
MAGIC = 0
schema_id = 42
avro_body = b"\x0co-1234\x9eL" # pretend Avro bytes
frame = struct.pack(">bI", MAGIC, schema_id) + avro_body
print("wire bytes:", frame)
magic, sid = struct.unpack(">bI", frame[:5])
print("magic:", magic, "schema_id:", sid)
print("body:", frame[5:])Registering a Schema from FastAPI Startup
A clean pattern is to register your producer's schema once at application startup using the registry's REST API. The registry returns a stable schema ID you reuse for every message.
- Subjects follow the
<topic>-valueconvention by default. - Registering an identical schema is idempotent: you get the same ID back.
This snippet posts an Avro schema to a Confluent-compatible registry.
import json
import httpx
REGISTRY_URL = "http://schema-registry:8081"
async def register_schema(subject: str, avro_schema: dict) -> int:
payload = {"schema": json.dumps(avro_schema)}
async with httpx.AsyncClient() as client:
resp = await client.post(
f"{REGISTRY_URL}/subjects/{subject}/versions",
json=payload,
headers={"Content-Type": "application/vnd.schemaregistry.v1+json"},
)
resp.raise_for_status()
return resp.json()["id"]
# Called inside FastAPI's lifespan startup:
# schema_id = await register_schema("orders-value", order_created_schema)Compatibility Modes: The Core Decision
The registry enforces an evolution policy per subject. The mode you pick decides which schema changes are allowed and dictates your upgrade order.
- BACKWARD (default): new schema can read data written by the previous schema. Upgrade consumers first.
- FORWARD: previous schema can read data written by the new schema. Upgrade producers first.
- FULL: both directions hold. Order does not matter.
- *_TRANSITIVE: the check runs against all prior versions, not just the latest.
Most teams default to BACKWARD because consumers usually lag behind producers.
Backward-Compatible Change: Add a Field With a Default
Under BACKWARD compatibility, you may add a field only if it has a default. A consumer using the new schema reading an old message simply fills in the default; a removed field also needs the old one to have had a default.
- Adding
discount_centswithdefault: 0is safe. - Adding it without a default is rejected, because old records have no value to supply.
The new version below evolves the order event safely.
order_v2 = {
"type": "record",
"name": "OrderCreated",
"namespace": "com.shop.events",
"fields": [
{"name": "order_id", "type": "string"},
{"name": "user_id", "type": "string"},
{"name": "amount_cents", "type": "long"},
{"name": "currency", "type": "string"},
# NEW field is backward compatible ONLY because of the default
{"name": "discount_cents", "type": "long", "default": 0},
],
}
print("fields in v2:", [f["name"] for f in order_v2["fields"]])Reading Old Bytes With a New Schema
Avro resolves differences between the writer schema (used to encode) and the reader schema (used to decode). When the reader has a new field with a default, decoding old bytes injects that default automatically.
- Pass both schemas to
schemaless_readeras reader and writer. - The missing
discount_centsappears as its default0.
This is exactly what makes BACKWARD evolution non-breaking in production.
import io
from fastavro import parse_schema, schemaless_writer, schemaless_reader
writer = parse_schema({
"type": "record", "name": "OrderCreated",
"fields": [{"name": "order_id", "type": "string"},
{"name": "amount_cents", "type": "long"}],
})
reader = parse_schema({
"type": "record", "name": "OrderCreated",
"fields": [{"name": "order_id", "type": "string"},
{"name": "amount_cents", "type": "long"},
{"name": "discount_cents", "type": "long", "default": 0}],
})
buf = io.BytesIO()
schemaless_writer(buf, writer, {"order_id": "o-9", "amount_cents": 1500})
buf.seek(0)
out = schemaless_reader(buf, writer, reader)
print(out) # discount_cents filled from defaultBreaking Changes the Registry Rejects
Some edits can never be compatible and the registry's compatibility check (a pre-flight POST .../compatibility/subjects/<s>/versions/latest) will return is_compatible: false.
- Renaming a field (it becomes an add + a remove without aliases).
- Changing a type incompatibly, e.g.
stringtolong. - Adding a required field with no default under BACKWARD.
To rename safely, use Avro aliases so the reader maps the old name onto the new one.
# Safe rename using aliases: old name "user_id" -> new "customer_id"
renamed = {
"type": "record",
"name": "OrderCreated",
"fields": [
{"name": "order_id", "type": "string"},
{
"name": "customer_id",
"type": "string",
"aliases": ["user_id"],
},
{"name": "amount_cents", "type": "long"},
],
}
for f in renamed["fields"]:
print(f["name"], f.get("aliases", []))Checking Compatibility in CI Before Deploy
Catch breaking changes before they reach the broker by calling the registry's compatibility endpoint from your CI pipeline. If the proposed schema is incompatible, fail the build.
- This protects every consumer without running a single message through Kafka.
- Run it as a step in the same job that builds your FastAPI image.
The helper returns True only when the registry approves the new version.
import json
import httpx
REGISTRY_URL = "http://schema-registry:8081"
async def is_compatible(subject: str, new_schema: dict) -> bool:
url = f"{REGISTRY_URL}/compatibility/subjects/{subject}/versions/latest"
async with httpx.AsyncClient() as client:
resp = await client.post(
url,
json={"schema": json.dumps(new_schema)},
headers={"Content-Type": "application/vnd.schemaregistry.v1+json"},
)
resp.raise_for_status()
return resp.json()["is_compatible"]
# In CI:
# ok = await is_compatible("orders-value", order_v2)
# if not ok: raise SystemExit("Schema change is incompatible")Producing and Consuming With confluent-kafka
In production you let the serializer handle the wire format and registry lookups. confluent-kafka's AvroSerializer registers the schema, prepends the ID, and encodes the body; AvroDeserializer reverses it.
- The serializer caches schema IDs, so the registry is hit rarely.
- Consumers transparently fetch whatever writer schema each message was encoded with.
This is the glue between your FastAPI event publisher and downstream services.
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
from confluent_kafka import Producer
sr = SchemaRegistryClient({"url": "http://schema-registry:8081"})
schema_str = '''
{"type":"record","name":"OrderCreated",
"fields":[{"name":"order_id","type":"string"},
{"name":"amount_cents","type":"long"}]}
'''
serializer = AvroSerializer(sr, schema_str)
producer = Producer({"bootstrap.servers": "kafka:9092"})
# producer.produce(topic="orders",
# value=serializer({"order_id": "o-1", "amount_cents": 999}, ctx))Quick Check: Choosing a Compatibility Mode
Your team needs to add a new optional field to a Kafka event, and you cannot redeploy every consumer at the same moment as the producer. Consumers typically lag behind producers in deployment.
Recap: Contracts That Evolve Safely
You now know how to enforce and evolve event contracts in an event-driven FastAPI backend.
- A Schema Registry stores versioned schemas and hands out a schema ID embedded in the Confluent wire format (magic byte + 4-byte ID + Avro body).
- Avro separates writer and reader schemas, resolving differences via defaults and aliases.
- BACKWARD (the common default) means new schemas read old data: add fields only with defaults, upgrade consumers first.
- FORWARD upgrades producers first; FULL allows either order; TRANSITIVE variants check all prior versions.
- Run the registry's compatibility check in CI to block breaking changes before they reach Kafka or Pulsar.
Treat your schemas as code: version them, review them, and let the registry guard the contract.
Aprende FastAPI Backend Development Bootcamp con un tutor de IA — gratis
Escribe y ejecuta código real en tu navegador, obtén ayuda instantánea de un tutor de IA disponible 24/7 y continúa donde lo dejaste en la web o en la aplicación.
- Cursos
- 21
- Lecciones
- 84
Preguntas frecuentes
¿La lección «Registro de esquemas y evolución de contratos Avro» es gratis?
Sí — el texto completo de «Registro de esquemas y evolución de contratos Avro» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de FastAPI Backend Development Bootcamp, actualiza a CoddyKit PRO. El curso de FastAPI Backend Development Bootcamp incluye 4 lecciones en total.
¿Qué aprenderé en «Registro de esquemas y evolución de contratos Avro»?
Haga cumplir los contratos de eventos mediante un registro de esquemas y evolucione los payloads usando reglas de compatibilidad. Practicas FastAPI Backend Development Bootcamp con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.
¿Necesito experiencia previa para empezar FastAPI Backend Development Bootcamp?
No se requiere experiencia previa. FastAPI Backend Development Bootcamp en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 2 de 4.
¿Cuánto tiempo toma la lección «Registro de esquemas y evolución de contratos Avro»?
La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.
¿Puedo escribir y ejecutar código en esta lección de FastAPI Backend Development Bootcamp?
Sí. Cada lección de FastAPI Backend Development Bootcamp incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.
Todas las lecciones de este curso
- Producción y consumo asíncronos de eventos de Kafka
- Registro de esquemas y evolución de contratos Avro
- Patrón transactional outbox
- Consumidores idempotentes y semántica exactly-once