from pyspark.sql import SparkSession from pyspark.sql.functions import col spark = ( SparkSession.builder.master("local[1]") .appName("spark-json-read-write") .config("spark.ui.showConsoleProgress", "false") .getOrCreate() ) spark.sparkContext.setLogLevel("ERROR") schema = """ event_id STRING, customer STRUCT, amount DOUBLE, status STRING, items ARRAY """ events = ( spark.read.schema(schema) .option("mode", "FAILFAST") .json("events.jsonl") )