Skip to content

pudl_diff.rows

Comparing the rows of two tables, with or without a primary key.

MAX_ROWS_PER_PARTITION = 100000000 module-attribute

Tables with more rows than this have their key hashes reduced and joined in several partitions, one after the other, to bound the memory the comparison needs. At around 75 bytes of peak memory per row, 100 million rows is roughly 8 GB.

KeyedRowDiff dataclass

Row-level differences between two tables that share a primary key.

Like RowSetDiff, the differing rows are backed by temporary Parquet files (owned by pk_diff) rather than held in memory.

Source code in src/pudl_diff/rows.py
@dataclass(frozen=True)
class KeyedRowDiff:
    """Row-level differences between two tables that share a primary key.

    Like `RowSetDiff`, the differing rows are backed by temporary Parquet
    files (owned by `pk_diff`) rather than held in memory.
    """

    pk_diff: RowSetDiff
    """Symmetric difference of primary key values: full rows present in only
    one of the two tables."""

    column_changes: dict[str, int]
    """Maps each non-primary-key column to the number of shared-primary-key
    rows where its value differs between the two tables. Columns with no
    changes are omitted."""

    changed_row_count: int
    """Number of shared-primary-key rows with at least one differing
    non-primary-key value."""

    changed_left: pl.LazyFrame
    """Rows sharing a primary key whose non-primary-key values differ,
    holding the left table's values, with the same schema as the left
    table."""

    changed_right: pl.LazyFrame
    """Same as `changed_left`, but holding the right table's values
    for those same primary keys, with the same schema as the right table."""

    @property
    def is_identical(self) -> bool:
        """Whether every shared-key row matches and no keys are one-sided."""
        return self.pk_diff.is_identical and not self.column_changes

    def cleanup(self) -> None:
        """Delete the temporary files that back the differing rows."""
        self.pk_diff.cleanup()

changed_left instance-attribute

Rows sharing a primary key whose non-primary-key values differ, holding the left table's values, with the same schema as the left table.

changed_right instance-attribute

Same as changed_left, but holding the right table's values for those same primary keys, with the same schema as the right table.

changed_row_count instance-attribute

Number of shared-primary-key rows with at least one differing non-primary-key value.

column_changes instance-attribute

Maps each non-primary-key column to the number of shared-primary-key rows where its value differs between the two tables. Columns with no changes are omitted.

is_identical property

Whether every shared-key row matches and no keys are one-sided.

pk_diff instance-attribute

Symmetric difference of primary key values: full rows present in only one of the two tables.

cleanup()

Delete the temporary files that back the differing rows.

Source code in src/pudl_diff/rows.py
def cleanup(self) -> None:
    """Delete the temporary files that back the differing rows."""
    self.pk_diff.cleanup()

RowSetDiff dataclass

The multiset symmetric difference of rows between two tables.

A row only counts as "the same" if the key columns match on both sides - every column, for tables with no primary key, or just the primary key columns, for tables that have one. For tables with no primary key, the number of copies of each row matters too: a row appearing 3 times on the left and once on the right contributes 2 rows to only_in_left.

The differing rows themselves are backed by temporary Parquet files, so they can be far larger than memory. The files are deleted by cleanup(), or else when this object (and any other object sharing spill_dir) is garbage collected; call .collect() on a frame to load it.

Source code in src/pudl_diff/rows.py
@dataclass(frozen=True)
class RowSetDiff:
    """The multiset symmetric difference of rows between two tables.

    A row only counts as "the same" if the key columns match on both sides -
    every column, for tables with no primary key, or just the primary key columns,
    for tables that have one. For tables with no primary key, the number of
    copies of each row matters too: a row appearing 3 times on the left and once
    on the right contributes 2 rows to `only_in_left`.

    The differing rows themselves are backed by temporary Parquet files, so they
    can be far larger than memory. The files are deleted by `cleanup()`, or else
    when this object (and any other object sharing `spill_dir`) is garbage
    collected; call `.collect()` on a frame to load it.
    """

    only_in_left: pl.LazyFrame
    only_in_right: pl.LazyFrame
    only_in_left_count: int
    only_in_right_count: int
    multiplicity_changed_row_count: int
    """Number of distinct rows present on both sides but a different number of
    times (only ever non-zero for tables with no primary key)."""
    spill_dir: SpillDir = field(repr=False, compare=False)

    @property
    def is_identical(self) -> bool:
        """Whether every row in one table has a matching row in the other."""
        return self.only_in_left_count == 0 and self.only_in_right_count == 0

    def cleanup(self) -> None:
        """Delete the temporary files that back the differing rows."""
        self.spill_dir.cleanup()

is_identical property

Whether every row in one table has a matching row in the other.

multiplicity_changed_row_count instance-attribute

Number of distinct rows present on both sides but a different number of times (only ever non-zero for tables with no primary key).

cleanup()

Delete the temporary files that back the differing rows.

Source code in src/pudl_diff/rows.py
def cleanup(self) -> None:
    """Delete the temporary files that back the differing rows."""
    self.spill_dir.cleanup()

SpillDir

A temporary directory holding the Parquet files of a comparison's rows.

Unlike tempfile.TemporaryDirectory, it is not a mistake to leave it to be cleaned up implicitly, so doing so doesn't warn: it is removed when cleanup() is called, or else when this object is garbage collected or the interpreter exits, whichever comes first.

Source code in src/pudl_diff/rows.py
class SpillDir:
    """A temporary directory holding the Parquet files of a comparison's rows.

    Unlike `tempfile.TemporaryDirectory`, it is not a mistake to leave it to be
    cleaned up implicitly, so doing so doesn't warn: it is removed when
    `cleanup()` is called, or else when this object is garbage collected or the
    interpreter exits, whichever comes first.
    """

    def __init__(self) -> None:
        """Make the directory."""
        self.name: str = tempfile.mkdtemp(prefix="pudl_diff_")
        self._finalizer = weakref.finalize(self, shutil.rmtree, self.name, True)

    def cleanup(self) -> None:
        """Remove the directory and everything in it, if that hasn't been done."""
        self._finalizer()

__init__()

Make the directory.

Source code in src/pudl_diff/rows.py
def __init__(self) -> None:
    """Make the directory."""
    self.name: str = tempfile.mkdtemp(prefix="pudl_diff_")
    self._finalizer = weakref.finalize(self, shutil.rmtree, self.name, True)

cleanup()

Remove the directory and everything in it, if that hasn't been done.

Source code in src/pudl_diff/rows.py
def cleanup(self) -> None:
    """Remove the directory and everything in it, if that hasn't been done."""
    self._finalizer()

compare_rows_with_pk(left, right, pk_cols, *, rtol=1e-05, atol=1e-08)

Compare two tables sharing a primary key.

Parameters:

Name Type Description Default
left LazyFrame

The "left" table to compare. Must have the same columns as right, with matching dtypes for the primary key columns.

required
right LazyFrame

The "right" table to compare against left.

required
pk_cols Sequence[str]

Names of the primary key columns shared by both tables, e.g. from PudlDiffDataset.primary_key(table_name).

required
rtol float

Relative tolerance used to treat two floating point values as equal, matching numpy.isclose()'s default.

1e-05
atol float

Absolute tolerance used to treat two floating point values as equal, matching numpy.isclose()'s default.

