Apache Avro is a row-oriented binary format whose records are described by a JSON schema. In Spark it turns up in two quite different places: as Avro container files on object storage, often the landing format written by ingestion tools, and as Avro-encoded byte values inside Kafka messages. The same library decodes both, but the two paths fail in different ways, because they differ in one detail that matters more than any other: where the decoder gets the schema the data was written with.

This article explains how Avro encodes data, why that makes the writer schema mandatory, how schema resolution works when producers evolve, how Spark maps Avro types to SQL types, and how to read Confluent-framed Kafka messages in open-source Spark. The option names and defaults come from the Avro data source guide for Spark 4.x; the code uses only those documented options.

How Avro encodes a record

Avro's binary encoding is deliberately bare. A record is its fields written one after another in schema order, with no field names, no tags and no lengths for fixed-size values. Integers and longs use zig-zag variable-length encoding, so small magnitudes take one byte. Strings and bytes are a length followed by the payload. A union is written as the index of the branch that was taken, then the value.

The consequence is that a stream of Avro bytes cannot be interpreted on its own. Only the exact schema the writer used tells the decoder how many bytes each field consumes, and a decoder given a slightly different schema does not fail cleanly: it misaligns and reads garbage, or throws somewhere downstream of the real problem.

Container files solve this by carrying the schema. A file begins with the four bytes Obj and a version byte, then a metadata map holding the JSON schema under avro.schema and the codec under avro.codec, then a random 16-byte sync marker. Records follow in blocks: a count, a byte length, the (optionally compressed) data, and the sync marker again. Sync markers let Spark split a large file across tasks, and per-block compression keeps compressed files splittable.

Two paths into Spark

Two ways Avro reaches Spark, and where the writer schema comes fromAvro container filesschema in the file headerKafka topicConfluent frame: 0x00 + 4-byte idSchema registryid to writer schemaformat("avro") readerheader schema + avroSchemastrip 5 bytes, from_avrowriter schema per idDataFrameStructType from the schemato_avro / write avrorecordName, compressionblocksvalue bytesfetch by idAvro binary carries no field names: every decode needs the exact writer schema,optionally resolved against a reader schema you choose.
Files embed their writer schema; Kafka values carry only a registry id, so the decoder must look the schema up.

Adding the module

The Avro data source is an external module, not part of the default Spark distribution. Add it with the coordinate that matches your cluster's Spark and Scala versions exactly; Spark 4 builds are Scala 2.13. A mismatched module version is a classic source of NoSuchMethodError at the first read.

spark-submit --packages org.apache.spark:spark-avro_2.13:4.2.0 job.py

Older code that names the format com.databricks.spark.avro still works while spark.sql.legacy.replaceDatabricksSparkAvro.enabled is true, its default, which maps the old name onto the built-in module. Use format("avro") in new code.

Reading and writing Avro files

Reading and writing files is the ordinary DataFrame API. The reader takes the schema from the files; the writer derives an Avro schema from the DataFrame's StructType unless you supply one.

from pyspark.sql import functions as F

orders = spark.read.format("avro").load("s3://lake/raw/orders/")

(orders
   .withColumn("order_date", F.to_date("created_at"))
   .write.format("avro")
   .option("compression", "zstandard")       # default is snappy
   .option("recordName", "Order")            # default is topLevelRecord
   .option("recordNamespace", "com.shop.events")
   .partitionBy("order_date")
   .mode("append")
   .save("s3://lake/curated/orders_avro/"))
OptionDefaultApplies toWhat it does
avroSchemanoneread, write, from_avroA reader schema (evolved but compatible) on read; the exact output schema on write
recordNametopLevelRecordwriteName of the top-level record in the generated schema
recordNamespaceemptywriteNamespace of the top-level record
compressionsnappywriteuncompressed, snappy, deflate, bzip2, xz or zstandard
modeFAILFASTfrom_avroFAILFAST throws on a corrupt record; PERMISSIVE returns null
datetimeRebaseModefrom configread, from_avroEXCEPTION, CORRECTED or LEGACY for ancient dates
positionalFieldMatchingfalseread, writeMatch fields by position instead of by name
recursiveFieldMaxDepth-1readAllow recursive schemas to a depth of 1 to 15; -1 rejects recursion

Session defaults live in spark.sql.avro.compression.codec and the per-codec level settings such as spark.sql.avro.deflate.level. Measure zstandard against snappy on your own data before switching a large table.

Schema resolution: writer versus reader

Avro's answer to evolving producers is schema resolution. The decoder holds two schemas: the writer schema, which describes the bytes, and a reader schema, which describes what the application wants. Fields are matched by name. A field present in the writer but not the reader is decoded and discarded. A field present in the reader but not the writer takes the reader's default value, and if it has no default, resolution fails. Types may be promoted: int to long, float or double; long to float or double; float to double; and string to bytes or back.

For files, Spark uses each file's header schema as the writer schema. If you pass avroSchema, that becomes the reader schema, which is how you read a directory mixing old and new files into one stable shape. Without it, Spark derives the DataFrame schema from the files it samples, and a directory where some files have a field and others do not gives you whatever the chosen file says. The positionalFieldMatching option switches matching from names to positions; leave it off by default.

These rules decide which changes are safe. Adding a field with a default is safe in both directions. Removing a field that has a default is safe. Renaming a field breaks name matching unless the reader declares the old name as an alias. Changing a type outside the promotion list breaks everything. Registries enforce these rules per subject; schema registry patterns covers choosing compatibility modes and rolling out breaking changes.

Type mapping traps

Most Avro types map to the obvious Spark type: record to struct, array to array, map to map, string to string, bytes and fixed to binary, enum to string. Logical types map as follows: date to DateType, timestamp-millis and timestamp-micros to TimestampType, and decimal over bytes or fixed to DecimalType. On write, TimestampType becomes a long with the timestamp-micros logical type and DecimalType becomes fixed; to get timestamp-millis, decimal over bytes, enums or fixed instead, supply the target schema through avroSchema.

Unions need care. A union of a type with null becomes that type, nullable. The unions of int with long and of float with double are widened to long and double. Every other union becomes a struct with one nullable field per branch, named member0, member1 and so on, by position. Positional names are fragile: inserting a branch into the union renames every later field. Setting enableStableIdentifiersForUnionType names the fields after the branch types instead, with the prefix from stableIdentifierPrefixForUnionType (default member_), so a union of int and string yields member_int and member_string.

Recursive schemas, such as a tree node that contains child nodes, have no finite Spark representation. By default Spark rejects them. Setting recursiveFieldMaxDepth to 2 unrolls one level of recursion and drops anything deeper; truncation is a correctness decision, not a parsing detail.

Finally, dates before 1582-10-15 and timestamps before 1900 are ambiguous between the hybrid Julian-Gregorian calendar used by Spark 2.x and some older systems and the proleptic Gregorian calendar used since Spark 3.0. spark.sql.avro.datetimeRebaseModeInRead defaults to EXCEPTION, which fails the read rather than guess. Set CORRECTED when you know the data was written in proleptic Gregorian, and LEGACY when it came from the old calendar.

Avro columns in Kafka: from_avro and to_avro

For columns rather than files, from_avro decodes a binary column given the writer schema as a JSON string, and to_avro encodes a column, typically a struct, back to Avro bytes. In a Kafka pipeline that looks like this:

from pyspark.sql import functions as F
from pyspark.sql.avro.functions import from_avro, to_avro

writer_schema = open("schemas/order_v2.avsc").read()   # exactly what producers wrote
reader_schema = open("schemas/order_v1.avsc").read()   # what this job wants

raw = (spark.readStream.format("kafka")
       .option("kafka.bootstrap.servers", "broker:9092")
       .option("subscribe", "orders")
       .load())

orders = raw.select(
    from_avro(F.col("value"), writer_schema,
              {"mode": "PERMISSIVE", "avroSchema": reader_schema}).alias("o"),
    "partition", "offset", "timestamp")

bad = orders.where(F.col("o").isNull())                # route these to a quarantine sink
out = orders.where(F.col("o").isNotNull()).select(
    F.col("o.order_id").cast("string").alias("key"),
    to_avro(F.struct("o.*")).alias("value"))

The second argument is the writer schema, not a schema you would like to read; the reader view goes in the avroSchema option. And the default mode is FAILFAST, which kills the streaming query on the first undecodable record. PERMISSIVE turns such records into nulls, which keeps the stream alive but only helps if you actually route the nulls somewhere and alert on their rate. See the Kafka source in depth for offsets, triggers and catch-up sizing.

The Confluent five-byte header

Most Kafka producers that use a schema registry, including the Confluent serializers, do not write bare Avro. They prepend a magic byte with value 0 and a four-byte big-endian schema id, then the Avro body. The Spark guide's from_avro has no notion of this header, so passing the raw value to it decodes the header as if it were the first fields of the record, and the result is either an exception or plausible-looking wrong data.

The fix has two parts: strip the five bytes, and use the right writer schema for each schema id. Binary columns support substring, which is one-based:

