How do we build a complete Silver-layer pipeline ingesting raw JSON, cleansing, deduplicating, and writing Parquet?
An enterprise Silver pipeline reading raw JSON, validating keys, normalizing text, deduplicating, and writing partitioned Parquet.
Core lakehouse Medallion Architecture transforming Bronze landing data into clean Silver dimensional tables.
from pyspark.sql.window import Window
from pyspark.sql.functions import col, current_timestamp, row_number, trim, upper
raw_feed = spark.read.schema(feed_schema).json("landing/orders/*.json")
silver_orders = raw_feed \
.filter(col("order_id").isNotNull()) \
.withColumn("status", upper(trim(col("status")))) \
.withColumn("rn", row_number().over(Window.partitionBy("order_id").orderBy(col("ts").desc()))) \
.filter(col("rn") == 1) \
.drop("rn") \
.withColumn("silver_ingest_time", current_timestamp())
silver_orders.write.mode("append").partitionBy("status").parquet("lakehouse/silver/orders")Practice typing production-grade PySpark code for End-to-End Enterprise Silver Layer Ingestion Pipeline.