Skip to content
Satya Prakash Solanki

In Part 3 of this series I showed how to connect Great Expectations to Trino, write an expectation suite and run a checkpoint. That is the unit test. This part is about everything around it: which rules exist and why, who owns them, where in the pipeline they run, how to catch problems no one wrote a rule for, and how to alert without training people to ignore alerts.

The shift is from a set of checks to a system that makes decisions. A suite that runs and produces an HTML report is useful. A system that blocks a bad partition, tells the right team, and stays quiet the rest of the time is what production needs.

The rule catalogue

Every check belongs in one catalogue, version-controlled and reviewed like code. Great Expectations suites, PyDeequ constraints and hand-written reconciliation SQL are implementations. The catalogue is the specification they implement.

I give each rule the same anatomy:

catalogue/finance.payments.yaml
- id: payments.amount_non_negative
dataset: lake.finance.payments
dimension: validity # completeness | validity | uniqueness | consistency | timeliness
stage: transform # ingest | transform | serve
check:
type: expectation
expectation: expect_column_values_to_be_between
kwargs: { column: amount, min_value: 0 }
mostly: 0.999
severity: high
owner: payments-data-eng
consumers: [finance-reporting, fraud-features]
runbook: runbooks/payments-amount.md
added: 2026-03-02
reason: Refunds are modelled as separate rows; negatives indicate a sign bug upstream.

Three fields do most of the work.

  • reason. A rule without a reason cannot be safely changed or removed. When a threshold starts firing, the reason tells the on-call engineer whether the data or the rule is wrong.
  • owner. One team, not a person and not “data platform”. The owner fixes the data or changes the rule.
  • consumers. Who is harmed when it fails. This drives severity and tells you who to notify when you publish suspect data.

The catalogue also answers coverage questions mechanically. Which critical tables have no timeliness rule? Which rules have not fired in a year and may be dead? Which owner holds forty high-severity rules and probably needs help?

Severity decides behaviour

Severity is not a colour on a dashboard. It is the action the pipeline takes.

Severity Pipeline action Notification
Critical Stop the publish. Previous version stays live. Page the owner
High Write to a quarantine location; dependants wait Ticket to the owner
Medium Publish, mark the partition as suspect in metadata Team channel
Low Publish, record the result Weekly digest

Two rules keep this honest. First, severity is set by consumer impact, not by how confident the author is in the check. A wrong total in a regulatory report is critical even if the check is crude. Second, a rule that fires at the same severity every week without action is wrong. Either it is critical and the pipeline should be blocked, or it is not and it should be demoted.

Where checks run

The same rule costs very different amounts depending on where it runs.

  1. 01Ingest: contract and schema
  2. 02Land: counts and reconciliation
  3. 03Transform: business rules
  4. 04Serve: freshness and SLAs
Figure 1. Checks by stage. Each stage catches what the previous one cannot see.
  • Ingest. Validate structure against the contract: schema, required fields, types, enum values. These checks are cheap, run per batch or per message, and reject at the boundary where the producer can still be told.
  • Land. Reconcile with the source: counts, checksums, duplicates. This is covered in depth in the source-to-target reconciliation case study.
  • Transform. Business rules and cross-table integrity: totals that must balance, foreign keys, derived fields. These run on the output of each PySpark job, before it is published.
  • Serve. What the consumer experiences: freshness, availability, and query-level checks on the views that dashboards and models read through Trino.

A practical test for placement: run a check at the earliest stage that has the information it needs. A null-rate check on a raw field belongs at ingest. A check that revenue equals the sum of line items belongs after the join that produces it. Running everything at the serving layer finds problems late and makes them expensive to trace.

Anomaly detection without a data science project

Rules catch what you anticipated. Monitors catch what you did not. Three monitors cover most of the gap, and none needs machine learning to be useful.

I compute them from a metrics table that every load writes to: table, partition, row count, max event timestamp, null counts and a few column profiles. PyDeequ’s metrics repository, which I described in Part 2, is one way to populate it. A plain Iceberg table written by the pipeline works just as well.

Volume

Compare today’s row count with the same weekday over recent weeks, using the median and the median absolute deviation (MAD). Both are resistant to the occasional outlier day that would drag a mean and standard deviation around.

volume_anomaly.sql
WITH hist AS (
SELECT table_name, row_count
FROM dq.load_metrics
WHERE partition_date BETWEEN current_date - INTERVAL '56' DAY AND current_date - INTERVAL '1' DAY
AND day_of_week(partition_date) = day_of_week(current_date)
),
med AS (
SELECT table_name, approx_percentile(row_count, 0.5) AS median_n
FROM hist GROUP BY table_name
),
mad AS (
SELECT h.table_name, m.median_n,
approx_percentile(abs(h.row_count - m.median_n), 0.5) AS mad_n
FROM hist h JOIN med m ON h.table_name = m.table_name
GROUP BY h.table_name, m.median_n
)
SELECT t.table_name, t.row_count, b.median_n, b.mad_n,
(t.row_count - b.median_n) / nullif(1.4826 * b.mad_n, 0) AS mad_z
FROM dq.load_metrics t
JOIN mad b ON t.table_name = b.table_name
WHERE t.partition_date = current_date
AND abs((t.row_count - b.median_n) / nullif(1.4826 * b.mad_n, 0)) > 4;

