Back to blog
Data EngineeringSeptember 1, 2026

Data Quality Framework: How We Catch Bad Data Before It Reaches the Dashboard

Framework for catching bad data before it reaches the dashboard.

Data Quality Framework: How We Catch Bad Data Before It Reaches the Dashboard

The Problem

Your CEO opens the Monday morning dashboard and sees revenue is down 40% week-over-week. The marketing team panics and pauses all campaigns. The product team calls an emergency meeting. Three hours later, someone discovers an API change in your payment processor that caused duplicate transactions to be excluded from the warehouse.

The revenue didn't drop. The data did.

Bad data is worse than no data. No data forces you to admit ignorance. Bad data lets you confidently make terrible decisions.


The Cost of Bad Data

ImpactExampleCost
Wrong decisionsPause campaigns based on bad revenue data$50k+ in lost revenue
Lost trustStakeholders stop believing dashboardsMonths of credibility rebuilding
ReworkAnalysts spend time debugging instead of analyzing30-50% of analyst time
ComplianceIncorrect reporting to regulatorsFines, legal exposure
Customer experienceWrong personalization from bad profile dataChurn, negative reviews

The Framework: Five Layers of Defense


Layer 1: Schema Validation (at ingestion)

    |

Layer 2: dbt Tests (in transformation)

    |

Layer 3: Anomaly Detection (in warehouse)

    |

Layer 4: Reconciliation (vs source systems)

    |

Layer 5: Semantic Validation (business logic)

Layer 1: Schema Validation

JSON Schema Enforcement


{

  "$id": "https://[PROPERTY].com/schemas/transaction",

  "type": "object",

  "required": ["transaction_id", "amount", "currency", "timestamp"],

  "properties": {

    "transaction_id": {

      "type": "string",

      "pattern": "^txn_[a-z0-9]{16}$"

    },

    "amount": {

      "type": "number",

      "minimum": 0,

      "maximum": 1000000

    },

    "currency": {

      "type": "string",

      "enum": ["USD", "EUR", "GBP", "AUD"]

    },

    "timestamp": {

      "type": "string",

      "format": "date-time"

    },

    "customer_id": {

      "type": "string"

    },

    "items": {

      "type": "array",

      "items": {

        "type": "object",

        "required": ["sku", "quantity", "price"],

        "properties": {

          "sku": { "type": "string" },

          "quantity": { "type": "integer", "minimum": 1 },

          "price": { "type": "number", "minimum": 0 }

        }

      }

    }

  }

}

Python Validation


from jsonschema import validate, ValidationError

import json



def validate_event(event, schema):

    try:

        validate(instance=event, schema=schema)

        return {"valid": True, "error": None}

    except ValidationError as e:

        return {"valid": False, "error": str(e)}



# Usage

result = validate_event(incoming_event, transaction_schema)

if not result["valid"]:

    send_to_dead_letter_queue(event, result["error"])

Layer 2: dbt Tests

Built-In Tests


# schema.yml

version: 2



models:

  - name: fct_transactions

    columns:

      - name: transaction_id

        tests:

          - unique

          - not_null



      - name: amount

        tests:

          - not_null

          - dbt_utils.accepted_range:

              min_value: 0

              max_value: 1000000



      - name: currency

        tests:

          - accepted_values:

              values: ['USD', 'EUR', 'GBP', 'AUD']



      - name: customer_id

        tests:

          - not_null

          - relationships:

              to: ref('dim_customers')

              field: customer_id



      - name: created_at

        tests:

          - not_null

          - dbt_utils.accepted_range:

              min_value: "'2024-01-01'"

              max_value: "current_timestamp + interval '1 day'"

Custom Tests


-- tests/transaction_amount_matches_items.sql

-- Verify transaction amount equals sum of line items



select

  transaction_id,

  amount as transaction_amount,

  sum(items.value:price::numeric * items.value:quantity::numeric) as items_total

from fct_transactions,

