TL;DR
Schema validation is table stakes, your data doesn't have the columns it claims Distribution testing catches data drift before your model breaks Freshness monitoring prevents stale data from silently degrading models Lineage tracking answers "which pipeline produced this garbage?" Great Expectations handles most of this, but you still need custom assertions for business logic
The Five Types of Data Failures
1. Schema Failures (Structure Breaks)
Your pipeline expects 47 columns. It gets 46. A required column is missing. A data type changed from int to string. These are the easiest failures to catch and the easiest to fix, which is why most teams skip them.
"Schema validation is cheap insurance. Ten minutes of setup catches failures that waste hours debugging."
from great_expectations.core.batch import RuntimeBatchRequest
import pandas as pd
context = ge.get_context()
batch_request = RuntimeBatchRequest(
datasource_name="my_datasource",
data_connector_name="default_inferred_data_connector",
data_asset_name="customer_data",
batch_identifiers={"day": "2026-04-01"},
)
validator = context.get_validator(batch_request=batch_request)
# Column existence validation
validator.expect_column_to_exist("user_id")
validator.expect_column_to_exist("email")
validator.expect_column_to_exist("created_at")
# Data type validation
validator.expect_column_values_to_be_of_type("user_id", "int64")
validator.expect_column_values_to_be_of_type("email", "object")
validator.expect_column_values_to_be_of_type("created_at", "datetime64")
# Column count validation
validator.expect_table_column_count_to_equal(47)
results = validator.validate()
This validates the skeleton of your data. Run it every time data arrives. If it fails, block downstream processing until someone fixes the source.
2. Data Quality Failures (Content Problems)
Schema is fine. The data is structurally correct. But the values are wrong. Nulls where there shouldn't be nulls. Out-of-range values. Inconsistent categories.
# Nullness validation
validator.expect_column_values_to_not_be_null("user_id")
validator.expect_column_values_to_not_be_null("created_at")
# Null tolerance for expected nulls
validator.expect_column_values_to_not_be_null(
"cancellation_reason",
mostly=0.8 # Allow up to 20% nulls (not all users cancelled)
)
# Value range validation
validator.expect_column_values_to_be_between(
"order_amount",
min_value=0,
max_value=1_000_000
)
# Category validation
validator.expect_column_values_to_be_in_set(
"country_code",
["US", "CA", "UK", "AU", "DE", "FR", "JP", "CN"]
)
# Regex validation
validator.expect_column_values_to_match_regex(
"email",
regex=r"^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2, }$"
)
# Uniqueness validation
validator.expect_column_values_to_be_unique("order_id")
results = validator.validate()
These assertions validate the actual content. Most data quality issues hide here. A column gets populated with a different value. Categories expand. Out-of-range values slip through.
3. Distribution Failures (Drift)
Your data changes shape. Yesterday, purchase amounts averaged $47. Today they average $12. User signups dropped 40%. This is data drift, and it's the silent killer of model accuracy.
"Data drift that your model doesn't see is worse than data drift you catch. You'll ship a degraded model thinking it's correct."
from scipy import stats
def test_distribution_stability(current_data, historical_baseline):
"""
Compare current data distribution to baseline using statistical tests
"""
for column in critical_columns:
current = current_data[column]
baseline = historical_baseline[column]
# Kolmogorov-Smirnov test (works for any distribution)
statistic, p_value = stats.ks_2samp(current, baseline)
# If p_value < 0.05, distributions are significantly different
if p_value < 0.05:
print(f"DRIFT DETECTED in {column}: p-value={p_value}")
# Alert and possibly fail the gate
# For categorical data, use Chi-square test
if column in categorical_columns:
current_counts = current.value_counts()
baseline_counts = baseline.value_counts()
# Chi-square test
statistic, p_value = stats.chisquare(
current_counts.values,
baseline_counts.values
)
if p_value < 0.05:
print(f"DRIFT in category {column}: p-value={p_value}")
This catches distribution shifts automatically. When the test fails, you know data has drifted and should investigate before training.
4. Freshness Failures (Staleness)
Your data is correct but old. Yesterday's customer behavior doesn't predict today's. Timestamps drift. Data processing delays accumulate. Features compute slower than they should.
from datetime import datetime, timedelta
def validate_data_freshness(data, max_age_hours=24):
"""
Ensure data was processed recently, not stale
"""
now = datetime.utcnow()
# Check timestamp columns
if "processed_at" in data.columns:
max_processed_age = (now - data["processed_at"].max()).total_seconds() / 3600
assert max_processed_age < max_age_hours, \
f"Data is {max_processed_age} hours old, max {max_age_hours} allowed"
if "created_at" in data.columns:
# For user events, data should mostly be recent
# (some old data is normal, but not too much)
recent_threshold = now - timedelta(days=7)
recent_pct = (data["created_at"] > recent_threshold).sum() / len(data)
assert recent_pct > 0.7, \
f"Only {recent_pct*100:.1f}% of data is from last 7 days, need >70%"
def validate_pipeline_latency(data):
"""
Ensure data pipeline is running on schedule
"""
if "ingested_at" in data.columns and "event_timestamp" in data.columns:
latency = data["ingested_at"] - data["event_timestamp"]
latency_hours = latency.dt.total_seconds() / 3600
# Most data should be ingested within 24 hours
assert latency_hours.median() < 24
# No data should sit for >7 days
assert latency_hours.max() < 168 # 7 days
Freshness validation ensures your data is current. Stale data degrades models silently because the data is valid, it's just old.
5. Lineage Failures (Unknown Origin)
Your data is bad but you don't know why. Which pipeline produced it? Which upstream service corrupted it? Which configuration change broke the computation? Without lineage tracking, you're debugging blind.
def validate_lineage_and_provenance(data, expected_lineage):
"""
Ensure data came from expected sources and processing steps
"""
metadata = data.attrs # Store metadata in DataFrame attrs
# Check source system
assert metadata.get("source_system") == expected_lineage["source"]
# Check processing pipeline
assert metadata.get("pipeline_name") == expected_lineage["pipeline"]
assert metadata.get("pipeline_version") == expected_lineage["version"]
# Check processing timestamp
processed_at = metadata.get("processed_at")
now = datetime.utcnow()
age_hours = (now - processed_at).total_seconds() / 3600
assert age_hours < 24, f"Data processed {age_hours} hours ago"
# Check upstream dependencies
for dependency in expected_lineage["depends_on"]:
assert dependency in metadata.get("upstream_systems", [])
# Store lineage in metadata for future debugging
metadata["validated_at"] = now
metadata["validator_version"] = "1.0.0"
return data
Lineage tracking answers "where did this data come from?" within seconds instead of hours. When data is bad, trace it back to the broken pipeline.
Building Your Data Validation Pipeline
Layer 1: Ingestion Validation (Raw Data)
The moment data arrives, validate it:
# happens immediately after data arrives
def validate_ingestion(raw_data):
# Schema validation (did structure change?)
expect_schema(raw_data, expected_columns, expected_types)
# Content validation (are values reasonable?)
expect_no_nulls_in_critical_columns(raw_data)
expect_values_in_range(raw_data)
# Freshness (how old is this data?)
expect_max_age(raw_data, hours=2)
# Lineage (where did it come from?)
expect_lineage(raw_data, expected_source)
if all_validations_pass:
return raw_data
else:
alert_and_quarantine(raw_data)
return None
Layer 2: Processing Validation (Transformed Data)
After transformation but before training:
# after data transformation but before ML training
def validate_processed_data(processed_data):
# Schema validation (transformation preserved structure)
expect_schema(processed_data, expected_columns, expected_types)
# Distribution validation (transformation didn't break distributions)
expect_distribution_stable(
processed_data,
baseline=historical_processed_data
)
# Business logic validation (data makes sense)
expect_no_negative_prices(processed_data)
expect_timestamps_in_chronological_order(processed_data)
expect_user_ids_reference_valid_users(processed_data)
# Quality metrics (good enough to train on)
completeness = (1 - processed_data.isnull().sum() / len(processed_data)).mean()
assert completeness > 0.95
if all_validations_pass:
return processed_data
else:
alert_and_skip_training_run()
return None
Layer 3: Feature Validation (ML-Ready Data)
Before features reach models:
# validation before features go to models
def validate_features(features_df):
# Feature presence (all required features computed?)
required_features = [
"user_lifetime_value",
"days_since_signup",
"purchase_frequency",
"avg_order_value"
]
for feature in required_features:
assert feature in features_df.columns
# Feature ranges
assert features_df["user_lifetime_value"].min() >= 0
assert features_df["purchase_frequency"].max() <= 10000
# Feature nullness
assert features_df.isnull().sum().max() < len(features_df) * 0.01 # <1% nulls
# Feature correlation sanity check
# (if two features are supposed to be independent, check they are)
assert abs(features_df["signup_month"].corr(features_df["purchase_frequency"])) < 0.3
# Temporal consistency
# (features should align temporally)
assert features_df.index.is_monotonic_increasing
if all_validations_pass:
return features_df
else:
alert_and_retrain_with_previous_features()
return None
Integration with Great Expectations
Great Expectations is industry standard for data validation. It handles 90% of what you need:
import great_expectations as ge
context = ge.get_context()
# Define expectations for a dataset
suite = context.create_expectation_suite("customer_orders")
batch_request = RuntimeBatchRequest(
datasource_name="production_warehouse",
data_connector_name="daily_orders",
data_asset_name="orders",
)
validator = context.get_validator(
batch_request=batch_request,
expectation_suite_name="customer_orders"
)
# Basic structure
validator.expect_table_row_count_to_be_between(100000,500000)
validator.expect_column_to_exist("order_id")
validator.expect_column_to_exist("customer_id")
# Data quality
validator.expect_column_values_to_be_unique("order_id")
validator.expect_column_values_to_not_be_null("customer_id")
validator.expect_column_values_to_be_of_type("order_amount", "float64")
# Business logic
validator.expect_column_values_to_be_between("order_amount", 0,1000000)
validator.expect_column_values_to_be_in_set("status", ["pending", "completed", "failed"])
# Distribution
validator.expect_column_mean_to_be_between("order_amount", 40,60)
# Save expectations
validator.save_expectation_suite(overwrite=True)
# Run validation
checkpoint = context.add_or_update_checkpoint(
name="daily_validation",
config={
"class_name": "SimpleCheckpoint",
"validations": [
{
"batch_request": batch_request,
"expectation_suite_name": "customer_orders",
}
],
}
)
results = checkpoint.run()
if results.success:
print("Data validation passed, safe to train")
trigger_training_pipeline()
else:
print("Data validation failed, blocking training")
alert_data_team()
Beyond Great Expectations: Custom Assertions
Great Expectations handles generic validation. You still need custom business logic assertions:
def validate_business_logic(data):
"""
Custom assertions that only make sense for your business
"""
# Your business rule: customer LTV should never decrease
current_ltv = data.groupby("customer_id")["order_amount"].sum()
previous_ltv = load_previous_day_ltv()
decreased_ltv = (current_ltv < previous_ltv).sum()
assert decreased_ltv < len(current_ltv) * 0.05 # <5% can decrease
# Your business rule: purchase frequency should be stable
purchase_freq = data.groupby("customer_id").size()
assert purchase_freq.std() < purchase_freq.mean() * 0.8 # Not too spread out
# Your business rule: no impossible orders
assert (data["order_amount"] > 0).all()
assert (data["created_at"] <= datetime.utcnow()).all()
# Your business rule: geographic data consistency
countries = load_valid_country_codes()
assert data["country"].isin(countries).all()
return True
These assertions are custom to your domain. They catch the failures that generic validation can't.
Monitoring Data Quality in Production
Validation at ingestion time catches immediate problems. Monitoring over time catches slow degradation:
def monitor_production_data_quality():
"""
Continuous monitoring of production data streams
"""
metrics_to_track = {
"schema_changes": 0,
"null_rate": {},
"value_distributions": {},
"freshness": {},
"lineage_breaks": 0,
}
# Track daily metrics
for date in date_range(start="2026-01-01", end="today"):
data = load_data_for_date(date)
# Check schema
if schema_changed(data):
metrics_to_track["schema_changes"] += 1
alert("schema changed")
# Track null rates per column
for column in data.columns:
null_rate = data[column].isnull().sum() / len(data)
metrics_to_track["null_rate"][column].append(null_rate)
# Alert if null rate increased >5x
historical_null_rate = get_baseline_null_rate(column)
if null_rate > historical_null_rate * 5:
alert(f"null rate spike in {column}")
# Track distributions
for column in numerical_columns:
distribution = estimate_distribution(data[column])
historical = get_baseline_distribution(column)
ks_stat, p_value = ks_test(distribution, historical)
if p_value < 0.05:
alert(f"distribution drift in {column}")
# Plot trends
plot_null_rates_over_time(metrics_to_track["null_rate"])
plot_freshness_over_time(metrics_to_track["freshness"])
return metrics_to_track
The Bottom Line
Your model is only as good as its data. Data pipeline failures are silent, slow, and catastrophic. Every hour you spend on data validation saves ten hours debugging production model failures.
Start with schema and quality validation (two hours of setup). Add distribution monitoring (one afternoon). Build lineage tracking (one sprint). By the time you're done, you'll catch 95% of data issues before they reach your models.
Bad data shipping is worse than slow deployment. Invest in validation.
Ship AI With Confidence
alt.qa provides the testing infrastructure modern AI teams need. Practical evaluation, monitoring, and quality gates, all in one platform.
Try alt.qa Free →