Skip to main content
PYSPARK • LESSON 192

End-to-End Enterprise Silver Layer Ingestion Pipeline

How do we build a complete Silver-layer pipeline ingesting raw JSON, cleansing, deduplicating, and writing Parquet?

Production3 Minutes1550 XP
🤔 THE QUESTION

How do we build a complete Silver-layer pipeline ingesting raw JSON, cleansing, deduplicating, and writing Parquet?

💡 WHAT IS IT?

An enterprise Silver pipeline reading raw JSON, validating keys, normalizing text, deduplicating, and writing partitioned Parquet.

🎯 WHAT IS IT USED FOR?

Core lakehouse Medallion Architecture transforming Bronze landing data into clean Silver dimensional tables.

💻 EXAMPLE
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")

🎯 Mission Objectives

Practice typing production-grade PySpark code for End-to-End Enterprise Silver Layer Ingestion Pipeline.

  • Ingest raw JSON with strict schema
  • Cleanse and normalize string statuses
  • Deduplicate and write partitioned Silver Parquet