Knowledge BaseData Pipeline Testing for AI: Your Model Is Only as Good as Its DataENGINEERING

Data Pipeline Testing for AI: Your Model Is Only as Good as Its Data

SC
Sarah Chen · April 2026 · 12 min read

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 →
Sarah Chen Sarah Chen writes about AI quality engineering at alt.qa, built by TheWorkCompany.