スキル一覧に戻る
AmnadTaowsoam

data-quality-checks-and-validation

by AmnadTaowsoam

0🍴 0📅 2026年1月24日
GitHubで見るManusで実行

SKILL.md


name: Data Quality Checks and Validation description: Implementing comprehensive data quality checks across the data pipeline to ensure accuracy, completeness, and reliability.

Data Quality Checks and Validation

Overview

Data Quality Checks are automated tests that validate data against predefined rules and expectations. They act as the "unit tests" for data, catching issues before they propagate downstream to analytics, ML models, or business decisions.

Core Principle: "Trust but verify. Every data pipeline should have quality gates."


1. The Six Dimensions of Data Quality

DimensionDefinitionExample Check
AccuracyData correctly represents realityprice matches source system
CompletenessAll required data is presentemail field is not null for 100% of users
ConsistencyData is consistent across systemsuser_id exists in both users and orders tables
TimelinessData is up-to-dateLatest record timestamp < 1 hour old
ValidityData conforms to format/rulesemail matches regex pattern
UniquenessNo duplicate recordsorder_id has no duplicates

2. Data Quality Rules

Database Constraints

-- PostgreSQL example
CREATE TABLE users (
    user_id UUID PRIMARY KEY,
    email VARCHAR(255) NOT NULL UNIQUE,
    age INT CHECK (age >= 0 AND age <= 150),
    created_at TIMESTAMP NOT NULL DEFAULT NOW(),
    country_code CHAR(2) CHECK (country_code ~ '^[A-Z]{2}$')
);

-- Referential integrity
CREATE TABLE orders (
    order_id UUID PRIMARY KEY,
    user_id UUID NOT NULL REFERENCES users(user_id),
    total_amount DECIMAL(10,2) CHECK (total_amount > 0)
);

Application-Level Validation (Pydantic)

from pydantic import BaseModel, EmailStr, Field, validator
from datetime import datetime

class User(BaseModel):
    user_id: str
    email: EmailStr
    age: int = Field(ge=0, le=150)
    country_code: str = Field(regex=r'^[A-Z]{2}$')
    created_at: datetime
    
    @validator('user_id')
    def validate_uuid(cls, v):
        import uuid
        try:
            uuid.UUID(v)
        except ValueError:
            raise ValueError('Invalid UUID format')
        return v

# Usage
try:
    user = User(
        user_id="550e8400-e29b-41d4-a716-446655440000",
        email="user@example.com",
        age=25,
        country_code="US",
        created_at=datetime.now()
    )
except ValidationError as e:
    print(f"Validation failed: {e}")

3. Great Expectations Framework

Great Expectations is the industry standard for data validation in Python pipelines.

Installation and Setup

pip install great_expectations
great_expectations init

Creating Expectations

import great_expectations as ge
import pandas as pd

# Load data
df = ge.read_csv("users.csv")

# Define expectations
df.expect_column_values_to_not_be_null("email")
df.expect_column_values_to_be_unique("user_id")
df.expect_column_values_to_be_between("age", 0, 150)
df.expect_column_values_to_match_regex("email", r"^[\w\.-]+@[\w\.-]+\.\w+$")
df.expect_column_values_to_be_in_set("country_code", ["US", "CA", "GB", "DE"])

# Validate
validation_result = df.validate()

if not validation_result["success"]:
    print("Data quality check failed!")
    for result in validation_result["results"]:
        if not result["success"]:
            print(f"Failed: {result['expectation_config']['expectation_type']}")

Great Expectations Suite Configuration

# great_expectations/expectations/users_suite.json
{
  "expectation_suite_name": "users_suite",
  "expectations": [
    {
      "expectation_type": "expect_column_to_exist",
      "kwargs": {"column": "user_id"}
    },
    {
      "expectation_type": "expect_column_values_to_not_be_null",
      "kwargs": {"column": "email"}
    },
    {
      "expectation_type": "expect_column_values_to_be_between",
      "kwargs": {
        "column": "age",
        "min_value": 0,
        "max_value": 150
      }
    }
  ]
}

4. dbt Data Quality Tests

Built-in Tests

# models/schema.yml
version: 2

models:
  - name: users
    columns:
      - name: user_id
        tests:
          - unique
          - not_null
      
      - name: email
        tests:
          - unique
          - not_null
          - dbt_utils.email_format
      
      - name: age
        tests:
          - dbt_utils.accepted_range:
              min_value: 0
              max_value: 150
      
      - name: created_at
        tests:
          - not_null
          - dbt_utils.recency:
              datepart: hour
              interval: 24

