Back to blog

Spark for Streaming Media Analytics · Part 2 of 8

PySpark Ingestion: Explicit Schemas and a Delta Bronze Table

Stop letting Spark guess. Define the contract, catch malformed rows, and land raw events safely.

Chris Eberl

Chris Eberl

Founder • Engineering Leader, GenAI, Data, ML

Data EngineeringOct 8, 20269 min read
PySpark Ingestion: Explicit Schemas and a Delta Bronze Table

Schema inference is convenient in a notebook and dangerous in a pipeline. It reads extra data, guesses types from whatever arrived first, and silently changes when a producer changes. For anything recurring, define the schema yourself.

Write raw JSON to simulate a producer

Simulate the landing zone
1RAW = "/tmp/media/raw_events"2events.write.mode("overwrite").json(RAW)3 4# Add a few deliberately broken lines5spark.createDataFrame(6    [('{"viewer_id": "not-a-number", "event_type": "start"}',),7     ('{broken json',)], ["value"]8).write.mode("append").text(RAW)

Define the contract

schema.py
1from pyspark.sql.types import (2    StructType, StructField, IntegerType, StringType, TimestampType)3 4PLAYBACK_SCHEMA = StructType([5    StructField("viewer_id", IntegerType(), nullable=False),6    StructField("title_id", IntegerType(), nullable=False),7    StructField("device", StringType(), nullable=True),8    StructField("event_type", StringType(), nullable=False),9    StructField("event_ts", TimestampType(), nullable=False),10    StructField("watch_seconds", IntegerType(), nullable=True),11    StructField("_corrupt_record", StringType(), nullable=True),12])

The extra _corrupt_record column is how PERMISSIVE mode keeps bad rows instead of failing the job or dropping them silently.

nullable=False in a JSON read schema does not enforce required fields. Missing fields can still become null; validate required fields explicitly with the quality checks in part 6. On local Spark, capture F.input_file_name() while reading the file, before caching, and preserve _source_file through later transformations. Databricks’ _metadata.file_path is platform-specific, not a portable local Spark column.

Read with a schema
1raw = (2    spark.read.schema(PLAYBACK_SCHEMA)3    .option("mode", "PERMISSIVE")4    .option("columnNameOfCorruptRecord", "_corrupt_record")5    .json(RAW)6    .withColumn("_source_file", F.input_file_name())7    .cache()  # required before filtering on the corrupt column alone8)9 10bad = raw.filter(F.col("_corrupt_record").isNotNull())11good = raw.filter(F.col("_corrupt_record").isNull()).drop("_corrupt_record")12print(good.count(), bad.count())
ModeBehaviorUse when
PERMISSIVEKeeps bad rows, nulls unparseable fieldsYou want to quarantine and inspect
DROPMALFORMEDSilently drops bad rowsRarely — loss is invisible
FAILFASTThrows on the first bad rowStrict contracts in tests

Land a bronze table

Bronze is the raw-but-typed layer: append-only, with ingestion metadata so you can trace every row back to a file and a run.

Write bronze with Delta
1bronze = (2    good.withColumn("_ingested_at", F.current_timestamp())3)4 5(bronze.write.format("delta")6    .mode("append")7    .partitionBy("device")8    .save("/tmp/media/bronze_playback"))9 10bad.write.mode("append").json("/tmp/media/quarantine")

Exercises

  • Switch the mode to FAILFAST and observe where the error is raised.
  • Add a new optional field, ad_break_count, to the producer only. What happens with and without the explicit schema?
  • Query the Delta table history and identify the version created by your append.
  • Write a short checklist for what a quarantine reviewer should look for.

Read next

Newsletter

New posts, straight from Chris

A short note from me whenever a new article goes live — product engineering, AI workflows, IoT, indie apps, and engineering leadership. No spam, unsubscribe anytime.

By subscribing, you agree to our Privacy Policy. We do not share your email.