1e-08
Source code in src/pudl_diff/rows.py
def compare_rows_with_pk(
    left: pl.LazyFrame,
    right: pl.LazyFrame,
    pk_cols: Sequence[str],
    *,
    rtol: float = 1e-5,
    atol: float = 1e-8,
) -> KeyedRowDiff:
    """Compare two tables sharing a primary key.

    Args:
        left: The "left" table to compare. Must have the same columns as
            `right`, with matching dtypes for the primary key columns.
        right: The "right" table to compare against `left`.
        pk_cols: Names of the primary key columns shared by both tables, e.g.
            from `PudlDiffDataset.primary_key(table_name)`.
        rtol: Relative tolerance used to treat two floating point values as
            equal, matching `numpy.isclose()`'s default.
        atol: Absolute tolerance used to treat two floating point values as
            equal, matching `numpy.isclose()`'s default.
    """
    schema = left.collect_schema()
    columns = list(schema.keys())
    non_pk_cols = [name for name in schema if name not in pk_cols]
    left_cols = [f"{c}_left" for c in non_pk_cols]
    right_cols = [f"{c}_right" for c in non_pk_cols]

    spill_dir = SpillDir()
    delta = _diff_by_key_hash(
        left,
        right,
        schema,
        pk_cols,
        rtol,
        atol,
        spill_dir,
        value_columns=non_pk_cols,
        multiset=False,
    )

    # Only keys present on both sides whose non-key values hash differently can
    # have changed, so compare just those rows, exactly and with tolerance.
    on_both_sides = (pl.col(_LEFT_COUNT_COL) > 0) & (pl.col(_RIGHT_COUNT_COL) > 0)
    left_candidates = _spill(
        _semi_rows(delta.left_keyed, delta.affected, on_both_sides).select(columns),
        spill_dir,
        "left_candidates",
    )
    right_candidates = _spill(
        _semi_rows(delta.right_keyed, delta.affected, on_both_sides).select(columns),
        spill_dir,
        "right_candidates",
    )
    left_renamed = left_candidates.rename(
        dict(zip(non_pk_cols, left_cols, strict=True))
    )
    right_renamed = right_candidates.rename(
        dict(zip(non_pk_cols, right_cols, strict=True))
    )
    joined = left_renamed.join(right_renamed, on=list(pk_cols), how="inner")
    diff_flag_cols = [f"_pudl_diff_flag_{c}" for c in non_pk_cols]
    joined = joined.with_columns(
        [
            _values_differ(lc, rc, schema[c], rtol, atol).alias(flag_col)
            for c, lc, rc, flag_col in zip(
                non_pk_cols, left_cols, right_cols, diff_flag_cols, strict=True
            )
        ]
    )

    # Keep only the rows that really changed (both sides' values and the
    # per-column flags). Everything below - the counts and both output frames -
    # is derived from this small spilled result.
    any_diff = pl.any_horizontal(diff_flag_cols) if diff_flag_cols else pl.lit(False)
    changed = _spill(joined.filter(any_diff), spill_dir, "changed")

    counts = (
        changed.select([pl.col(fc).sum() for fc in diff_flag_cols]).collect(
            engine="streaming"
        )
        if non_pk_cols
        else pl.DataFrame()
    )
    column_changes = {
        c: count
        for c, fc in zip(non_pk_cols, diff_flag_cols, strict=True)
        if (count := counts[fc][0]) > 0
    }

    changed_left = (
        changed.select([*pk_cols, *left_cols])
        .rename(dict(zip(left_cols, non_pk_cols, strict=True)))
        .select(columns)
    )
    changed_right = (
        changed.select([*pk_cols, *right_cols])
        .rename(dict(zip(right_cols, non_pk_cols, strict=True)))
        .select(columns)
    )

    return KeyedRowDiff(
        pk_diff=delta.row_set_diff,
        column_changes=column_changes,
        changed_row_count=count_rows(changed),
        changed_left=changed_left,
        changed_right=changed_right,
    )

compare_rows_without_pk(left, right, *, rtol=1e-05, atol=1e-08)

Compare two tables with no primary key by taking the symmetric difference of rows.

The comparison is of multisets: a row that appears a different number of times in the two tables counts as a difference, and the surplus copies are reported as only in the table that has more of them.

Parameters:

Name Type Description Default
left LazyFrame

The "left" table to compare. Must have the same columns as right, with the same dtypes.

required
right LazyFrame

The "right" table to compare against left.

required
rtol float

Relative tolerance used to treat two floating point values as equal, matching numpy.isclose()'s default. Set to 0 along with atol to require exact float equality.

1e-05
atol float

Absolute tolerance used to treat two floating point values as equal, matching numpy.isclose()'s default.

1e-08

Raises:

Type Description
SchemaError

if a column's dtype differs between left and right.

Source code in src/pudl_diff/rows.py
def compare_rows_without_pk(
    left: pl.LazyFrame,
    right: pl.LazyFrame,
    *,
    rtol: float = 1e-5,
    atol: float = 1e-8,
) -> RowSetDiff:
    """Compare two tables with no primary key by taking the symmetric difference of rows.

    The comparison is of multisets: a row that appears a different number of times
    in the two tables counts as a difference, and the surplus copies are reported
    as only in the table that has more of them.

    Args:
        left: The "left" table to compare. Must have the same columns as
            `right`, with the same dtypes.
        right: The "right" table to compare against `left`.
        rtol: Relative tolerance used to treat two floating point values as
            equal, matching `numpy.isclose()`'s default. Set to `0` along
            with `atol` to require exact float equality.
        atol: Absolute tolerance used to treat two floating point values as
            equal, matching `numpy.isclose()`'s default.

    Raises:
        polars.exceptions.SchemaError: if a column's dtype differs between
            `left` and `right`.
    """
    schema = left.collect_schema()
    spill_dir = SpillDir()
    return _diff_by_key_hash(
        left, right, schema, schema.keys(), rtol, atol, spill_dir, multiset=True
    ).row_set_diff