Custom dbt Tests

-- tests/assert_positive_revenue.sql
SELECT
    order_id,
    total_amount
FROM {{ ref('orders') }}
WHERE total_amount <= 0

5. Data Validation in Pipelines

Pre-Processing Validation

def validate_input_data(df: pd.DataFrame) -> bool:
    """Validate data before processing"""
    checks = {
        'row_count': len(df) > 0,
        'no_nulls_in_id': df['user_id'].notna().all(),
        'valid_emails': df['email'].str.match(r'^[\w\.-]+@[\w\.-]+\.\w+$').all(),
        'age_range': df['age'].between(0, 150).all()
    }
    
    failed_checks = [k for k, v in checks.items() if not v]
    
    if failed_checks:
        raise ValueError(f"Validation failed: {failed_checks}")
    
    return True

# Usage in pipeline
try:
    validate_input_data(raw_data)
    processed_data = transform(raw_data)
except ValueError as e:
    logger.error(f"Data validation failed: {e}")
    # Quarantine bad data
    raw_data.to_csv(f"quarantine/bad_data_{datetime.now()}.csv")
    raise

Post-Load Validation

def validate_loaded_data(table_name: str, db_connection):
    """Validate data after loading to database"""
    
    # Check row count
    expected_count = get_source_count()
    actual_count = db_connection.execute(
        f"SELECT COUNT(*) FROM {table_name}"
    ).fetchone()[0]
    
    assert actual_count == expected_count, \
        f"Row count mismatch: expected {expected_count}, got {actual_count}"
    
    # Check for nulls in critical columns
    null_check = db_connection.execute(f"""
        SELECT COUNT(*) 
        FROM {table_name} 
        WHERE user_id IS NULL OR email IS NULL
    """).fetchone()[0]
    
    assert null_check == 0, f"Found {null_check} rows with null critical fields"
    
    # Check data freshness
    latest_timestamp = db_connection.execute(f"""
        SELECT MAX(created_at) FROM {table_name}
    """).fetchone()[0]
    
    age_hours = (datetime.now() - latest_timestamp).total_seconds() / 3600
    assert age_hours < 2, f"Data is {age_hours} hours old (threshold: 2 hours)"

6. Anomaly Detection

Statistical Methods

import numpy as np
from scipy import stats

def detect_anomalies_zscore(data: pd.Series, threshold: float = 3.0):
    """Detect anomalies using Z-score method"""
    z_scores = np.abs(stats.zscore(data))
    anomalies = data[z_scores > threshold]
    return anomalies

def detect_anomalies_iqr(data: pd.Series):
    """Detect anomalies using Interquartile Range"""
    Q1 = data.quantile(0.25)
    Q3 = data.quantile(0.75)
    IQR = Q3 - Q1
    
    lower_bound = Q1 - 1.5 * IQR
    upper_bound = Q3 + 1.5 * IQR
    
    anomalies = data[(data < lower_bound) | (data > upper_bound)]
    return anomalies

# Usage
daily_revenue = df.groupby('date')['revenue'].sum()
revenue_anomalies = detect_anomalies_zscore(daily_revenue)

if len(revenue_anomalies) > 0:
    alert(f"Revenue anomalies detected: {revenue_anomalies}")

ML-Based Anomaly Detection

from sklearn.ensemble import IsolationForest

def detect_anomalies_ml(df: pd.DataFrame, features: list):
    """Detect anomalies using Isolation Forest"""
    model = IsolationForest(contamination=0.01, random_state=42)
    
    # Fit and predict
    predictions = model.fit_predict(df[features])
    
    # -1 indicates anomaly
    anomalies = df[predictions == -1]
    return anomalies

# Usage
anomalies = detect_anomalies_ml(
    df, 
    features=['order_count', 'total_revenue', 'avg_order_value']
)

7. Data Profiling

import pandas as pd

def profile_dataframe(df: pd.DataFrame) -> dict:
    """Generate comprehensive data profile"""
    profile = {
        'row_count': len(df),
        'column_count': len(df.columns),
        'columns': {}
    }
    
    for col in df.columns:
        profile['columns'][col] = {
            'dtype': str(df[col].dtype),
            'null_count': df[col].isna().sum(),
            'null_percentage': df[col].isna().mean() * 100,
            'unique_count': df[col].nunique(),
            'cardinality': df[col].nunique() / len(df) * 100
        }
        
        # Numeric columns
        if pd.api.types.is_numeric_dtype(df[col]):
            profile['columns'][col].update({
                'min': df[col].min(),
                'max': df[col].max(),
                'mean': df[col].mean(),
                'median': df[col].median(),
                'std': df[col].std()
            })
        
        # String columns
        elif pd.api.types.is_string_dtype(df[col]):
            profile['columns'][col].update({
                'min_length': df[col].str.len().min(),
                'max_length': df[col].str.len().max(),
                'avg_length': df[col].str.len().mean()
            })
    
    return profile

8. Handling Data Quality Failures

Strategy 1: Fail Pipeline

def process_with_strict_validation(data):
    """Fail entire pipeline on any validation error"""
    if not validate_data(data):
        raise DataQualityError("Validation failed - stopping pipeline")
    
    return transform(data)

Strategy 2: Quarantine Bad Data

def process_with_quarantine(data):
    """Separate good and bad data"""
    valid_data = data[validate_rows(data)]
    invalid_data = data[~validate_rows(data)]
    
    if len(invalid_data) > 0:
        invalid_data.to_csv(f"quarantine/{datetime.now()}.csv")
        alert(f"Quarantined {len(invalid_data)} invalid rows")
    
    return transform(valid_data)

Strategy 3: Alert and Continue

def process_with_alerts(data):
    """Log issues but continue processing"""
    validation_results = validate_data(data)
    
    if not validation_results['success']:
        alert(f"Data quality issues: {validation_results['failures']}")
        log_to_monitoring(validation_results)
    
    # Continue processing anyway
    return transform(data)

9. Data Quality Metrics

def calculate_dq_score(validation_results: dict) -> float:
    """Calculate overall data quality score (0-100)"""
    total_checks = len(validation_results)
    passed_checks = sum(1 for r in validation_results.values() if r['passed'])
    
    return (passed_checks / total_checks) * 100

# Track over time
dq_scores = []
for date, data in daily_data.items():
    results = validate_data(data)
    score = calculate_dq_score(results)
    dq_scores.append({'date': date, 'score': score})

# Alert if score drops
if score < 95:
    alert(f"Data quality score dropped to {score}%")

10. Tools Comparison

ToolBest ForProsCons
Great ExpectationsPython pipelinesComprehensive, well-documentedLearning curve
dbt testsSQL transformationsIntegrated with dbt, simpleLimited to SQL
SodaCollaborationBusiness-friendly, SaaSPaid for advanced features
Monte CarloObservabilityML-based anomaly detectionExpensive

11. Real-World Data Quality Issues

Case Study: The Missing Orders

  • Problem: 10% of orders missing from data warehouse
  • Root Cause: ETL pipeline skipped records with null shipping_address
  • Solution: Added validation to fail pipeline if > 1% of records are skipped
  • Prevention: Implemented pre-load row count validation

Case Study: The Duplicate Customers

  • Problem: Same customer appearing multiple times with different IDs
  • Root Cause: No uniqueness check on email during ingestion
  • Solution: Added UNIQUE constraint on email, deduplication logic
  • Prevention: Implemented fuzzy matching for duplicate detection

12. Data Quality Checklist

  • Completeness: Are all required fields populated?
  • Uniqueness: Are primary keys truly unique?
  • Validity: Do all fields match expected formats?
  • Consistency: Is data consistent across related tables?
  • Freshness: Is data updated within SLA?
  • Accuracy: Have we validated against source systems?
  • Monitoring: Are we tracking DQ metrics over time?
  • Alerting: Do we get notified of quality degradation?

  • 43-data-reliability/data-contracts
  • 43-data-reliability/schema-management
  • 43-data-reliability/data-lineage

スコア

総合スコア

60/100

リポジトリの品質指標に基づく評価

SKILL.md

SKILL.mdファイルが含まれている

+20
LICENSE

ライセンスが設定されている

+10
説明文

100文字以上の説明がある

0/10
人気

GitHub Stars 100以上

0/15
最近の活動

3ヶ月以内に更新がある

0/10
フォーク

10回以上フォークされている

0/5
Issue管理

オープンIssueが50未満

+5
言語

プログラミング言語が設定されている

+5
タグ

1つ以上のタグが設定されている

0/5

レビュー

💬

レビュー機能は近日公開予定です