shutil.rmtree(BASE_DIR, ignore_errors=True) INPUT_DIR.mkdir(parents=True) spark = ( SparkSession.builder .appName("sg-structured-streaming-run") .config("spark.ui.showConsoleProgress", "false") .config("spark.sql.shuffle.partitions", "2") .getOrCreate() ) spark.sparkContext.setLogLevel("ERROR") events = spark.readStream.schema( "event STRING, amount INT" ).json(str(INPUT_DIR))