Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,24 @@ adheres to [Semantic Versioning](https://semver.org/).
a dict, or every one a tuple). Such a column now scores like the equivalent
nested Arrow column. Columns that mix kinds, such as strings with numbers or
lists with scalars, are still flagged.
- `StreamingCleaner` no longer raises `TypeError: Cannot compare tz-naive and
tz-aware timestamps` when a datetime column is tz-naive in one batch and
tz-aware in another (including a naive batch followed by strings with mixed
UTC offsets). Running datetime statistics and time-series watermarks compare
in UTC, reading naive values as UTC, and keep the zone each value arrived in.
Time-series mode warns when a timestamp or event-time column changes
awareness between batches, and `state_` marks such columns `tz_mixed`.
- `cdc_profile` no longer reads small numbers such as row numbers `1..100` as
epoch seconds and reports a 56-year-old batch as passing. When the epoch unit
is inferred and the median event time falls before 1990-01-01, it reports an
`event_time_implausible` error and leaves `freshness_seconds` unset, without
evaluating lateness or ordering. An explicit `event_time_unit=` is trusted as
before, and epoch columns in s, ms, us and ns from 1990 on are unchanged.
Numeric epoch event times before 1990-01-01 now need an explicit
`event_time_unit`; without it such a batch fails with
`event_time_implausible`.
Time-series streaming warns about such columns when `timestamp_unit` is not
set.

## [2.1.0] - 2026-09-15

Expand Down
109 changes: 85 additions & 24 deletions src/freshdata/cdc.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,14 @@

import pandas as pd

from .streaming._timeseries import TIMESTAMP_UNITS, parse_timestamps, to_timedelta
from .streaming._timeseries import (
PLAUSIBLE_EPOCH_START,
TIMESTAMP_UNITS,
describe_implausible_epoch,
implausible_epoch,
parse_timestamps,
to_timedelta,
)

__all__ = ["CDCDefect", "CDCReport", "cdc_profile"]

