fix(ees): newest KS4 release owns every year it contains (C2)
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
0e987ef06e
commit
4db1131d0f
3 files changed
+142
No files matched your search
@@ -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
|
||||
@@ -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),
|
||||
|
||||
Reference in new issue
Block a user