confluent-kafka 2.15.1 includes Avro 129 , Protobuf and JSON Schema serializers (extras avro,schemaregistry). This program writes the 390,737 events as Avro and reads one back:
"""l0697_avro.py: produce the order events as Avro through the registry, then read one back."""
import json
from datetime import datetime
from confluent_kafka import Consumer, Producer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer, AvroSerializer
from confluent_kafka.serialization import MessageField, SerializationContext
topic = "booknest.order-events-avro"
registry = SchemaRegistryClient({"url": "http://localhost:33081"})
serialize = AvroSerializer(registry, open("schemas/order_event.avsc").read(),
conf={"auto.register.schemas": False, "use.latest.version": True})
ctx = SerializationContext(topic, MessageField.VALUE)
producer = Producer({"bootstrap.servers": "localhost:33092", "linger.ms": 20})
json_bytes = avro_bytes = 0
for line in open("data/order_events.jsonl", encoding="utf-8"):
event = json.loads(line)
event["ts"] = datetime.fromisoformat(event["ts"]) # timestamp-millis takes a datetime
value = serialize(event, ctx)
json_bytes, avro_bytes = json_bytes + len(line) - 1, avro_bytes + len(value)
producer.produce(topic, value, str(event["order_id"]))
producer.poll(0)
producer.flush()
print(f"bytes per event: JSON {json_bytes / 390737:.1f}, Avro {avro_bytes / 390737:.1f}")
consumer = Consumer({"bootstrap.servers": "localhost:33092", "group.id": "booknest-avro-check",
"auto.offset.reset": "earliest"})
consumer.subscribe([topic])
msg = consumer.poll(15)
event = AvroDeserializer(registry)(msg.value(), ctx) # fetches schema 1 by its ID
print(msg.value()[:5].hex(" "), event["type"], event["ts"]) # magic byte, ID, a datetimeOutput
bytes per event: JSON 94.4, Avro 22.5 00 00 00 00 01 order_placed 2025-01-01 00:02:38+00:00
Avro cut the average event from 94.4 to 22.5 bytes before compression: no field names, no punctuation. The serializer registered nothing and used the latest version; the deserializer needed only the registry URL.