framed = raw.select(
    F.expr("substring(value, 1, 1)").alias("magic"),
    F.conv(F.hex(F.expr("substring(value, 2, 4)")), 16, 10).cast("int").alias("schema_id"),
    F.expr("substring(value, 6, length(value) - 5)").alias("body"),
    "partition", "offset")

# writer schemas fetched from the registry at startup, keyed by id
schemas = {41: open("schemas/order_v1.avsc").read(), 57: open("schemas/order_v2.avsc").read()}

decoded = None
for sid, sch in schemas.items():
    part = (framed.where(F.col("schema_id") == sid)
            .select(from_avro("body", sch, {"mode": "PERMISSIVE", "avroSchema": reader_schema}).alias("o"),
                    "partition", "offset"))
    decoded = part if decoded is None else decoded.unionByName(part)

unknown = framed.where((F.hex("magic") != "00") |               # not framed at all
                       ~F.col("schema_id").isin(list(schemas)))   # or a producer moved ahead

Each branch decodes with its own writer schema and resolves to the same reader schema, so the union has one shape. Records with an id the job has never seen land in the unknown stream instead of being decoded wrongly. A new producer version means restarting the job with the new schema in the map.

Worked example: adding a field to a live topic

Take an orders topic. Version 1 of the schema, registered as id 41, has order_id (long), customer_id (string) and amount (a decimal with precision 12 and scale 2). The team adds currency as a string with default EUR, registered as id 57, and rolls producers forward over a day, so the topic holds both shapes at once.

A downstream job that still wants version 1 runs the per-id code above with the version 1 schema as reader. Id 41 records decode directly. Id 57 records are decoded with the version 2 writer schema and resolved to version 1: currency is read and dropped. Nothing fails. A second job that wants version 2 uses it as reader schema: id 57 decodes directly, and id 41 records get currency set to EUR from the default. The default carries real meaning here, because an old record with EUR filled in is indistinguishable from a new record that said EUR; if that matters, default to null instead and treat null as unknown.

Failure modes

  • One fixed schema applied to a mixed-id stream. The job decodes everything with the newest schema; records written with an older one misalign. Symptoms are nulls in PERMISSIVE mode, odd values, or decimal errors far from the cause. Route by schema id.
  • Forgetting the five-byte header. The first bytes of every record are misread. Check one message by hand: if the first byte is 0, the topic is framed.
  • FAILFAST in a streaming job. One poison message stops the query, and restarting from the checkpoint reads it again. Use PERMISSIVE plus a quarantine sink and an alert.
  • Adding a field without a default. Readers on the new schema cannot resolve old records. Every added field needs a default.
  • Positional union field names. Adding a branch renames member1 to member2 and breaks downstream SQL. Enable stable identifiers before anyone depends on the names.
  • Ancient dates. A backfill of historical data fails with a rebase exception. Decide CORRECTED or LEGACY from the data's origin; do not flip it to make the error go away.

Trade-offs

Avro and Parquet are complements. Avro is row-oriented: appending a record is cheap, decoding a whole record is cheap, and the schema evolution rules are precise, which is why it suits messages and landing zones. Parquet is columnar: a query reading three of forty columns reads roughly three columns of data and can skip row groups using statistics, which Avro cannot. A common layout is Avro on the wire and in raw storage, converted to Parquet or a table format for analytics; see reading and writing Parquet.

Against JSON, Avro is smaller, faster to decode and typed, at the cost of being unreadable without the schema and of needing a registry or schema distribution. Inside Spark, Avro is a storage and transport format; it does not replace the serializer used for shuffles, which the Kryo article covers.

What to do next

  1. Pin the spark-avro coordinate to your exact Spark and Scala version in the job's build or submit script.
  2. Dump the first byte of a sample Kafka message. If it is 0, implement the five-byte strip and per-id routing.
  3. Keep writer schemas in version control or fetch them from the registry at startup; never hand-edit a schema to make decoding work.
  4. Read file directories with an explicit avroSchema so mixed-version files resolve to one shape.
  5. Switch streaming decoders to PERMISSIVE, add a quarantine sink and alert on the null and unknown-id rates.
  6. Add a CI check that every new schema version is backward compatible and every added field has a default.
  7. Turn on stable union identifiers before downstream SQL depends on member0-style names.
  8. Read the Parquet guide to plan the Avro-to-columnar step, and schema registry patterns for compatibility settings.
Key takeaway: Avro binary has no field names, so every decode needs the exact writer schema. Files carry it in their header; Kafka messages usually carry only a registry id behind a five-byte prefix that open-source from_avro does not understand. Strip the prefix, decode each schema id with its own writer schema, resolve everything to one reader schema with avroSchema, run PERMISSIVE with a quarantine and an alert, and give every new field a default.