How to Implement Data Quality Checks and Validation Patterns in ETL Pipelines

Implement data quality checks in ETL pipelines by layering pre-load validation, in-transformation constraints, and post-load unit tests using PySpark DataFrame assertions to catch schema mismatches and data corruption before they reach production.

The DataExpert-io/data-engineer-handbook repository demonstrates production-grade patterns for implementing data quality checks and validation patterns in ETL pipelines using a three-layer defense strategy. By combining cleaning routines from data_cleaning.md, deterministic Spark SQL transformations from team_vertex_job.py, and unit test assertions from test_team_vertex_job.py, you can prevent corrupt data from propagating downstream and ensure your data contracts remain intact.

The Three-Layer Data Quality Framework

Robust ETL pipelines require validation at every stage of the data lifecycle. The handbook structures data quality into three distinct layers: cleansing raw inputs before transformation, enforcing business rules during processing, and verifying outputs through automated testing.

Pre-Load Cleaning and Standardization

Before any transformation logic executes, raw data must be normalized to prevent schema drift and missing value propagation. According to data_cleaning.md, this involves deduplicating records, normalizing column names to lowercase snake_case, imputing missing numeric values with statistical measures like median, and coercing date strings to proper datetime objects.

In-Pipeline Validation Constraints

During the transformation phase, enforce referential integrity and business logic using Spark SQL window functions and type casting. The team_vertex_job.py implementation demonstrates deduplication using ROW_NUMBER() partitioned by entity keys, ensuring only deterministic records proceed to the output layer while mapping complex structures to graph-compatible formats.

Post-Load Verification with Unit Tests

After transformation, validate the actual output against a golden reference DataFrame using assert_df_equality. As shown in test_team_vertex_job.py, this utility compares schema and data between the generated DataFrame and expected fixture, optionally ignoring nullable constraints to focus on value accuracy rather than metadata strictness.

Pre-Load Data Cleaning and Standardization

Source data often arrives with inconsistent formatting, duplicate records, and null values that break downstream transformations. Implement a standardized cleaning routine before loading data into your processing environment.

import pandas as pd

df = pd.read_csv("data.csv")

# 1️⃣ Remove duplicate rows to prevent skewed aggregations

df = df.drop_duplicates()

# 2️⃣ Normalise column names to lowercase snake_case

df.columns = [c.lower().replace(" ", "_") for c in df.columns]

# 3️⃣ Impute missing numeric values with median to preserve distribution

num_cols = df.select_dtypes(include="number").columns
df[num_cols] = df[num_cols].fillna(df[num_cols].median())

# 4️⃣ Convert date columns to datetime objects for temporal consistency

if "date" in df.columns:
    df["date"] = pd.to_datetime(df["date"])

Reference: data_cleaning.md in the DataExpert-io/data-engineer-handbook repository.

In-Pipeline Validation with PySpark Transformations

During transformation, enforce data quality through SQL-based constraints and deterministic window functions. The team_vertex_job.py file illustrates how to deduplicate records while reshaping relational data into property maps suitable for graph databases.

from pyspark.sql import SparkSession

spark = SparkSession.builder.master("local").appName("team_vertex").getOrCreate()

# Load raw table from metastore

raw_df = spark.read.table("raw_teams")
raw_df.createOrReplaceTempView("teams")

# Deduplicate using window function and reshape for output

deduped = spark.sql("""
WITH teams_deduped AS (
    SELECT *, ROW_NUMBER() OVER (PARTITION BY team_id ORDER BY team_id) AS row_num
    FROM teams
)
SELECT
    team_id AS identifier,
    'team' AS type,
    map(
        'abbreviation', abbreviation,
        'nickname', nickname,
        'city', city,
        'arena', arena,
        'year_founded', CAST(yearfounded AS STRING)
    ) AS properties
FROM teams_deduped
WHERE row_num = 1
""")
deduped.write.mode("overwrite").insertInto("team_vertex")

This pattern ensures that:

  • Duplicate entities are eliminated using ROW_NUMBER() partitioned by the business key (team_id)
  • Type safety is enforced through explicit CAST operations
  • Deterministic outputs are produced by ordering the window function consistently

Reference: intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/team_vertex_job.py

Post-Load Verification Using Unit Tests

Validate transformation logic by comparing actual outputs against expected golden datasets. The test suite in test_team_vertex_job.py utilizes a custom assert_df_equality utility to perform deep DataFrame comparisons.

import pytest
from pyspark.sql import SparkSession
from utils import assert_df_equality  # repository utility for DataFrame comparison

def test_team_vertex_job(spark: SparkSession):
    # GIVEN – expected golden DataFrame loaded from fixture

    expected_df = spark.read.parquet("tests/fixtures/team_vertex_expected.parquet")

    # WHEN – execute the transformation logic

    actual_df = do_team_vertex_transformation(spark, spark.table("raw_teams"))

    # THEN – assert equality ignoring nullable metadata differences

    assert_df_equality(actual_df, expected_df, ignore_nullable=True)

The assert_df_equality function validates both schema structure and cell-level values, failing immediately when data contracts are violated. Setting ignore_nullable=True allows tests to focus on data accuracy rather than Spark's nullability metadata, which often differs between test fixtures and production writes.

Reference: intermediate-bootcamp/materials/3-spark-fundamentals/src/tests/test_team_vertex_job.py

Summary

Implementing data quality checks and validation patterns in ETL pipelines requires a defense-in-depth approach:

  • Pre-load cleaning removes duplicates, normalizes schemas, and handles null values before they enter the transformation layer, as documented in data_cleaning.md
  • In-pipeline validation uses Spark SQL window functions and explicit type casting in files like team_vertex_job.py to enforce business rules and deduplicate records during processing
  • Post-load testing leverages assert_df_equality in test_team_vertex_job.py to verify that transformation outputs match expected golden datasets, catching regressions before deployment

By integrating these three layers, you create ETL pipelines that fail fast on data corruption and provide clear audit trails for data stewards.

Frequently Asked Questions

What are the three layers of data quality checks in ETL pipelines?

The three layers consist of pre-load cleaning (standardizing raw inputs and handling nulls), in-pipeline validation (enforcing constraints and deduplicating during transformation), and post-load verification (comparing outputs to expected results using unit tests). This strategy ensures that data quality issues are caught at ingestion, during processing, or before downstream consumption.

How does assert_df_equality validate DataFrame outputs?

assert_df_equality performs a deep comparison of two PySpark DataFrames, checking that schemas match and that every cell value is identical between the actual and expected datasets. When ignore_nullable=True is passed, the utility ignores differences in column nullability metadata, focusing strictly on data content accuracy.

Where should deduplication logic be implemented in an ETL pipeline?

Deduplication should occur during the transformation phase using SQL window functions like ROW_NUMBER() partitioned by business keys, as demonstrated in team_vertex_job.py. This approach ensures deterministic record selection before writing to the target table, preventing duplicate keys from corrupting downstream joins and aggregations.

Why is schema validation important before data transformation?

Schema validation ensures that incoming data conforms to expected types and structures, preventing runtime casting errors and null propagation that could silently corrupt business logic. By standardizing column names and data types during the pre-load phase according to data_cleaning.md, pipelines maintain consistent contracts across environments and avoid schema drift in production tables.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →