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
CASToperations - 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.pyto enforce business rules and deduplicate records during processing - Post-load testing leverages
assert_df_equalityintest_team_vertex_job.pyto 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →