From e236669fdec531d25c023180c608042a8b10f90d Mon Sep 17 00:00:00 2001 From: Tudor Date: Mon, 31 Aug 2026 21:36:43 +0100 Subject: [PATCH] fix(destinations): school rows and the England reference are different grains MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The annual DAG died with a BrokenPipeError from Meltano's log writer, which is several frames from the cause: target-postgres exited first and the tap saw its stdout close. The tap declared primary_keys = [urn, ...] while emitting urn=None for the national rows, and target-postgres turns primary_keys into a NOT NULL constraint. The first national row of the run failed the insert and took the loader with it. Every other tap in this repo keys on non-null columns. Carrying two grains in one stream was the actual mistake, so the fix is to separate them rather than paper over the null: four streams now, with ees_ks4/ks5_destinations_national carrying no urn column at all — a school identifier that is null in every row is a grain mismatch, not a column. The staging models split the same way and the national mart reads the new pair instead of filtering `where urn is null`. Verified against the live API: the school stream yields 135,240 rows over 4,508 schools with no duplicate keys, no null key columns and all 31,382 suppression sentinels intact; the national streams yield 30 and 33 rows with no urn column. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01BvdDKvFFSZuMVDH5fEyTob --- pipeline/dags/school_data_pipeline.py | 2 +- .../tap_uk_ees_destinations/tap.py | 92 +++++++++++++++---- .../tap-uk-ees-destinations/tests/test_tap.py | 65 +++++++++++++ .../marts/fact_destination_national.sql | 6 +- .../transform/models/staging/_stg_sources.yml | 11 +++ .../staging/stg_ees_ks4_destinations.sql | 11 ++- .../stg_ees_ks4_destinations_national.sql | 38 ++++++++ .../staging/stg_ees_ks5_destinations.sql | 13 +-- .../stg_ees_ks5_destinations_national.sql | 38 ++++++++ 9 files changed, 242 insertions(+), 34 deletions(-) create mode 100644 pipeline/transform/models/staging/stg_ees_ks4_destinations_national.sql create mode 100644 pipeline/transform/models/staging/stg_ees_ks5_destinations_national.sql diff --git a/pipeline/dags/school_data_pipeline.py b/pipeline/dags/school_data_pipeline.py index 990ab79..a8b5a2a 100644 --- a/pipeline/dags/school_data_pipeline.py +++ b/pipeline/dags/school_data_pipeline.py @@ -190,7 +190,7 @@ with DAG( dbt_build_ees = BashOperator( task_id="dbt_build", - bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ees_ks2+ stg_legacy_ks2+ stg_ees_ks4+ stg_legacy_ks4+ stg_ees_census+ stg_ees_admissions+ stg_ees_ks2_national+ stg_ees_ks4_national+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+", + bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ees_ks2+ stg_legacy_ks2+ stg_ees_ks4+ stg_legacy_ks4+ stg_ees_census+ stg_ees_admissions+ stg_ees_ks2_national+ stg_ees_ks4_national+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+ stg_ees_ks4_destinations_national+ stg_ees_ks5_destinations_national+", ) sync_typesense_ees = BashOperator( diff --git a/pipeline/plugins/extractors/tap-uk-ees-destinations/tap_uk_ees_destinations/tap.py b/pipeline/plugins/extractors/tap-uk-ees-destinations/tap_uk_ees_destinations/tap.py index c502fe5..bdb7aed 100644 --- a/pipeline/plugins/extractors/tap-uk-ees-destinations/tap_uk_ees_destinations/tap.py +++ b/pipeline/plugins/extractors/tap-uk-ees-destinations/tap_uk_ees_destinations/tap.py @@ -188,21 +188,25 @@ class DestinationsStream(Stream): _national_pinned: list[str] = [] _indicators: dict[str, str] = {} - schema = th.PropertiesList( - th.Property("urn", th.StringType), + # School rows and the England reference are DIFFERENT GRAINS, so they are + # different streams. Carrying both in one table meant a null `urn` inside + # the primary key, which target-postgres turns into a NOT NULL constraint: + # the first national row killed the loader mid-run, and the tap saw only a + # BrokenPipeError on its stdout. + _MEASURE_PROPERTIES = ( th.Property("time_period", th.StringType), th.Property("pupil_group", th.StringType), th.Property("destination_measure", th.StringType), th.Property("cohort_pupils", th.StringType), th.Property("pupils_raw", th.StringType), th.Property("percentage_raw", th.StringType), - ).to_dict() + ) + _MEASURE_KEYS = ["time_period", "pupil_group", "destination_measure"] - primary_keys = ["urn", "time_period", "pupil_group", "destination_measure"] replication_key = None def _query(self, period: str, pinned: list[str], keep_level: str, - urn_by_location: dict[str, str]): + urn_by_location: dict[str, str]): # noqa: D401 """Page one period of one geographic level, yielding Singer records.""" criteria = [ {"filters": {"in": list(self._destination_slugs)}}, @@ -248,19 +252,51 @@ class DestinationsStream(Stream): "%s: %d school locations, %d time periods", self.name, len(urn_by_location), len(periods), ) - - # Two queries per period, because school rows and the England reference - # need different establishment pins — see KS4_NATIONAL_PINNED. The API - # cannot filter by geographic level, so each pass keeps its own and - # discards the local-authority, district, regional and constituency - # rows that come with them. for period in periods: - yield from self._query(period, self._pinned, "SCH", urn_by_location) - yield from self._query(period, self._national_pinned, "NAT", urn_by_location) + yield from self._query( + period, self._pinned_for_level, self._level, urn_by_location, + ) -class KS4DestinationsStream(DestinationsStream): - name = "ees_ks4_destinations" +class SchoolDestinationsStream(DestinationsStream): + """School-level rows. `urn` is part of the key and is never null.""" + + _level = "SCH" + + @property + def _pinned_for_level(self): + return self._pinned + + schema = th.PropertiesList( + th.Property("urn", th.StringType), + *DestinationsStream._MEASURE_PROPERTIES, + ).to_dict() + + primary_keys = ["urn", *DestinationsStream._MEASURE_KEYS] + + +class NationalDestinationsStream(DestinationsStream): + """The England reference. No `urn` column at all — a school identifier that + is always null is not a column, it is a grain mismatch.""" + + _level = "NAT" + + @property + def _pinned_for_level(self): + return self._national_pinned + + schema = th.PropertiesList(*DestinationsStream._MEASURE_PROPERTIES).to_dict() + + primary_keys = list(DestinationsStream._MEASURE_KEYS) + + def post_process(self, row, context=None): + # row_to_record emits urn=None for national rows; drop the key rather + # than ship a column that is null in every row. + row.pop("urn", None) + return row + + +class _KS4Config(DestinationsStream): _dataset_id = KS4_DATASET _destination_slugs = KS4_DESTINATION_SLUGS _pupil_group_slugs = KS4_PUPIL_GROUP_SLUGS @@ -269,8 +305,7 @@ class KS4DestinationsStream(DestinationsStream): _indicators = KS4_INDICATORS -class KS5DestinationsStream(DestinationsStream): - name = "ees_ks5_destinations" +class _KS5Config(DestinationsStream): _dataset_id = KS5_DATASET _destination_slugs = KS5_DESTINATION_SLUGS _pupil_group_slugs = KS5_PUPIL_GROUP_SLUGS @@ -279,12 +314,33 @@ class KS5DestinationsStream(DestinationsStream): _indicators = KS5_INDICATORS +class KS4DestinationsStream(SchoolDestinationsStream, _KS4Config): + name = "ees_ks4_destinations" + + +class KS5DestinationsStream(SchoolDestinationsStream, _KS5Config): + name = "ees_ks5_destinations" + + +class KS4NationalDestinationsStream(NationalDestinationsStream, _KS4Config): + name = "ees_ks4_destinations_national" + + +class KS5NationalDestinationsStream(NationalDestinationsStream, _KS5Config): + name = "ees_ks5_destinations_national" + + class TapUKEESDestinations(Tap): name = "tap-uk-ees-destinations" config_jsonschema = th.PropertiesList().to_dict() def discover_streams(self): - return [KS4DestinationsStream(self), KS5DestinationsStream(self)] + return [ + KS4DestinationsStream(self), + KS5DestinationsStream(self), + KS4NationalDestinationsStream(self), + KS5NationalDestinationsStream(self), + ] if __name__ == "__main__": diff --git a/pipeline/plugins/extractors/tap-uk-ees-destinations/tests/test_tap.py b/pipeline/plugins/extractors/tap-uk-ees-destinations/tests/test_tap.py index 07da87c..c17bb20 100644 --- a/pipeline/plugins/extractors/tap-uk-ees-destinations/tests/test_tap.py +++ b/pipeline/plugins/extractors/tap-uk-ees-destinations/tests/test_tap.py @@ -132,3 +132,68 @@ def test_establishment_dimensions_are_not_pinned(): establishment_totals = {"4369U", "rgHcN", "EfHQq", "S4ROV"} assert not set(KS4_PINNED) & establishment_totals assert not set(KS5_PINNED) & establishment_totals + + +# ── Grain separation ──────────────────────────────────────────────────────── +# +# School rows and the England reference were originally one stream with a +# nullable `urn` in the primary key. target-postgres turns primary_keys into a +# NOT NULL constraint, so the first national row killed the loader mid-run and +# the tap saw only a BrokenPipeError on its stdout — a symptom several frames +# away from the cause. + +def _streams(): + from tap_uk_ees_destinations.tap import TapUKEESDestinations + return TapUKEESDestinations(config={}, validate_config=False).discover_streams() + + +def test_no_stream_has_a_nullable_primary_key_column(): + for stream in _streams(): + props = stream.schema["properties"] + for key in stream.primary_keys: + assert key in props, f"{stream.name}: key {key} is not in the schema" + if "urn" in stream.primary_keys: + assert stream._level == "SCH", ( + f"{stream.name} keys on urn but does not emit school rows" + ) + + +def test_national_streams_carry_no_urn_column_at_all(): + for stream in _streams(): + if not stream.name.endswith("_national"): + continue + assert "urn" not in stream.schema["properties"], ( + "a school identifier that is null in every row is a grain " + "mismatch, not a column" + ) + assert "urn" not in stream.primary_keys + + +def test_national_post_process_drops_the_null_urn(): + from tap_uk_ees_destinations.tap import KS4NationalDestinationsStream, TapUKEESDestinations + tap = TapUKEESDestinations(config={}, validate_config=False) + stream = KS4NationalDestinationsStream(tap) + row = {"urn": None, "time_period": "202223", "pupil_group": "all", + "destination_measure": "school_sixth_form", "cohort_pupils": "1", + "pupils_raw": "1", "percentage_raw": "1"} + assert "urn" not in stream.post_process(dict(row)) + + +def test_school_and_national_streams_exist_for_both_phases(): + names = {s.name for s in _streams()} + assert names == { + "ees_ks4_destinations", "ees_ks5_destinations", + "ees_ks4_destinations_national", "ees_ks5_destinations_national", + } + + +def test_each_stream_queries_its_own_geographic_level_with_its_own_pins(): + """The national pass needs establishment pinned to Total; the school pass + must not pin it at all, or every school returns zero rows.""" + for stream in _streams(): + if stream.name.endswith("_national"): + assert stream._level == "NAT" + assert stream._pinned_for_level is stream._national_pinned + else: + assert stream._level == "SCH" + assert stream._pinned_for_level is stream._pinned diff --git a/pipeline/transform/models/marts/fact_destination_national.sql b/pipeline/transform/models/marts/fact_destination_national.sql index 811078f..34cb3b8 100644 --- a/pipeline/transform/models/marts/fact_destination_national.sql +++ b/pipeline/transform/models/marts/fact_destination_national.sql @@ -22,8 +22,7 @@ select pupils, percentage, status -from {{ ref('stg_ees_ks4_destinations') }} -where urn is null +from {{ ref('stg_ees_ks4_destinations_national') }} union all @@ -36,5 +35,4 @@ select pupils, percentage, status -from {{ ref('stg_ees_ks5_destinations') }} -where urn is null +from {{ ref('stg_ees_ks5_destinations_national') }} diff --git a/pipeline/transform/models/staging/_stg_sources.yml b/pipeline/transform/models/staging/_stg_sources.yml index 564927b..2d5880a 100644 --- a/pipeline/transform/models/staging/_stg_sources.yml +++ b/pipeline/transform/models/staging/_stg_sources.yml @@ -47,6 +47,17 @@ sources: 16-18 study leavers destinations. Same grain, same suppression caveat, and only institutions with post-16 provision appear. + - name: ees_ks4_destinations_national + description: > + England KS4 destination measures by pupil group. A separate table + from ees_ks4_destinations because it is a separate grain — no school, + so no urn column. Same 'c' suppression caveat. + + - name: ees_ks5_destinations_national + description: > + England 16-18 destination measures by pupil group. Same grain and + caveat as ees_ks4_destinations_national. + - name: ees_ks4_performance description: KS4 performance tables (long format — one row per school × breakdown × sex) diff --git a/pipeline/transform/models/staging/stg_ees_ks4_destinations.sql b/pipeline/transform/models/staging/stg_ees_ks4_destinations.sql index 95d1ca4..62b6d07 100644 --- a/pipeline/transform/models/staging/stg_ees_ks4_destinations.sql +++ b/pipeline/transform/models/staging/stg_ees_ks4_destinations.sql @@ -1,7 +1,9 @@ {{ config(materialized='table') }} --- Staging model: KS4 leavers destinations, school level plus the England --- reference (which carries a null urn). +-- Staging model: KS4 leavers destinations, school level. +-- +-- School rows only. The England reference is a different grain and lives in +-- stg_ees_ks4_destinations_national. -- -- DELIBERATELY DOES NOT USE safe_numeric. That macro maps every EES sentinel -- (z, c, x, q, u) to NULL, which is right for attainment — there, "suppressed" @@ -14,13 +16,12 @@ with source as ( select * from {{ source('raw', 'ees_ks4_destinations') }} - -- National rows carry a null urn and feed fact_destination_national. - where (urn is null or urn = '' or urn ~ '^[0-9]+$') + where urn ~ '^[0-9]+$' and time_period ~ '^[0-9]+$' ) select - case when urn ~ '^[0-9]+$' then cast(trim(urn) as integer) end as urn, + cast(trim(urn) as integer) as urn, cast(trim(time_period) as integer) as year, trim(pupil_group) as pupil_group, trim(destination_measure) as destination_measure, diff --git a/pipeline/transform/models/staging/stg_ees_ks4_destinations_national.sql b/pipeline/transform/models/staging/stg_ees_ks4_destinations_national.sql new file mode 100644 index 0000000..7e938f6 --- /dev/null +++ b/pipeline/transform/models/staging/stg_ees_ks4_destinations_national.sql @@ -0,0 +1,38 @@ +{{ config(materialized='table') }} + +-- Staging model: England KS4 destination measures — the national +-- reference the school sections compare against. +-- +-- A separate model because it is a separate grain: there is no school here, and +-- carrying these rows in the school table meant a null urn inside the primary +-- key, which the Postgres loader rejects. +-- +-- DELIBERATELY DOES NOT USE safe_numeric, for the same reason as the school +-- model: 'suppressed' and 'not applicable' are different claims. + +with source as ( + select * from {{ source('raw', 'ees_ks4_destinations_national') }} + where time_period ~ '^[0-9]+$' +) + +select + cast(trim(time_period) as integer) as year, + trim(pupil_group) as pupil_group, + trim(destination_measure) as destination_measure, + + case when cohort_pupils ~ '^[0-9]+$' + then cast(cohort_pupils as integer) end as cohort_pupils, + + case when pupils_raw ~ '^[0-9]+$' + then cast(pupils_raw as integer) end as pupils, + + case when percentage_raw ~ '^-?[0-9]+(\.[0-9]+)?$' + then cast(percentage_raw as numeric) end as percentage, + + case + when pupils_raw ~ '^[0-9]+$' then 'published' + when lower(trim(pupils_raw)) = 'c' then 'suppressed' + else 'not_applicable' + end as status + +from source diff --git a/pipeline/transform/models/staging/stg_ees_ks5_destinations.sql b/pipeline/transform/models/staging/stg_ees_ks5_destinations.sql index bce4202..7a71752 100644 --- a/pipeline/transform/models/staging/stg_ees_ks5_destinations.sql +++ b/pipeline/transform/models/staging/stg_ees_ks5_destinations.sql @@ -1,8 +1,10 @@ {{ config(materialized='table') }} --- Staging model: 16-18 study leavers destinations, institution level plus the --- England reference (which carries a null urn). Only sixth forms and colleges --- appear here, so a secondary with no post-16 provision has no rows at all. +-- Staging model: 16-18 study leavers destinations, institution level. +-- +-- Only sixth forms and colleges appear, so a secondary with no post-16 +-- provision has no rows at all. The England reference is a different grain and +-- lives in stg_ees_ks5_destinations_national. -- -- DELIBERATELY DOES NOT USE safe_numeric. That macro maps every EES sentinel -- (z, c, x, q, u) to NULL, which is right for attainment — there, "suppressed" @@ -15,13 +17,12 @@ with source as ( select * from {{ source('raw', 'ees_ks5_destinations') }} - -- National rows carry a null urn and feed fact_destination_national. - where (urn is null or urn = '' or urn ~ '^[0-9]+$') + where urn ~ '^[0-9]+$' and time_period ~ '^[0-9]+$' ) select - case when urn ~ '^[0-9]+$' then cast(trim(urn) as integer) end as urn, + cast(trim(urn) as integer) as urn, cast(trim(time_period) as integer) as year, trim(pupil_group) as pupil_group, trim(destination_measure) as destination_measure, diff --git a/pipeline/transform/models/staging/stg_ees_ks5_destinations_national.sql b/pipeline/transform/models/staging/stg_ees_ks5_destinations_national.sql new file mode 100644 index 0000000..a31a537 --- /dev/null +++ b/pipeline/transform/models/staging/stg_ees_ks5_destinations_national.sql @@ -0,0 +1,38 @@ +{{ config(materialized='table') }} + +-- Staging model: England KS5 destination measures — the national +-- reference the school sections compare against. +-- +-- A separate model because it is a separate grain: there is no school here, and +-- carrying these rows in the school table meant a null urn inside the primary +-- key, which the Postgres loader rejects. +-- +-- DELIBERATELY DOES NOT USE safe_numeric, for the same reason as the school +-- model: 'suppressed' and 'not applicable' are different claims. + +with source as ( + select * from {{ source('raw', 'ees_ks5_destinations_national') }} + where time_period ~ '^[0-9]+$' +) + +select + cast(trim(time_period) as integer) as year, + trim(pupil_group) as pupil_group, + trim(destination_measure) as destination_measure, + + case when cohort_pupils ~ '^[0-9]+$' + then cast(cohort_pupils as integer) end as cohort_pupils, + + case when pupils_raw ~ '^[0-9]+$' + then cast(pupils_raw as integer) end as pupils, + + case when percentage_raw ~ '^-?[0-9]+(\.[0-9]+)?$' + then cast(percentage_raw as numeric) end as percentage, + + case + when pupils_raw ~ '^[0-9]+$' then 'published' + when lower(trim(pupils_raw)) = 'c' then 'suppressed' + else 'not_applicable' + end as status + +from source -- 2.54.0