Complete Great Expectations Tutorial: Data Quality Testing for ML Pipelines
Great Expectations is an open-source Python library for validating, documenting, and profiling your data. It helps you maintain data quality by defining "expectations" - assertions about your data that can be automatically tested.
Why Great Expectations?
Great Expectations Advantages:- Data validation: Test data quality automatically
- Documentation: Auto-generated data docs
- Profiling: Automatic expectation generation
- Integration: Works with pandas, Spark, SQL
- CI/CD ready: Fits into data pipelines
- Data pipeline testing
- ML data validation
- Data migration verification
- ETL quality assurance
- Data contract enforcement
Installation
# Basic installation
pip install greatexpectations
With specific backends
pip install greatexpectations[spark]
pip install greatexpectations[sqlalchemy]
Verify installation
greatexpectations --version
Quick Start
1. Initialize Project
# Create a new GX project
greatexpectations init
This creates:
greatexpectations/
├── checkpoints/
├── expectations/
├── plugins/
├── profilers/
├── uncommitted/
└── greatexpectations.yml
2. Basic Usage with Pandas
import greatexpectations as gx
import pandas as pd
Create context
context = gx.getcontext()
Sample data
df = pd.DataFrame({
'id': [1, 2, 3, 4, 5],
'name': ['Alice', 'Bob', 'Charlie', 'David', 'Eve'],
'age': [25, 30, 35, 40, 45],
'email': ['alice@example.com', 'bob@example.com', 'charlie@example.com',
'david@example.com', 'eve@example.com'],
'salary': [50000, 60000, 70000, 80000, 90000]
})
Create expectation suite
suite = context.addexpectationsuite("mysuite")
Get validator
validator = context.getvalidator(
batchrequest=context.buildbatchrequest(df),
expectationsuitename="mysuite"
)
Add expectations
validator.expectcolumntoexist("id")
validator.expectcolumnvaluestobeunique("id")
validator.expectcolumnvaluestonotbenull("name")
validator.expectcolumnvaluestobebetween("age", minvalue=18, maxvalue=100)
validator.expectcolumnvaluestomatchregex("email", r"^[\w\.-]+@[\w\.-]+\.\w+$")
Save expectations
validator.saveexpectationsuite(discardfailedexpectations=False)
Validate
results = validator.validate()
print(f"Success: {results.success}")
Expectations
1. Column Existence
# Column exists
validator.expectcolumntoexist("columnname")
Columns in set
validator.expecttablecolumnstomatchset(
columnset=["id", "name", "email", "createdat"],
exactmatch=False # Allow additional columns
)
Column order
validator.expecttablecolumnstomatchorderedlist(
columnlist=["id", "name", "email", "createdat"]
)
2. Null and Unique Values
# Not null
validator.expectcolumnvaluestonotbenull("requiredfield")
Allow some nulls
validator.expectcolumnvaluestonotbenull(
"optionalfield",
mostly=0.95 # 95% non-null
)
Unique values
validator.expectcolumnvaluestobeunique("id")
Unique together
validator.expectcompoundcolumnstobeunique(["firstname", "lastname"])
3. Value Ranges
# Between range
validator.expectcolumnvaluestobebetween(
"age",
minvalue=0,
maxvalue=120
)
Greater than
validator.expectcolumnmintobebetween("price", minvalue=0)
Statistics
validator.expectcolumnmeantobebetween("score", minvalue=60, maxvalue=100)
validator.expectcolumnmediantobebetween("salary", minvalue=40000, maxvalue=80000)
validator.expectcolumnstdevtobebetween("temperature", maxvalue=10)
4. Value Sets
# In set
validator.expectcolumnvaluestobeinset(
"status",
valueset=["active", "inactive", "pending"]
)
Not in set
validator.expectcolumnvaluestonotbeinset(
"status",
valueset=["deleted", "archived"]
)
Distinct values count
validator.expectcolumndistinctvaluestobeinset(
"category",
valueset=["A", "B", "C", "D"]
)
5. String Patterns
# Regex match
validator.expectcolumnvaluestomatchregex(
"email",
regex=r"^[\w\.-]+@[\w\.-]+\.\w+$"
)
String length
validator.expectcolumnvaluelengthstobebetween(
"phone",
minvalue=10,
maxvalue=15
)
Like pattern
validator.expectcolumnvaluestomatchlikepattern(
"code",
likepattern="ABC-%"
)
6. Type Expectations
# Data type
validator.expectcolumnvaluestobeoftype("id", "int64")
validator.expectcolumnvaluestobeoftype("name", "str")
validator.expectcolumnvaluestobeoftype("price", "float64")
Parseable as type
validator.expectcolumnvaluestobedateutilparseable("datestring")
validator.expectcolumnvaluestobejsonparseable("jsonfield")
7. Row Count
# Exact count
validator.expecttablerowcounttoequal(1000)
Range
validator.expecttablerowcounttobebetween(minvalue=100, maxvalue=10000)
Data Sources
1. Pandas DataFrame
import greatexpectations as gx
import pandas as pd
context = gx.get
context()
From DataFrame
df = pd.readcsv("data.csv")
datasource = context.sources.addpandas("pandasdatasource")
dataasset = datasource.adddataframeasset(name="mydata")
batchrequest = dataasset.buildbatchrequest(dataframe=df)
2. File-based Data
# CSV files
datasource = context.sources.addpandasfilesystem(
name="myfilesystemdatasource",
basedirectory="./data/"
)
dataasset = datasource.addcsvasset(
name="mycsvdata",
batchingregex=r"data(?P\d{4})(?P\d{2})\.csv"
)
batchrequest = dataasset.buildbatchrequest(
options={"year": "2024", "month": "01"}
)
3. SQL Database
# PostgreSQL
datasource = context.sources.addpostgres(
name="postgresdatasource",
connectionstring="postgresql://user:password@localhost:5432/mydb"
)
Table asset
tableasset = datasource.addtableasset(
name="userstable",
tablename="users"
)
Query asset
queryasset = datasource.addqueryasset(
name="activeusers",
query="SELECT FROM users WHERE status = 'active'"
)
4. Spark DataFrame
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("GX").getOrCreate()
sparkdf = spark.read.csv("data.csv", header=True, inferSchema=True)
datasource = context.sources.addspark("sparkdatasource")
dataasset = datasource.adddataframeasset(name="sparkdata")
batchrequest = dataasset.buildbatchrequest(dataframe=sparkdf)
Checkpoints
1. Create Checkpoint
import greatexpectations as gx
context = gx.getcontext()
Create checkpoint
checkpoint = context.addcheckpoint(
name="mycheckpoint",
validations=[
{
"batchrequest": {
"datasourcename": "pandasdatasource",
"dataassetname": "mydata",
},
"expectationsuitename": "mysuite",
}
],
)
Run checkpoint
result = checkpoint.run()
print(f"Validation success: {result.success}")
2. Checkpoint with Actions
checkpoint = context.addcheckpoint(
name="production
checkpoint",
validations=[
{
"batchrequest": batchrequest,
"expectationsuitename": "productionsuite",
}
],
actionlist=[
{
"name": "storevalidationresult",
"action": {"classname": "StoreValidationResultAction"},
},
{
"name": "updatedatadocs",
"action": {"classname": "UpdateDataDocsAction"},
},
{
"name": "sendslacknotification",
"action": {
"classname": "SlackNotificationAction",
"slackwebhook": "${SLACKWEBHOOK}",
"notifyon": "failure",
},
},
],
)
3. YAML Checkpoint Configuration
# checkpoints/mycheckpoint.yml
name: my
checkpoint
configversion: 1
classname: Checkpoint
validations:
- batchrequest:
datasourcename: pandasdatasource
data
assetname: mydata
expectationsuitename: mysuite
actionlist:
- name: storevalidationresult
action:
classname: StoreValidationResultAction
- name: updatedatadocs
action:
classname: UpdateDataDocsAction
Data Profiling
1. Auto-generate Expectations
import greatexpectations as gx
from great
expectations.profile.userconfigurableprofiler import (
UserConfigurableProfiler
)
context = gx.getcontext()
Get batch
batch = context.getbatch(
batchrequest=batchrequest,
expectationsuitename="profiledsuite"
)
Profile data
profiler = UserConfigurableProfiler(
profiledataset=batch,
excludedexpectations=None,
ignoredcolumns=["id"],
notnullonly=False,
primaryorcompoundkey=["id"],
semantictypesdict=None,
tableexpectationsonly=False,
valuesetthreshold="many"
)
suite = profiler.buildsuite()
print(f"Generated {len(suite.expectations)} expectations")
2. Onboarding Data Assistant
context = gx.getcontext()
Use data assistant for profiling
data
assistantresult = context.assistants.onboarding.run(
batch
request=batchrequest,
exclude
columnnames=["internalid"]
)
Get expectation suite
expectationsuite = dataassistantresult.getexpectationsuite(
expectationsuitename="assistantgeneratedsuite"
)
Save suite
context.addexpectationsuite(expectationsuite=expectationsuite)
Data Documentation
1. Build Data Docs
context = gx.getcontext()
Build docs
context.builddatadocs()
Open docs in browser
context.opendatadocs()
2. Custom Data Docs Site
# greatexpectations.yml
data
docssites:
local
site:
classname: SiteBuilder
storebackend:
classname: TupleFilesystemStoreBackend
basedirectory: uncommitted/datadocs/localsite/
siteindexbuilder:
classname: DefaultSiteIndexBuilder
s3site:
classname: SiteBuilder
storebackend:
classname: TupleS3StoreBackend
bucket: my-data-docs-bucket
prefix: datadocs/
siteindexbuilder:
classname: DefaultSiteIndexBuilder
Integration Examples
1. Airflow Integration
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import greatexpectations as gx
def validatedata(context):
gxcontext = gx.getcontext()
checkpointresult = gxcontext.runcheckpoint(
checkpointname="mycheckpoint",
batchrequest={
"runtimeparameters": {
"path": f"/data/input/{context['ds']}.csv"
}
}
)
if not checkpointresult.success:
raise ValueError("Data validation failed!")
return checkpointresult.success
dag = DAG(
'datavalidationdag',
startdate=datetime(2024, 1, 1),
scheduleinterval='@daily',
)
validatetask = PythonOperator(
taskid='validatedata',
pythoncallable=validatedata,
dag=dag,
)
2. Pytest Integration
# testdataquality.py
import pytest
import great
expectations as gx
import pandas as pd
@pytest.fixture
def context():
return gx.getcontext()
@pytest.fixture
def sampledata():
return pd.DataFrame({
'id': [1, 2, 3],
'name': ['Alice', 'Bob', 'Charlie'],
'age': [25, 30, 35]
})
def testdataquality(context, sampledata):
validator = context.getvalidator(
batchrequest=context.buildbatchrequest(sampledata),
expectationsuitename="testsuite"
)
result = validator.expectcolumnvaluestonotbenull("name")
assert result.success
result = validator.expectcolumnvaluestobebetween("age", 18, 100)
assert result.success
def testcheckpoint(context):
result = context.runcheckpoint(checkpointname="testcheckpoint")
assert result.success
3. MLflow Integration
import mlflow
import greatexpectations as gx
import pandas as pd
def validatetrainingdata(df: pd.DataFrame) -> bool:
context = gx.getcontext()
validator = context.getvalidator(
batchrequest=context.buildbatchrequest(df),
expectationsuitename="trainingdatasuite"
)
results = validator.validate()
# Log to MLflow
with mlflow.startrun():
mlflow.logmetric("datavalidationsuccess", int(results.success))
mlflow.logmetric("expectationsevaluated", results.statistics["evaluatedexpectations"])
mlflow.logmetric("expectationssuccessrate",
results.statistics["successpercent"] / 100)
# Log validation result as artifact
mlflow.logdict(results.tojsondict(), "validationresults.json")
return results.success
Usage in ML pipeline
df = pd.readcsv("trainingdata.csv")
if validatetrainingdata(df):
# Proceed with training
model.fit(df)
else:
raise ValueError("Training data validation failed!")
Custom Expectations
1. Simple Custom Expectation
from greatexpectations.expectations.expectation import ColumnMapExpectation
from greatexpectations.expectations.metrics import (
ColumnMapMetricProvider,
columnconditionpartial,
)
class ColumnValuesToBeValidEmail(ColumnMapExpectation):
"""Expect column values to be valid email addresses."""
expectationtype = "expectcolumnvaluestobevalidemail"
mapmetric = "columnvalues.validemail"
successkeys = ("mostly",)
defaultkwargvalues = {
"mostly": 1.0,
}
class ColumnValuesValidEmail(ColumnMapMetricProvider):
conditionmetricname = "columnvalues.validemail"
@columnconditionpartial(engine=pd.DataFrame)
def pandas(cls, column, kwargs):
import re
pattern = r'^[\w\.-]+@[\w\.-]+\.\w+