Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c26b65246f | ||
|
|
3728a63275 | ||
|
|
b6e48c4930 | ||
|
|
9f4f2507cc | ||
|
|
94bfac9caf | ||
|
|
65a2619e1d |
No files matched your search
@@ -106,9 +106,12 @@ print(f'Validation passed: {{count}} GIAS rows')
|
|||||||
""",
|
""",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Marts fed by annual EES staging models are rebuilt by the EES DAG, even
|
||||||
|
# when they join dim_school. Selecting them here fails in any database
|
||||||
|
# where that DAG hasn't run (pipeline/tests/test_dag_selectors.py).
|
||||||
dbt_build = BashOperator(
|
dbt_build = BashOperator(
|
||||||
task_id="dbt_build",
|
task_id="dbt_build",
|
||||||
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_gias_establishments+ stg_gias_links+ gias_code_names+ --exclude int_ks2_with_lineage+ int_ks4_with_lineage+",
|
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_gias_establishments+ stg_gias_links+ gias_code_names+ --exclude int_ks2_with_lineage+ int_ks4_with_lineage+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
|
||||||
)
|
)
|
||||||
|
|
||||||
sync_typesense = BashOperator(
|
sync_typesense = BashOperator(
|
||||||
@@ -143,7 +146,7 @@ with DAG(
|
|||||||
|
|
||||||
dbt_build_ofsted = BashOperator(
|
dbt_build_ofsted = BashOperator(
|
||||||
task_id="dbt_build",
|
task_id="dbt_build",
|
||||||
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ofsted_inspections+ int_ofsted_latest+ fact_ofsted_inspection+ dim_school+",
|
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ofsted_inspections+ int_ofsted_latest+ fact_ofsted_inspection+ dim_school+ --exclude stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
|
||||||
)
|
)
|
||||||
|
|
||||||
sync_typesense_ofsted = BashOperator(
|
sync_typesense_ofsted = BashOperator(
|
||||||
|
|||||||
@@ -0,0 +1,26 @@
|
|||||||
|
"""Read a GIAS extract from the raw bytes of the download.
|
||||||
|
|
||||||
|
GIAS writes its CSVs in Windows-1252 and sends no charset, so `resp.text`
|
||||||
|
leaves requests to guess the codec. On 3 Oct 2026 it guessed windows-1250 and
|
||||||
|
"à" became "ŕ". Decode the bytes ourselves instead.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import io
|
||||||
|
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
GIAS_ENCODING = "cp1252"
|
||||||
|
|
||||||
|
|
||||||
|
def read_gias_csv(content: bytes, logger=None) -> pd.DataFrame:
|
||||||
|
"""Every column as a string; a blank cell stays ''."""
|
||||||
|
# Windows-1252 leaves five bytes undefined. One stray byte must not stop
|
||||||
|
# the daily refresh of every school, so it becomes U+FFFD and is logged.
|
||||||
|
text = content.decode(GIAS_ENCODING, errors="replace")
|
||||||
|
undecodable = text.count("�")
|
||||||
|
if undecodable and logger is not None:
|
||||||
|
logger.warning("%d byte(s) in the GIAS extract could not be decoded as %s",
|
||||||
|
undecodable, GIAS_ENCODING)
|
||||||
|
return pd.read_csv(io.StringIO(text), dtype=str, keep_default_na=False)
|
||||||
@@ -7,6 +7,8 @@ from datetime import date, timedelta
|
|||||||
from singer_sdk import Stream, Tap
|
from singer_sdk import Stream, Tap
|
||||||
from singer_sdk import typing as th
|
from singer_sdk import typing as th
|
||||||
|
|
||||||
|
from tap_uk_gias.gias_csv import read_gias_csv
|
||||||
|
|
||||||
GIAS_URL_TEMPLATE = (
|
GIAS_URL_TEMPLATE = (
|
||||||
"https://ea-edubase-api-prod.azurewebsites.net"
|
"https://ea-edubase-api-prod.azurewebsites.net"
|
||||||
"/edubase/downloads/public/edubasealldata{date}.csv"
|
"/edubase/downloads/public/edubasealldata{date}.csv"
|
||||||
@@ -74,9 +76,6 @@ class GIASEstablishmentsStream(Stream):
|
|||||||
|
|
||||||
def get_records(self, context):
|
def get_records(self, context):
|
||||||
"""Download GIAS CSV and yield rows."""
|
"""Download GIAS CSV and yield rows."""
|
||||||
import io
|
|
||||||
|
|
||||||
import pandas as pd
|
|
||||||
import requests
|
import requests
|
||||||
|
|
||||||
today = date.today()
|
today = date.today()
|
||||||
@@ -94,12 +93,7 @@ class GIASEstablishmentsStream(Stream):
|
|||||||
|
|
||||||
resp.raise_for_status()
|
resp.raise_for_status()
|
||||||
|
|
||||||
df = pd.read_csv(
|
df = read_gias_csv(resp.content, self.logger)
|
||||||
io.StringIO(resp.text),
|
|
||||||
encoding="latin-1",
|
|
||||||
dtype=str,
|
|
||||||
keep_default_na=False,
|
|
||||||
)
|
|
||||||
|
|
||||||
for _, row in df.iterrows():
|
for _, row in df.iterrows():
|
||||||
record = row.to_dict()
|
record = row.to_dict()
|
||||||
@@ -126,9 +120,6 @@ class GIASLinksStream(Stream):
|
|||||||
|
|
||||||
def get_records(self, context):
|
def get_records(self, context):
|
||||||
"""Download GIAS links CSV and yield rows."""
|
"""Download GIAS links CSV and yield rows."""
|
||||||
import io
|
|
||||||
|
|
||||||
import pandas as pd
|
|
||||||
import requests
|
import requests
|
||||||
|
|
||||||
today = date.today()
|
today = date.today()
|
||||||
@@ -146,12 +137,7 @@ class GIASLinksStream(Stream):
|
|||||||
|
|
||||||
resp.raise_for_status()
|
resp.raise_for_status()
|
||||||
|
|
||||||
df = pd.read_csv(
|
df = read_gias_csv(resp.content, self.logger)
|
||||||
io.StringIO(resp.text),
|
|
||||||
encoding="latin-1",
|
|
||||||
dtype=str,
|
|
||||||
keep_default_na=False,
|
|
||||||
)
|
|
||||||
|
|
||||||
for _, row in df.iterrows():
|
for _, row in df.iterrows():
|
||||||
record = row.to_dict()
|
record = row.to_dict()
|
||||||
|
|||||||
@@ -0,0 +1,98 @@
|
|||||||
|
"""Every scheduled dbt build must only build models whose parents exist.
|
||||||
|
|
||||||
|
The daily GIAS build selects `stg_gias_establishments+`, so any mart that joins
|
||||||
|
dim_school joins the daily build too. When such a mart also reads a staging
|
||||||
|
model that only the manually triggered EES DAG builds, the daily build fails in
|
||||||
|
any database where that DAG has not run since. Sync and cache invalidation then
|
||||||
|
never run either. The destinations marts did this from late August 2026.
|
||||||
|
|
||||||
|
The graph is read from the model SQL, because CI has no dbt.
|
||||||
|
"""
|
||||||
|
import re
|
||||||
|
from collections import defaultdict
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
PIPELINE = Path(__file__).resolve().parents[1]
|
||||||
|
MODELS = PIPELINE / 'transform' / 'models'
|
||||||
|
DAG_FILE = PIPELINE / 'dags' / 'school_data_pipeline.py'
|
||||||
|
|
||||||
|
REF = re.compile(r"ref\(\s*'([a-z0-9_]+)'\s*\)")
|
||||||
|
DBT_BUILD = re.compile(r'dbt_build\w*\s*=\s*BashOperator\(.*?build --profiles-dir \. --target production ([^"]+)"', re.S)
|
||||||
|
DAG_ID = re.compile(r'dag_id="([a-z0-9_]+)"')
|
||||||
|
|
||||||
|
DAILY = 'school_data_daily'
|
||||||
|
|
||||||
|
# dim_school reads int_ofsted_latest only when the relation exists
|
||||||
|
# (adapter.get_relation), so a missing table is not a failure.
|
||||||
|
OPTIONAL_PARENTS = {'int_ofsted_latest'}
|
||||||
|
|
||||||
|
|
||||||
|
def model_parents():
|
||||||
|
"""{model: models it refs}. Seeds are left out: they are loaded once and always exist."""
|
||||||
|
sql = {p.stem: p.read_text() for p in MODELS.rglob('*.sql')}
|
||||||
|
return {name: set(REF.findall(text)) & set(sql) for name, text in sql.items()}
|
||||||
|
|
||||||
|
|
||||||
|
def downstream(node, children):
|
||||||
|
seen, stack = {node}, [node]
|
||||||
|
while stack:
|
||||||
|
for child in children[stack.pop()]:
|
||||||
|
if child not in seen:
|
||||||
|
seen.add(child)
|
||||||
|
stack.append(child)
|
||||||
|
return seen
|
||||||
|
|
||||||
|
|
||||||
|
def expand(tokens, children):
|
||||||
|
out = set()
|
||||||
|
for token in tokens:
|
||||||
|
out |= downstream(token[:-1], children) if token.endswith('+') else {token}
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def scheduled_builds():
|
||||||
|
"""{dag_id: dbt selection arguments} for every dbt build in the DAG file."""
|
||||||
|
text = DAG_FILE.read_text()
|
||||||
|
starts = [(m.start(), m.group(1)) for m in DAG_ID.finditer(text)]
|
||||||
|
builds = {}
|
||||||
|
for i, (start, dag_id) in enumerate(starts):
|
||||||
|
end = starts[i + 1][0] if i + 1 < len(starts) else len(text)
|
||||||
|
found = DBT_BUILD.search(text, start, end)
|
||||||
|
if found:
|
||||||
|
builds[dag_id] = found.group(1)
|
||||||
|
return builds
|
||||||
|
|
||||||
|
|
||||||
|
def selected_models(args, parents):
|
||||||
|
children = defaultdict(set)
|
||||||
|
for model, ps in parents.items():
|
||||||
|
for p in ps:
|
||||||
|
children[p].add(model)
|
||||||
|
select = re.search(r'--select (.+?)(?= --exclude|$)', args).group(1).split()
|
||||||
|
excluded = re.search(r'--exclude (.+)$', args)
|
||||||
|
exclude = excluded.group(1).split() if excluded else []
|
||||||
|
return (expand(select, children) - expand(exclude, children)) & set(parents)
|
||||||
|
|
||||||
|
|
||||||
|
PARENTS = model_parents()
|
||||||
|
BUILDS = scheduled_builds()
|
||||||
|
DAILY_MODELS = selected_models(BUILDS[DAILY], PARENTS)
|
||||||
|
|
||||||
|
|
||||||
|
def test_every_dag_with_a_dbt_build_is_parsed():
|
||||||
|
assert set(BUILDS) == {
|
||||||
|
'school_data_daily', 'school_data_monthly_ofsted', 'school_data_annual_ees',
|
||||||
|
'school_data_annual_idaci', 'school_data_annual_distance',
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize('dag_id', sorted(BUILDS))
|
||||||
|
def test_selected_models_only_read_models_that_exist(dag_id):
|
||||||
|
selected = selected_models(BUILDS[dag_id], PARENTS)
|
||||||
|
# The daily build is the base layer: other DAGs may rely on what it builds.
|
||||||
|
available = selected | OPTIONAL_PARENTS | (DAILY_MODELS if dag_id != DAILY else set())
|
||||||
|
missing = {model: sorted(PARENTS[model] - available) for model in sorted(selected)
|
||||||
|
if PARENTS[model] - available}
|
||||||
|
assert missing == {}, f'{dag_id} builds models whose parents it never builds: {missing}'
|
||||||
@@ -0,0 +1,66 @@
|
|||||||
|
"""GIAS publishes its extracts in Windows-1252 and declares no charset.
|
||||||
|
|
||||||
|
The tap used to hand pandas `resp.text`, so requests guessed the codec.
|
||||||
|
On 3 Oct 2026 it guessed windows-1250, and "St Thomas à Becket" was stored
|
||||||
|
as "St Thomas ŕ Becket". The `encoding=` passed to read_csv did nothing,
|
||||||
|
because the text was already decoded.
|
||||||
|
"""
|
||||||
|
import importlib.util
|
||||||
|
import logging
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
MODULE = (Path(__file__).resolve().parents[1] / 'plugins' / 'extractors' / 'tap-uk-gias'
|
||||||
|
/ 'tap_uk_gias' / 'gias_csv.py')
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def gias_csv():
|
||||||
|
spec = importlib.util.spec_from_file_location('gias_csv', MODULE)
|
||||||
|
module = importlib.util.module_from_spec(spec)
|
||||||
|
spec.loader.exec_module(module)
|
||||||
|
return module
|
||||||
|
|
||||||
|
|
||||||
|
# Byte for byte as GIAS writes it: 0xE0 à, 0x92 ’, 0xE9 é, 0xB0 °, 0xE7 ç.
|
||||||
|
EXTRACT = (
|
||||||
|
b'"URN","EstablishmentName","HeadLastName"\r\n'
|
||||||
|
b'"138950","St Thomas \xe0 Becket Catholic Secondary School","Smith"\r\n'
|
||||||
|
b'"100000","The Dean and Chapter of St Paul\x92s Cathedral","Pr\xe9vert"\r\n'
|
||||||
|
b'"140677","North Star 180\xb0","Fran\xe7ois"\r\n'
|
||||||
|
b'"100001","No head recorded",""\r\n'
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_names_decode_as_windows_1252(gias_csv):
|
||||||
|
df = gias_csv.read_gias_csv(EXTRACT)
|
||||||
|
assert list(df['EstablishmentName']) == [
|
||||||
|
'St Thomas à Becket Catholic Secondary School',
|
||||||
|
'The Dean and Chapter of St Paul’s Cathedral',
|
||||||
|
'North Star 180°',
|
||||||
|
'No head recorded',
|
||||||
|
]
|
||||||
|
assert list(df['HeadLastName']) == ['Smith', 'Prévert', 'François', '']
|
||||||
|
|
||||||
|
|
||||||
|
def test_the_codec_requests_guessed_is_not_used(gias_csv):
|
||||||
|
# What the tap stored on 3 Oct: the same bytes read as windows-1250.
|
||||||
|
assert 'ŕ' in EXTRACT.decode('cp1250')
|
||||||
|
names = ' '.join(gias_csv.read_gias_csv(EXTRACT)['EstablishmentName'])
|
||||||
|
assert 'ŕ' not in names
|
||||||
|
|
||||||
|
|
||||||
|
def test_values_stay_strings(gias_csv):
|
||||||
|
df = gias_csv.read_gias_csv(EXTRACT)
|
||||||
|
assert df.loc[0, 'URN'] == '138950'
|
||||||
|
|
||||||
|
|
||||||
|
def test_a_byte_windows_1252_leaves_undefined_does_not_stop_the_load(gias_csv, caplog):
|
||||||
|
# 0x81 has no Windows-1252 character. One odd name must not block the daily
|
||||||
|
# refresh of every school, but it must be visible in the log.
|
||||||
|
extract = b'"URN","EstablishmentName"\r\n"100002","Odd \x81 Name"\r\n'
|
||||||
|
with caplog.at_level(logging.WARNING):
|
||||||
|
df = gias_csv.read_gias_csv(extract, logger=logging.getLogger('gias'))
|
||||||
|
assert df.loc[0, 'EstablishmentName'] == 'Odd � Name'
|
||||||
|
assert 'could not be decoded' in caplog.text
|
||||||
Reference in new issue
Block a user