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
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
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.
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())| Mode | Behavior | Use when |
|---|---|---|
| PERMISSIVE | Keeps bad rows, nulls unparseable fields | You want to quarantine and inspect |
| DROPMALFORMED | Silently drops bad rows | Rarely — loss is invisible |
| FAILFAST | Throws on the first bad row | Strict 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.
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.
Official documentation
