diff --git a/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/release_precedence.py b/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/release_precedence.py new file mode 100644 index 0000000..52080c6 --- /dev/null +++ b/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/release_precedence.py @@ -0,0 +1,44 @@ +"""Which release a row comes from when DfE re-publishes a year. + +DfE re-publishes earlier years inside later releases: the 2024/25 KS4 results +file holds 2022/23, 2023/24 and 2024/25. The 2023/24 release's own file, +re-issued on 10 March 2026 under older column names, was read after it and +overwrote every 2023/24 row with blanks (audit C2). A stream that opts in +treats the newest release as the authority for every year it contains. + +Free of the Singer SDK so CI's pytest, which installs only the backend's +requirements, can load it. +""" +from __future__ import annotations + +import pandas as pd + + +def newest_first(releases: list[dict]) -> list[dict]: + """Releases by time_period, newest first. A release whose time_period is + unknown goes last, in the order given.""" + dated = [r for r in releases if r.get("time_period")] + undated = [r for r in releases if not r.get("time_period")] + return sorted(dated, key=lambda r: r["time_period"], reverse=True) + undated + + +def periods_in(df: pd.DataFrame) -> set[str]: + """The years a release's rows cover.""" + if "time_period" not in df.columns: + return set() + return set(df["time_period"].astype(str).str.strip()) - {""} + + +def drop_owned_periods( + df: pd.DataFrame, owned: set[str] +) -> tuple[pd.DataFrame, dict[str, int]]: + """Drop the rows for years a newer release already supplied. + + Returns the rows kept and, for each year dropped, how many rows went. + """ + if "time_period" not in df.columns or not owned: + return df, {} + periods = df["time_period"].astype(str).str.strip() + dropped = periods.isin(owned) + skipped = {str(k): int(v) for k, v in periods[dropped].value_counts().items()} + return df[~dropped], skipped diff --git a/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/tap.py b/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/tap.py index b8bdf53..1b399df 100644 --- a/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/tap.py +++ b/pipeline/plugins/extractors/tap-uk-ees/tap_uk_ees/tap.py @@ -17,6 +17,8 @@ import requests from singer_sdk import Stream, Tap from singer_sdk import typing as th +from tap_uk_ees.release_precedence import drop_owned_periods, newest_first, periods_in + CONTENT_API_BASE = ( "https://content.explore-education-statistics.service.gov.uk/api" ) @@ -88,6 +90,8 @@ class EESDatasetStream(Stream): target CSV path inside the ZIP (substring match, not exact). Subclasses may set _column_renames to map messy CSV column names to clean Singer field names before yielding records. + Subclasses may set _newest_release_owns_period when DfE re-publishes + earlier years in later releases and the newest copy is the authority. """ replication_key = None @@ -96,6 +100,7 @@ class EESDatasetStream(Stream): _urn_column: str = "school_urn" # column name for URN in the CSV _encoding: str = "utf-8" # CSV file encoding (some DfE files use latin-1) _column_renames: dict = {} # CSV column name → Singer field name + _newest_release_owns_period: bool = False # see release_precedence.py def get_records(self, context): import pandas as pd @@ -110,6 +115,9 @@ class EESDatasetStream(Stream): self.logger.info( "Found %d release(s) for %s", len(releases), self._publication_slug ) + if self._newest_release_owns_period: + releases = newest_first(releases) + owned_periods: set[str] = set() for release in releases: release_id = release["id"] @@ -163,6 +171,15 @@ class EESDatasetStream(Stream): if urn_col in df.columns: df = df[df[urn_col].notna() & (df[urn_col] != "")] + if self._newest_release_owns_period: + df, skipped = drop_owned_periods(df, owned_periods) + for period, count in sorted(skipped.items()): + self.logger.info( + "Skipping %d rows for %s from release %s: a newer release supplied that year", + count, period, release_id, + ) + owned_periods |= periods_in(df) + self.logger.info("Emitting %d school-level rows from release %s", len(df), release_id) for _, row in df.iterrows(): @@ -251,6 +268,9 @@ class EESKS4PerformanceStream(EESDatasetStream): primary_keys = ["school_urn", "time_period", "breakdown_topic", "breakdown", "sex"] _publication_slug = "key-stage-4-performance" _target_filename = "performance_tables_schools" + # DfE's 2024/25 file re-publishes 2022/23 and 2023/24 under current names; + # the 2023/24 release's own file uses older ones (audit C2). + _newest_release_owns_period = True schema = th.PropertiesList( th.Property("time_period", th.StringType, required=True), th.Property("school_urn", th.StringType, required=True), diff --git a/pipeline/tests/test_ees_release_precedence.py b/pipeline/tests/test_ees_release_precedence.py new file mode 100644 index 0000000..09ce3ae --- /dev/null +++ b/pipeline/tests/test_ees_release_precedence.py @@ -0,0 +1,78 @@ +"""DfE re-publishes earlier years inside later KS4 releases. + +The 2024/25 results file holds 2022/23, 2023/24 and 2024/25 under current +column names. The 2023/24 release's own file, re-issued in March 2026 under +older names, was read after it and overwrote every 2023/24 row with blanks +(audit C2). For the KS4 results stream the newest release owns every year it +contains. +""" +import importlib.util +import re +from pathlib import Path + +import pandas as pd +import pytest + +TAP_DIR = (Path(__file__).resolve().parents[1] / 'plugins' / 'extractors' / 'tap-uk-ees' + / 'tap_uk_ees') + + +@pytest.fixture +def precedence(): + spec = importlib.util.spec_from_file_location( + 'release_precedence', TAP_DIR / 'release_precedence.py') + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def _release(period): + return {'id': f'release-{period}', 'time_period': period} + + +def test_releases_are_taken_newest_first_whatever_order_the_api_gives(precedence): + releases = [_release('202223'), _release('202425'), _release(None), _release('202324')] + ordered = precedence.newest_first(releases) + assert [r['time_period'] for r in ordered] == ['202425', '202324', '202223', None] + + +def test_a_year_a_newer_release_supplied_is_dropped_from_an_older_one(precedence): + newer = pd.DataFrame({'time_period': ['202425', '202324', '202223'], 'school_urn': ['1'] * 3}) + older = pd.DataFrame({'time_period': ['202324', '202324', '201920'], 'school_urn': ['1', '2', '1']}) + + kept, skipped = precedence.drop_owned_periods(older, precedence.periods_in(newer)) + + assert list(kept['time_period']) == ['201920'] + assert skipped == {'202324': 2} + + +def test_periods_match_despite_surrounding_spaces(precedence): + owned = precedence.periods_in(pd.DataFrame({'time_period': [' 202324 ']})) + kept, skipped = precedence.drop_owned_periods(pd.DataFrame({'time_period': ['202324']}), owned) + assert owned == {'202324'} + assert kept.empty + assert skipped == {'202324': 1} + + +def test_nothing_is_dropped_before_any_year_is_owned(precedence): + df = pd.DataFrame({'time_period': ['202324'], 'school_urn': ['1']}) + kept, skipped = precedence.drop_owned_periods(df, set()) + assert kept.equals(df) + assert skipped == {} + + +def test_a_file_without_time_period_is_left_alone(precedence): + df = pd.DataFrame({'school_urn': ['1']}) + kept, skipped = precedence.drop_owned_periods(df, {'202324'}) + assert kept.equals(df) + assert skipped == {} + assert precedence.periods_in(df) == set() + + +def test_only_the_ks4_results_stream_opts_in(): + # A general rule would wipe KS2: the 2024/25 KS2 file holds 98,448 of the + # 955,956 rows the 2023/24 release has for 2023/24. + source = (TAP_DIR / 'tap.py').read_text() + opted_in = [chunk.split('(')[0] for chunk in source.split('\nclass ')[1:] + if re.search(r'_newest_release_owns_period\s*=\s*True', chunk)] + assert opted_in == ['EESKS4PerformanceStream']