Expand All @@ -42,6 +49,7 @@
"missing_key",
"invalid_operation",
"missing_event_time",
"event_time_implausible",
"replay_risk",
]
Level = Literal["info", "warning", "error"]
Expand Down Expand Up @@ -157,8 +165,11 @@ def _sample_keys(
return tuple(str(v) for v in vals)


def _parse_event_time(raw: pd.Series, unit: str | None = None) -> pd.Series:
"""Parse the event-time column; unparseable values become ``NaT``.
def _parse_event_time(raw: pd.Series, unit: str | None = None) -> tuple[pd.Series, str | None]:
"""Parse the event-time column; return ``(parsed, epoch_unit)``.

Unparseable values become ``NaT``. ``epoch_unit`` is the unit numeric
values were read in, or ``None`` for string and datetime columns.

Datetime columns pass through. Mixed UTC offsets (e.g. across a DST change)
or a mix of naive and offset-aware values are normalised to UTC, reading
Expand All @@ -167,8 +178,7 @@ def _parse_event_time(raw: pd.Series, unit: str | None = None) -> pd.Series:
when *unit* is ``None``, as ``fd.clean_timeseries`` does
(see :func:`~freshdata.streaming._timeseries.parse_timestamps`).
"""
parsed, _ = parse_timestamps(raw, unit)
return parsed
return parse_timestamps(raw, unit)


def _to_event_tz(value: object, tz: Any) -> pd.Timestamp:
Expand Down Expand Up @@ -238,6 +248,30 @@ def _ordering_defects(
return {"out_of_order": n_ooo, "late": n_late}


def _watermark_defects(
df: pd.DataFrame,
ts: pd.Series,
key: str | None,
wm: pd.Timestamp,
defects: list[CDCDefect],
) -> dict[str, int]:
"""Detect records strictly before an explicit watermark (already in the events' zone)."""
late_mask = ts < wm
n_late = int(late_mask.sum())
if n_late:
defects.append(
CDCDefect(
kind="late",
level="error",
n_rows=n_late,
rationale=f"event time before the watermark {wm}",
sample_keys=_sample_keys(late_mask, df, key),
details={"watermark": str(wm)},
)
)
return {"out_of_order": 0, "late": n_late}


def _missing_key_defect(df: pd.DataFrame, key: str, defects: list[CDCDefect]) -> None:
"""Flag rows whose CDC key is null (a warning; the trust penalty is unchanged)."""
null_key_mask = df[key].isna()
Expand All @@ -255,6 +289,32 @@ def _missing_key_defect(df: pd.DataFrame, key: str, defects: list[CDCDefect]) ->
)


def _implausible_event_time_defect(
raw: pd.Series, name: str, unit: str, n_present: int, defects: list[CDCDefect]
) -> bool:
"""Flag numeric event times whose *inferred* epoch unit gives implausible
instants (see :func:`~freshdata.streaming._timeseries.implausible_epoch`)."""
median = implausible_epoch(raw, unit)
if median is None:
return False
defects.append(
CDCDefect(
kind="event_time_implausible",
level="error",
n_rows=n_present,
rationale=f"{name!r} "
+ describe_implausible_epoch(median, unit, "event_time_unit")
+ "; freshness, lateness and ordering were not evaluated",
details={
"inferred_unit": unit,
"median": median,
"plausible_from": str(PLAUSIBLE_EPOCH_START.date()),
},
)
)
return True


def _freshness(
ts: pd.Series,
now: pd.Timestamp | None,
Expand Down Expand Up @@ -343,7 +403,11 @@ def cdc_profile(
Epoch unit (``"s"``, ``"ms"``, ``"us"`` or ``"ns"``) of a *numeric*
``event_time`` column, such as Debezium's ``ts_ms``. ``None`` (default)
infers it from the magnitude of the values, as ``fd.clean_timeseries``
does. Ignored for string and datetime columns.
does. Ignored for string and datetime columns. When the unit is
inferred and the median value falls before 1990-01-01 in it (row
numbers or offsets such as ``1..100``), an ``event_time_implausible``
error is reported and freshness, lateness and ordering are not
evaluated; pass the unit explicitly to read such values as epochs.

Returns
-------
Expand All @@ -366,7 +430,7 @@ def cdc_profile(
if stale_after_td is not None and stale_after_td < pd.Timedelta(0):
raise ValueError(f"stale_after must not be negative, got {stale_after!r}")

ts = _parse_event_time(df[event_time], event_time_unit)
ts, epoch_unit = _parse_event_time(df[event_time], event_time_unit)
event_tz = ts.dt.tz
now_ts = _to_event_tz(now, event_tz) if now is not None else None

Expand All @@ -386,23 +450,20 @@ def cdc_profile(
)
)

# 1b) numeric event times that are not plausible epochs (row numbers, offsets).
# Freshness, lateness and ordering in seconds would be meaningless, so they
# are not evaluated, and the error keeps the gate from passing unverified.
implausible = epoch_unit is not None and event_time_unit is None and (
_implausible_event_time_defect(
df[event_time], event_time, epoch_unit, n_rows - n_missing, defects
)
)

# 2) ordering: out-of-order + late (relative to watermark / running watermark).
if watermark is not None:
wm = _to_event_tz(watermark, event_tz)
late_mask = ts < wm
n_late = int(late_mask.sum())
counts = {"out_of_order": 0, "late": n_late}
if n_late:
defects.append(
CDCDefect(
kind="late",
level="error",
n_rows=n_late,
rationale=f"event time before the watermark {wm}",
sample_keys=_sample_keys(late_mask, df, key),
details={"watermark": str(wm)},
)
)
if implausible:
counts = {"out_of_order": 0, "late": 0}
elif watermark is not None:
counts = _watermark_defects(df, ts, key, _to_event_tz(watermark, event_tz), defects)
else:
counts = _ordering_defects(df, ts, key, lateness_td, defects)

Expand Down Expand Up @@ -444,7 +505,7 @@ def cdc_profile(
)

# 5) freshness / stale batch.
freshness_seconds = _freshness(ts, now_ts, stale_after_td, defects)
freshness_seconds = None if implausible else _freshness(ts, now_ts, stale_after_td, defects)

# 6) replay-risk batch heuristic. A replay re-delivers changes, so at least
# one duplicate-key row is required; late rows alone never trigger it.
Expand Down
22 changes: 19 additions & 3 deletions src/freshdata/streaming/_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import pandas as pd
from pandas.api.types import is_bool_dtype, is_datetime64_any_dtype, is_numeric_dtype

from ..fieldcheck import _as_utc
from ._stats import BoundedCounter, ReservoirSampler, Welford

#: Cap on the retained trust-score history (keeps the state bounded on long streams).
Expand Down Expand Up @@ -86,7 +87,10 @@ def __init__(self, name: str, role: str, baseline_dtype: str, first_seen_batch:
self._dt_min: pd.Timestamp | None = None
self._dt_max: pd.Timestamp | None = None
self._dt_prev_max: pd.Timestamp | None = None
self._dt_tz_aware: bool | None = None # awareness of the latest batch
self.datetime_ordered = True # until a batch proves otherwise
#: True once both tz-naive and tz-aware batches were seen for this column.
self.datetime_tz_mixed = False

@property
def seen(self) -> int:
Expand Down Expand Up @@ -126,11 +130,21 @@ def _update_datetime(self, s: pd.Series) -> None:
if nonnull.empty:
return
batch_min, batch_max = nonnull.min(), nonnull.max()
self._dt_min = batch_min if self._dt_min is None else min(self._dt_min, batch_min)
self._dt_max = batch_max if self._dt_max is None else max(self._dt_max, batch_max)
# Batches may differ in tz-awareness (a naive batch, then an offset-aware
# one). Compare in UTC, reading naive values as UTC, and keep each stored
# bound in the zone it arrived in.
aware = batch_min.tzinfo is not None
if self._dt_tz_aware is not None and aware != self._dt_tz_aware:
self.datetime_tz_mixed = True
self._dt_tz_aware = aware
self._dt_min = (batch_min if self._dt_min is None
else min(self._dt_min, batch_min, key=_as_utc))
self._dt_max = (batch_max if self._dt_max is None
else max(self._dt_max, batch_max, key=_as_utc))
# Globally ordered iff each batch starts no earlier than the previous ended
# and is itself internally non-decreasing.
if self._dt_prev_max is not None and batch_min < self._dt_prev_max:
if (self._dt_prev_max is not None
and _as_utc(batch_min) < _as_utc(self._dt_prev_max)):
self.datetime_ordered = False
if not nonnull.is_monotonic_increasing:
self.datetime_ordered = False
Expand Down Expand Up @@ -188,6 +202,8 @@ def to_dict(self) -> dict[str, Any]:
"max": str(self._dt_max),
"ordered": self.datetime_ordered,
}
if self.datetime_tz_mixed:
payload["datetime"]["tz_mixed"] = True
return payload


Expand Down
Loading
Loading