The constant 1.4826 scales MAD to be comparable with a standard deviation for normally distributed data. The threshold of 4 is an example starting value; tune it per table against its history. Handle the nullif case explicitly: a table whose count never varies has a MAD of zero, and any change at all is worth a look.

Freshness

Freshness is the lag between now and the newest event in the table, compared with an agreed expectation. It is the monitor consumers feel first, and the one most often missing.

freshness.sql
SELECT 'lake.finance.payments' AS table_name,
max(event_ts) AS newest_event,
date_diff('minute', max(event_ts),
CAST(current_timestamp AT TIME ZONE 'UTC' AS timestamp(6))) AS lag_minutes
FROM lake.finance.payments
WHERE event_date >= current_date - INTERVAL '1' DAY;

Set the threshold from the consumer’s need, written as a service level objective, not from how fast the pipeline usually runs. If finance reads the table at 07:00, the objective is “loaded by 06:30”, and the alert should fire at 06:30, not when the job is a few minutes slower than average. The approach mirrors the SRE treatment of SLOs: measure what the user experiences and alert on that.

Distribution

Row counts can be perfect while the content shifts: a currency field that suddenly holds only one value, a category that disappears, an amount distribution that moves after an upstream unit change. For categorical and bucketed numeric columns I use the population stability index (PSI) against a reference window.

psi.py
import numpy as np
def psi(expected: dict, actual: dict, eps: float = 1e-6) -> float:
"""PSI between two bucket -> count maps."""
keys = set(expected) | set(actual)
e_total = sum(expected.values()) or 1
a_total = sum(actual.values()) or 1
score = 0.0
for k in keys:
e = max(expected.get(k, 0) / e_total, eps)
a = max(actual.get(k, 0) / a_total, eps)
score += (a - e) * np.log(a / e)
return score
# Bucket counts come from a Trino GROUP BY, not from pulling rows into Python.

Common rules of thumb treat PSI below about 0.1 as stable and above about 0.25 as a significant shift. They are conventions, not laws, so calibrate them on your own tables. Great Expectations offers a related check, expect_column_kl_divergence_to_be_less_than, which appeared in the expectation list in Part 3. Either works. What matters is a fixed reference window, buckets computed in the engine, and a threshold chosen from history.

Alerting without noise

Data-quality programmes rarely fail because they lack checks. They fail because the checks are ignored. A few habits keep alert volume proportional to real problems.

  • Alert on incidents, not on rule results. One late upstream table can fail freshness on twenty downstream tables. Group failures by root dataset using lineage, and send one alert that lists the affected dependants.
  • Suppress downstream of a known failure. If a critical rule blocked a publish, dependants will look stale. That is expected and should not page anyone else.
  • Require persistence for monitors. A medium-severity anomaly that clears on the next load is a digest item. Two consecutive breaches become an alert.
  • Track noise per rule. Record how each alert was closed. A rule whose alerts are mostly closed as “expected” gets retuned, demoted or deleted within a fixed period. Its owner is accountable for that, not the on-call engineer.
  • Make every alert actionable. It links to the runbook from the catalogue, the failing query, and the pinned source and target positions so it can be reproduced.

Data contracts: moving checks to the producer

Most data issues are created upstream of the team that finds them. A data contract makes the producer’s promises explicit and testable: schema, semantics, freshness and volume expectations, and how changes are announced.

contracts/payments.v2.yaml
dataset: payments
version: 2.1.0
owner: payments-platform
schema:
- { name: payment_id, type: string, required: true, unique: true }
- { name: amount, type: decimal(18,2), required: true, min: 0 }
- { name: currency, type: string, required: true, allowed: [AED, USD, EUR, INR] }
- { name: event_ts, type: timestamp, required: true, timezone: UTC }
service_levels:
freshness_minutes: 30
expected_daily_rows: { min: 50000, max: 400000 }
change_policy:
additive: minor version, 7 days notice
breaking: major version, parallel run 30 days

The values here are examples. The important properties are that the contract lives with the producer’s code, the producer’s CI validates sample output against it before release, and the consumer’s ingest stage validates against the same file. A breaking change then fails the producer’s build rather than the consumer’s dashboard.

Contracts do not replace monitoring. A producer can honour every field and still send half the usual volume because of an outage. They narrow what monitoring needs to catch, and they move the conversation about breaking changes to before the change, which is where it belongs.

Putting it together

The progression I recommend is incremental. Start with a catalogue and severities for the ten tables that matter most. Add freshness and volume monitors to all of them, because those two catch a large share of real incidents at low cost. Add reconciliation for tables fed from operational systems. Introduce contracts with the producers who change schemas most often. Then grow coverage using the catalogue’s own gap reports. Incremental and CDC-fed tables need additional checks, covered in the next part on validating CDC pipelines.

Further reading

Big data testing & data quality

Try “evaluation”, “red-teaming”, “governance” or “agents”.