lateral flatten(input => items) as items

group by 1, 2

having abs(transaction_amount - items_total) > 0.01

-- tests/no_orphaned_transactions.sql

-- Every transaction must have a valid customer



select transaction_id

from fct_transactions t

left join dim_customers c on t.customer_id = c.customer_id

where c.customer_id is null

  and t.created_at >= current_date - 7

-- tests/revenue_not_negative.sql

-- Daily revenue should never be negative



select

  date(created_at) as date,

  sum(amount) as daily_revenue

from fct_transactions

group by 1

having daily_revenue < 0

Layer 3: Anomaly Detection

Statistical Process Control


-- Flag days where revenue deviates > 3 standard deviations from 30-day mean

with daily_stats as (

  select

    date(created_at) as date,

    sum(amount) as revenue,

    avg(sum(amount)) over (order by date rows between 30 preceding and 1 preceding) as avg_revenue,

    stddev(sum(amount)) over (order by date rows between 30 preceding and 1 preceding) as std_revenue

  from fct_transactions

  group by 1

)



select

  date,

  revenue,

  avg_revenue,

  (revenue - avg_revenue) / nullif(std_revenue, 0) as z_score

from daily_stats

where abs(z_score) > 3

  and date >= current_date - 7

order by date desc;

Volume Monitoring


-- Alert if event volume drops > 30% vs same day last week

with hourly_comparison as (

  select

    date_trunc('hour', received_at) as hour,

    count(*) as current_count,

    lag(count(*)) over (partition by extract(dow from received_at), extract(hour from received_at) order by date_trunc('week', received_at)) as prior_count

  from raw.events

  where received_at >= current_date - 14

  group by 1

)



select

  hour,

  current_count,

  prior_count,

  (current_count - prior_count) / nullif(prior_count, 0) as pct_change

from hourly_comparison

where hour >= current_date - 1

  and abs(pct_change) > 0.30

order by hour desc;

Python Anomaly Detection


from scipy import stats

import numpy as np



def detect_anomalies(values, method='zscore', threshold=3):

    """

    Detect anomalous values in a time series.



    Methods:

    - zscore: values beyond threshold standard deviations

    - iqr: values beyond 1.5 * IQR

    - mad: median absolute deviation

    """

    if method == 'zscore':

        z_scores = np.abs(stats.zscore(values))

        return z_scores > threshold



    elif method == 'iqr':

        q1, q3 = np.percentile(values, [25, 75])

        iqr = q3 - q1

        lower = q1 - 1.5 * iqr

        upper = q3 + 1.5 * iqr

        return (values < lower) | (values > upper)



    elif method == 'mad':

        median = np.median(values)

        mad = np.median(np.abs(values - median))

        modified_z = 0.6745 * (values - median) / mad

        return np.abs(modified_z) > threshold



# Usage

daily_revenue = [12000, 12500, 11900, 12100, 45000, 12200, 11800]

anomalies = detect_anomalies(daily_revenue, method='zscore')

# Anomaly detected at index 4 (revenue spike to 45k)

Layer 4: Reconciliation

Source System Comparison


-- Compare warehouse totals to Stripe dashboard

with warehouse_totals as (

  select

    date_trunc('day', created_at) as date,

    count(*) as transaction_count,

    sum(amount) as total_amount

  from fct_transactions

  where source_system = 'stripe'

  group by 1

),



stripe_totals as (

  select

    date_trunc('day', created) as date,

    count(*) as transaction_count,

    sum(amount) / 100.0 as total_amount  -- Stripe stores cents

  from raw.stripe_charges

  group by 1

)



select

  w.date,

  w.transaction_count as warehouse_count,

  s.transaction_count as stripe_count,

  w.total_amount as warehouse_amount,

  s.total_amount as stripe_amount,

  abs(w.transaction_count - s.transaction_count) as count_diff,

  abs(w.total_amount - s.total_amount) as amount_diff

