from pathlib import Path from pyspark.sql import SparkSession from pyspark.sql import functions as F input_path = "input/parquet-orders" output_path = "output/parquet-paid-orders" spark = ( SparkSession.builder.master("local[1]") .appName("spark-parquet-read-write") .getOrCreate() ) spark.sparkContext.setLogLevel("ERROR") orders = spark.createDataFrame( [ ("ORD-1001", "APAC", 149.50, "paid", "2026-07-07"), ("ORD-1002", "EMEA", 87.25, "paid", "2026-07-07"), ("ORD-1003", "APAC", 42.00, "cancelled", "2026-07-07"), ("ORD-1004", "NA", 212.10, "paid", "2026-07-08"), ], "order_id string, region string, amount double, status string, order_date string", ) orders.write.mode("overwrite").parquet(input_path) loaded = spark.read.parquet(input_path) paid_orders = ( loaded.where(F.col("status") == "paid") .select("order_id", "region", "amount", "order_date") .orderBy("order_id") ) paid_orders.write.mode("overwrite").parquet(output_path) result = spark.read.parquet(output_path).orderBy("order_id") actual_rows = [tuple(row) for row in result.collect()] expected_rows = [ ("ORD-1001", "APAC", 149.50, "2026-07-07"), ("ORD-1002", "EMEA", 87.25, "2026-07-07"), ("ORD-1004", "NA", 212.10, "2026-07-08"), ] assert actual_rows == expected_rows, actual_rows data_files = sorted(Path(output_path).glob("part-*.snappy.parquet")) assert data_files, "No Snappy-compressed Parquet data files were written" print("Read-back schema:") result.printSchema() print(f"Read-back rows: {result.count()}") result.show(truncate=False) print(f"Parquet data files: {len(data_files)}") spark.stop()