3. Advanced rules¶
A migration review compares a pipeline's output before and after a rewrite: each difference must be declared by a rule or reported. This tutorial reviews 1,000 NYC taxi trips. The source is last night's warehouse extract, and the target is the same trips from the rewritten pipeline.
id, from 0 to 999, is the unique primary key. Every change is deterministic, so the recorded results are stable. Every run sets strict_types=True. Without it, each target column is cast to its source column's type, which would reconcile text codes and padded zone IDs without a rule. With it, each rule that aligns types is visibly needed.
veridelta.datasets.load_nyc_taxi() downloads the data once and caches it under ~/.cache/veridelta/datasets.
Open this tutorial in Google Colab. Its first cell installs Veridelta there.
# A kernel without Veridelta, such as Colab's, installs the release this tutorial shows.
import importlib.util
import subprocess
import sys
if importlib.util.find_spec("veridelta") is None:
subprocess.run(
[sys.executable, "-m", "pip", "install", "--quiet", "veridelta==0.35.1"],
check=True,
)
1. The rewritten pipeline¶
The rewrite changes the data in ways the review expects:
- fares and totals gain sub-cent rounding, and tips are recalculated;
- the store-and-forward flag is lowercased and padded with spaces;
- trip distance is text with a unit suffix, such as
1.2 mi; - dropoff times are text in
dd/mm/YYYYorder; - payment type is text, and pickup zone codes are zero-padded;
- rate code
99, a legacy placeholder, becomes NULL; Airport_feeis renamedairport_fee, and two_etl_*bookkeeping columns appear;- four trips are dropped, and three are cloned under new IDs.
One change is a defect: every trip with id % 100 == 7 has its pickup time shifted by one hour, and no rule should forgive those ten trips. The cell below loads the trips and applies every change:
import polars as pl
from veridelta import DiffConfig, DiffEngine, DiffRule
from veridelta.datasets import load_nyc_taxi
from veridelta.models import DiffResult
source = load_nyc_taxi()
target = (
source.with_columns(
fare_amount=pl.col("fare_amount") + 0.004,
total_amount=pl.col("total_amount") + 0.004,
tip_amount=pl.col("tip_amount") * 1.005,
store_and_fwd_flag=" " + pl.col("store_and_fwd_flag").str.to_lowercase() + " ",
trip_distance=pl.col("trip_distance").cast(pl.String) + pl.lit(" mi"),
tpep_dropoff_datetime=pl.col("tpep_dropoff_datetime").dt.strftime("%d/%m/%Y %H:%M:%S"),
payment_type=pl.col("payment_type").cast(pl.String),
PULocationID=pl.col("PULocationID").cast(pl.String).str.pad_start(3, "0"),
RatecodeID=pl.when(pl.col("RatecodeID") == 99).then(None).otherwise(pl.col("RatecodeID")),
tpep_pickup_datetime=pl.when((pl.col("id") % 100) == 7)
.then(pl.col("tpep_pickup_datetime") + pl.duration(hours=1))
.otherwise(pl.col("tpep_pickup_datetime")),
_etl_batch=pl.lit(1),
_etl_loaded_at=pl.lit("nightly"),
)
.rename({"Airport_fee": "airport_fee"})
.filter((pl.col("id") % 250) != 0)
)
clones = target.head(3).with_columns(id=pl.Series([1000, 1001, 1002], dtype=pl.Int64))
target = pl.concat([target, clones], how="vertical")
source_lf, target_lf = source.lazy(), target.lazy()
print(f"source={source.height} target={target.height}")
def verdict(label: str, result: DiffResult) -> None:
"""Print one run's status, row counts, and columns ranked by drift."""
summary = result.summary
status = "PASSED" if summary.is_match else "FAILED"
ranked = (
", ".join(
f"{name}={count}"
for name, count in sorted(
summary.column_mismatches.items(), key=lambda item: (-item[1], item[0])
)
)
or "none"
)
print(
f"{label}: {status} changed={summary.changed_count} "
f"added={summary.added_count} removed={summary.removed_count} drift=[{ranked}]"
)
# Output:
# source=1000 target=999
2. Baseline¶
The default schema_mode, intersection, compares only the columns both sides share. Airport_fee and airport_fee differ in case, so neither is compared. Without rules, every shared row is changed:
baseline = DiffEngine(
DiffConfig(primary_keys=["id"], strict_types=True), source_lf, target_lf
).run()
verdict("1 baseline", baseline)
print("airport_fee compared:", "airport_fee" in baseline.compared_columns)
print("Airport_fee compared:", "Airport_fee" in baseline.compared_columns)
# Output:
# 1 baseline: FAILED changed=996 added=3 removed=4 drift=[PULocationID=996, fare_amount=996, payment_type=996, store_and_fwd_flag=996, total_amount=996, tpep_dropoff_datetime=996, trip_distance=996, tip_amount=770, tpep_pickup_datetime=10, RatecodeID=3]
# airport_fee compared: False
# Airport_fee compared: False
3. Value rules¶
These rules forgive the expected changes in values. A pattern rule gives every *_amount column an absolute tolerance of one cent. A rule that names tip_amount wins over the pattern, with a relative tolerance of 1%. The flag ignores whitespace and case, trip distance loses its suffix and becomes a number, and rate code 99 becomes NULL. The drift that remains is in types and times:
value_rules = [
DiffRule(pattern=r".*_amount$", absolute_tolerance=0.01),
DiffRule(column_names=["tip_amount"], relative_tolerance=0.01),
DiffRule(
column_names=["store_and_fwd_flag"],
whitespace_mode="both",
case_insensitive=True,
),
DiffRule(
column_names=["trip_distance"],
regex_replace={r"\s*mi$": ""},
cast_to="Float64",
),
DiffRule(column_names=["RatecodeID"], null_values=[99]),
]
values = DiffEngine(
DiffConfig(primary_keys=["id"], strict_types=True, rules=value_rules),
source_lf,
target_lf,
).run()
verdict("2 value rules", values)
# Output:
# 2 value rules: FAILED changed=996 added=3 removed=4 drift=[PULocationID=996, payment_type=996, tpep_dropoff_datetime=996, tpep_pickup_datetime=10]
4. Type and time rules¶
cast_to, pad_zeros, and datetime_format align the remaining type differences. What is left is the shifted pickup time:
type_rules = [
*value_rules,
DiffRule(column_names=["payment_type"], cast_to="Int64"),
DiffRule(column_names=["PULocationID"], pad_zeros=3),
DiffRule(column_names=["tpep_dropoff_datetime"], datetime_format="%d/%m/%Y %H:%M:%S"),
]
typed = DiffEngine(
DiffConfig(primary_keys=["id"], strict_types=True, rules=type_rules),
source_lf,
target_lf,
).run()
verdict("3 type and time", typed)
# Output:
# 3 type and time: FAILED changed=10 added=3 removed=4 drift=[tpep_pickup_datetime=10]
5. Structure¶
rename_to maps Airport_fee onto airport_fee, and a pattern rule with ignore leaves out the _etl_* bookkeeping columns. schema_mode="allow_additions" then fails the run if the target loses a source column, and accepts the columns it adds. What remains is the defect:
structure_rules = [
*type_rules,
DiffRule(column_names=["Airport_fee"], rename_to="airport_fee"),
DiffRule(pattern=r"^_etl_", ignore=True),
]
final = DiffEngine(
DiffConfig(
primary_keys=["id"],
strict_types=True,
schema_mode="allow_additions",
rules=structure_rules,
),
source_lf,
target_lf,
).run()
verdict("4 structure", final)
print(final.get_mismatches("tpep_pickup_datetime"))
# Output:
# 4 structure: FAILED changed=10 added=3 removed=4 drift=[tpep_pickup_datetime=10]
# shape: (10, 3)
# ┌─────┬─────────────────────────────┬─────────────────────────────┐
# │ id ┆ tpep_pickup_datetime_source ┆ tpep_pickup_datetime_target │
# │ --- ┆ --- ┆ --- │
# │ i64 ┆ datetime[μs] ┆ datetime[μs] │
# ╞═════╪═════════════════════════════╪═════════════════════════════╡
# │ 7 ┆ 2026-01-01 00:34:28 ┆ 2026-01-01 01:34:28 │
# │ 107 ┆ 2026-01-01 00:20:33 ┆ 2026-01-01 01:20:33 │
# │ 207 ┆ 2026-01-01 00:58:44 ┆ 2026-01-01 01:58:44 │
# │ 307 ┆ 2026-01-01 00:24:47 ┆ 2026-01-01 01:24:47 │
# │ 407 ┆ 2026-01-01 00:58:52 ┆ 2026-01-01 01:58:52 │
# │ 507 ┆ 2026-01-01 00:28:51 ┆ 2026-01-01 01:28:51 │
# │ 607 ┆ 2026-01-01 00:22:13 ┆ 2026-01-01 01:22:13 │
# │ 707 ┆ 2026-01-01 00:14:57 ┆ 2026-01-01 01:14:57 │
# │ 807 ┆ 2026-01-01 00:05:21 ┆ 2026-01-01 01:05:21 │
# │ 907 ┆ 2026-01-01 00:59:11 ┆ 2026-01-01 01:59:11 │
# └─────┴─────────────────────────────┴─────────────────────────────┘
Ten trips remain, each with its pickup time shifted by exactly one hour. A threshold of 0.02 would pass this run: 10 changed, 4 removed, and 3 added rows are a 1.7% mismatch ratio on 1,000 source rows. A threshold forgives any difference up to a share of rows, the defect included. Rules forgive only the differences they name, so what they leave is the finding.
6. YAML for CI¶
The YAML and CLI tutorial runs a configuration file from the command line, and the HTML reports tutorial turns a result into a report. This cell writes both sides to Parquet, writes the review's configuration beside them, and runs the file, which reaches the verdict above:
import yaml
from veridelta import load_config
source_lf.collect().write_parquet("extract.parquet")
target_lf.collect().write_parquet("rewrite.parquet")
settings = DiffConfig(
primary_keys=["id"],
strict_types=True,
schema_mode="allow_additions",
rules=structure_rules,
).model_dump(mode="json", exclude_none=True, exclude_defaults=True)
with open("veridelta.yaml", "w", encoding="utf-8") as file:
yaml.safe_dump(
{"source": {"path": "extract.parquet"}, "target": {"path": "rewrite.parquet"}, **settings},
file,
sort_keys=False,
)
with open("veridelta.yaml", encoding="utf-8") as file:
print(file.read())
verdict("5 from the file", DiffEngine.run_from_configs(*load_config("veridelta.yaml")))
# Output:
# source:
# path: extract.parquet
# target:
# path: rewrite.parquet
# primary_keys:
# - id
# schema_mode: allow_additions
# strict_types: true
# rules:
# - pattern: .*_amount$
# absolute_tolerance: 0.01
# - column_names:
# - tip_amount
# relative_tolerance: 0.01
# - column_names:
# - store_and_fwd_flag
# case_insensitive: true
# whitespace_mode: both
# - column_names:
# - trip_distance
# regex_replace:
# \s*mi$: ''
# cast_to: Float64
# - column_names:
# - RatecodeID
# null_values:
# - 99
# - column_names:
# - payment_type
# cast_to: Int64
# - column_names:
# - PULocationID
# pad_zeros: 3
# - column_names:
# - tpep_dropoff_datetime
# datetime_format: '%d/%m/%Y %H:%M:%S'
# - column_names:
# - Airport_fee
# rename_to: airport_fee
# - pattern: ^_etl_
# ignore: true
#
# 5 from the file: FAILED changed=10 added=3 removed=4 drift=[tpep_pickup_datetime=10]