from warehouse_totals w

full outer join stripe_totals s on w.date = s.date

where abs(w.total_amount - s.total_amount) > 100

   or abs(w.transaction_count - s.transaction_count) > 10

order by w.date desc;

Row-Level Reconciliation


-- Find transactions in Stripe but missing from warehouse

select

  s.id as stripe_transaction_id,

  s.amount / 100.0 as stripe_amount,

  s.created as stripe_created

from raw.stripe_charges s

left join fct_transactions w on s.id = w.source_transaction_id

where w.transaction_id is null

  and s.created >= current_date - 1

Layer 5: Semantic Validation

Business Rule Checks


-- A customer's lifetime revenue should equal sum of their transactions

select

  c.customer_id,

  c.lifetime_revenue as reported_ltv,

  sum(t.amount) as calculated_ltv,

  abs(c.lifetime_revenue - sum(t.amount)) as discrepancy

from dim_customers c

left join fct_transactions t on c.customer_id = t.customer_id

group by 1, 2

having discrepancy > 0.01

-- Refund amount should not exceed original transaction

select

  r.refund_id,

  r.amount as refund_amount,

  t.amount as original_amount

from fct_refunds r

join fct_transactions t on r.original_transaction_id = t.transaction_id

where r.amount > t.amount

-- Active subscriptions should have a valid payment method

select

  s.subscription_id,

  s.customer_id,

  s.status

from fct_subscriptions s

left join dim_payment_methods p on s.customer_id = p.customer_id

  and p.is_default = true

where s.status = 'active'

  and p.payment_method_id is null

Data Quality Scorecard

Metrics to Track

MetricTargetAlert Threshold
Test pass rate>99%<95%
Schema validation rate>99.5%<98%
Source reconciliation accuracy100%<99.9%
Anomaly false positive rate<5%>10%
Time to detect data issue<1 hour>4 hours
Time to resolve data issue<4 hours>24 hours

Dashboard Query


-- Daily data quality summary

select

  current_date as report_date,

  count(*) as total_tests,

  sum(case when status = 'pass' then 1 else 0 end) as passed,

  sum(case when status = 'fail' then 1 else 0 end) as failed,

  sum(case when status = 'error' then 1 else 0 end) as errors,

  passed::float / nullif(total_tests, 0) as pass_rate

from dbt_test_results

where run_date = current_date;

Incident Response

Severity Levels

LevelCriteriaResponse
P1Revenue/transaction data incorrectPage on-call immediately
P2Dimension data stale (>24h)Slack alert, fix within 4h
P3Test failure, no business impactTicket, fix within 24h
P4Documentation/observability gapBacklog

Runbook Template


# Data Incident: [DESCRIPTION]



## Detection

- Alert: [which check failed]

- Time: [when detected]

- Impact: [which dashboards/reports affected]



## Investigation

1. Check source system health

2. Check pipeline execution logs

3. Identify failing records

4. Determine root cause



## Resolution

- Fix applied: [what changed]

- Data backfill required: [yes/no]

- Backfill query: [if applicable]



## Prevention

- Test added: [which test prevents recurrence]

- Documentation updated: [link]

Implementation Checklist

  • [ ] Schema registry with validation for all event types
  • [ ] dbt tests on all marts models (unique, not_null, relationships)
  • [ ] Custom tests for business logic
  • [ ] Anomaly detection on key metrics (revenue, users, conversions)
  • [ ] Source reconciliation queries (daily automated comparison)
  • [ ] Semantic validation rules
  • [ ] Data quality dashboard
  • [ ] Incident response runbook
  • [ ] On-call rotation for data issues
  • [ ] Monthly data quality review meeting

Key Principle

Data quality is not a one-time project. It's a continuous process that requires\

the same rigor as software engineering: tests, monitoring, incident response,\

and retrospectives.

The goal isn't perfect data. The goal is knowing when your data is imperfect\

before someone makes a decision based on it.