Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a82882cbdc |
No files matched your search
@@ -18,10 +18,7 @@
|
|||||||
# TYPESENSE_SEARCH_KEY — Typesense search-only key (exposed to frontend)
|
# TYPESENSE_SEARCH_KEY — Typesense search-only key (exposed to frontend)
|
||||||
# UNLEASH_URL — http://<unleash-ip>:4242/api (empty = all flags off)
|
# UNLEASH_URL — http://<unleash-ip>:4242/api (empty = all flags off)
|
||||||
# UNLEASH_API_TOKEN — Unleash *client* token, environment: development
|
# UNLEASH_API_TOKEN — Unleash *client* token, environment: development
|
||||||
# AIRFLOW_ADMIN_USER — Airflow admin username (default: admin)
|
# AIRFLOW_ADMIN_USER — Airflow admin username (password auto-generated, see api-server logs)
|
||||||
# AIRFLOW_ADMIN_PASSWORD — Airflow admin password. REQUIRED: the api-server
|
|
||||||
# refuses to start without it, rather than falling
|
|
||||||
# back to a generated one that changes on restart.
|
|
||||||
# STAGING_DB_IP — macvlan IP for staging Postgres (default 10.0.1.190)
|
# STAGING_DB_IP — macvlan IP for staging Postgres (default 10.0.1.190)
|
||||||
# STAGING_FRONTEND_IP — macvlan IP for staging frontend (default 10.0.1.151)
|
# STAGING_FRONTEND_IP — macvlan IP for staging frontend (default 10.0.1.151)
|
||||||
|
|
||||||
@@ -127,23 +124,7 @@ services:
|
|||||||
airflow-api-server:
|
airflow-api-server:
|
||||||
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:staging
|
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:staging
|
||||||
container_name: sc_staging_airflow_api
|
container_name: sc_staging_airflow_api
|
||||||
# The simple auth manager generates a random password on first start and
|
command: airflow api-server --port 8080
|
||||||
# writes it to a file, so every container restart invalidates the last one.
|
|
||||||
# Writing the file ourselves from an environment variable makes the login
|
|
||||||
# deterministic. Airflow does not generate anything when the file exists.
|
|
||||||
#
|
|
||||||
# Built with python rather than echo/printf so a password containing quotes,
|
|
||||||
# backslashes or spaces is escaped correctly by json.dumps. An unset
|
|
||||||
# AIRFLOW_ADMIN_PASSWORD raises KeyError and the container exits: falling
|
|
||||||
# back to a generated password would silently undo the point of this.
|
|
||||||
command:
|
|
||||||
- bash
|
|
||||||
- -c
|
|
||||||
- |
|
|
||||||
set -euo pipefail
|
|
||||||
mkdir -p /opt/airflow
|
|
||||||
python -c "import json, os, pathlib; pathlib.Path('/opt/airflow/simple_auth_manager_passwords.json').write_text(json.dumps({os.environ.get('AIRFLOW_ADMIN_USER', 'admin'): os.environ['AIRFLOW_ADMIN_PASSWORD']}))"
|
|
||||||
exec airflow api-server --port 8080
|
|
||||||
ports:
|
ports:
|
||||||
- "8081:8080"
|
- "8081:8080"
|
||||||
environment:
|
environment:
|
||||||
@@ -155,8 +136,6 @@ services:
|
|||||||
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-staging-airflow-jwt-secret-key-long-enough-for-sha512"
|
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-staging-airflow-jwt-secret-key-long-enough-for-sha512"
|
||||||
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "${AIRFLOW_ADMIN_USER:-admin}:admin"
|
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "${AIRFLOW_ADMIN_USER:-admin}:admin"
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_PASSWORDS_FILE: /opt/airflow/simple_auth_manager_passwords.json
|
|
||||||
AIRFLOW_ADMIN_PASSWORD: ${AIRFLOW_ADMIN_PASSWORD:?set AIRFLOW_ADMIN_PASSWORD in the Portainer stack environment}
|
|
||||||
AIRFLOW__LOGGING__BASE_LOG_FOLDER: /opt/airflow/logs
|
AIRFLOW__LOGGING__BASE_LOG_FOLDER: /opt/airflow/logs
|
||||||
PG_HOST: sc_database
|
PG_HOST: sc_database
|
||||||
PG_PORT: "5432"
|
PG_PORT: "5432"
|
||||||
|
|||||||
@@ -9,10 +9,7 @@
|
|||||||
# TYPESENSE_SEARCH_KEY — Typesense search-only key (exposed to frontend)
|
# TYPESENSE_SEARCH_KEY — Typesense search-only key (exposed to frontend)
|
||||||
# UNLEASH_URL — http://<unleash-ip>:4242/api (empty = all flags off)
|
# UNLEASH_URL — http://<unleash-ip>:4242/api (empty = all flags off)
|
||||||
# UNLEASH_API_TOKEN — Unleash *client* token, environment: production
|
# UNLEASH_API_TOKEN — Unleash *client* token, environment: production
|
||||||
# AIRFLOW_ADMIN_USER — Airflow admin username (default: admin)
|
# AIRFLOW_ADMIN_USER — Airflow admin username (password auto-generated, see api-server logs)
|
||||||
# AIRFLOW_ADMIN_PASSWORD — Airflow admin password. REQUIRED: the api-server
|
|
||||||
# refuses to start without it, rather than falling
|
|
||||||
# back to a generated one that changes on restart.
|
|
||||||
|
|
||||||
services:
|
services:
|
||||||
|
|
||||||
@@ -116,23 +113,7 @@ services:
|
|||||||
airflow-api-server:
|
airflow-api-server:
|
||||||
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:prod
|
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:prod
|
||||||
container_name: schoolcompare_airflow_api
|
container_name: schoolcompare_airflow_api
|
||||||
# The simple auth manager generates a random password on first start and
|
command: airflow api-server --port 8080
|
||||||
# writes it to a file, so every container restart invalidates the last one.
|
|
||||||
# Writing the file ourselves from an environment variable makes the login
|
|
||||||
# deterministic. Airflow does not generate anything when the file exists.
|
|
||||||
#
|
|
||||||
# Built with python rather than echo/printf so a password containing quotes,
|
|
||||||
# backslashes or spaces is escaped correctly by json.dumps. An unset
|
|
||||||
# AIRFLOW_ADMIN_PASSWORD raises KeyError and the container exits: falling
|
|
||||||
# back to a generated password would silently undo the point of this.
|
|
||||||
command:
|
|
||||||
- bash
|
|
||||||
- -c
|
|
||||||
- |
|
|
||||||
set -euo pipefail
|
|
||||||
mkdir -p /opt/airflow
|
|
||||||
python -c "import json, os, pathlib; pathlib.Path('/opt/airflow/simple_auth_manager_passwords.json').write_text(json.dumps({os.environ.get('AIRFLOW_ADMIN_USER', 'admin'): os.environ['AIRFLOW_ADMIN_PASSWORD']}))"
|
|
||||||
exec airflow api-server --port 8080
|
|
||||||
ports:
|
ports:
|
||||||
- "8080:8080"
|
- "8080:8080"
|
||||||
environment:
|
environment:
|
||||||
@@ -144,8 +125,6 @@ services:
|
|||||||
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-airflow-jwt-secret-key-long-enough-for-sha512"
|
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-airflow-jwt-secret-key-long-enough-for-sha512"
|
||||||
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "${AIRFLOW_ADMIN_USER:-admin}:admin"
|
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "${AIRFLOW_ADMIN_USER:-admin}:admin"
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_PASSWORDS_FILE: /opt/airflow/simple_auth_manager_passwords.json
|
|
||||||
AIRFLOW_ADMIN_PASSWORD: ${AIRFLOW_ADMIN_PASSWORD:?set AIRFLOW_ADMIN_PASSWORD in the Portainer stack environment}
|
|
||||||
AIRFLOW__LOGGING__BASE_LOG_FOLDER: /opt/airflow/logs
|
AIRFLOW__LOGGING__BASE_LOG_FOLDER: /opt/airflow/logs
|
||||||
PG_HOST: sc_database
|
PG_HOST: sc_database
|
||||||
PG_PORT: "5432"
|
PG_PORT: "5432"
|
||||||
|
|||||||
+1
-19
@@ -105,23 +105,7 @@ services:
|
|||||||
airflow-api-server:
|
airflow-api-server:
|
||||||
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:latest
|
image: privaterepo.sitaru.org/tudor/school_compare-pipeline:latest
|
||||||
container_name: schoolcompare_airflow_api
|
container_name: schoolcompare_airflow_api
|
||||||
# The simple auth manager generates a random password on first start and
|
command: airflow api-server --port 8080
|
||||||
# writes it to a file, so every container restart invalidates the last one.
|
|
||||||
# Writing the file ourselves from an environment variable makes the login
|
|
||||||
# deterministic. Airflow does not generate anything when the file exists.
|
|
||||||
#
|
|
||||||
# Built with python rather than echo/printf so a password containing quotes,
|
|
||||||
# backslashes or spaces is escaped correctly by json.dumps. An unset
|
|
||||||
# AIRFLOW_ADMIN_PASSWORD raises KeyError and the container exits: falling
|
|
||||||
# back to a generated password would silently undo the point of this.
|
|
||||||
command:
|
|
||||||
- bash
|
|
||||||
- -c
|
|
||||||
- |
|
|
||||||
set -euo pipefail
|
|
||||||
mkdir -p /opt/airflow
|
|
||||||
python -c "import json, os, pathlib; pathlib.Path('/opt/airflow/simple_auth_manager_passwords.json').write_text(json.dumps({os.environ.get('AIRFLOW_ADMIN_USER', 'admin'): os.environ['AIRFLOW_ADMIN_PASSWORD']}))"
|
|
||||||
exec airflow api-server --port 8080
|
|
||||||
ports:
|
ports:
|
||||||
- "8080:8080"
|
- "8080:8080"
|
||||||
environment: &airflow-env
|
environment: &airflow-env
|
||||||
@@ -133,8 +117,6 @@ services:
|
|||||||
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-airflow-jwt-secret-key-long-enough-for-sha512"
|
AIRFLOW__API_AUTH__JWT_SECRET: "school-compare-airflow-jwt-secret-key-long-enough-for-sha512"
|
||||||
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
AIRFLOW__API_AUTH__JWT_ISSUER: airflow
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "admin:admin"
|
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS: "admin:admin"
|
||||||
AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_PASSWORDS_FILE: /opt/airflow/simple_auth_manager_passwords.json
|
|
||||||
AIRFLOW_ADMIN_PASSWORD: ${AIRFLOW_ADMIN_PASSWORD:-admin}
|
|
||||||
PG_HOST: db
|
PG_HOST: db
|
||||||
PG_PORT: "5432"
|
PG_PORT: "5432"
|
||||||
PG_USER: schoolcompare
|
PG_USER: schoolcompare
|
||||||
|
|||||||
@@ -98,12 +98,6 @@ fail the E2E gate. That's the point: staging absorbs the risk.
|
|||||||
pr-checks status checks (frontend, backend, builds, ai-review) to pass.
|
pr-checks status checks (frontend, backend, builds, ai-review) to pass.
|
||||||
5. **Bootstrap staging data via Airflow** (no prod dump — staging populates
|
5. **Bootstrap staging data via Airflow** (no prod dump — staging populates
|
||||||
itself from source, exercising the pipeline image end-to-end):
|
itself from source, exercising the pipeline image end-to-end):
|
||||||
- Set `AIRFLOW_ADMIN_PASSWORD` in the stack environment first. The
|
|
||||||
api-server refuses to start without it. Airflow's simple auth manager
|
|
||||||
otherwise generates a password on first start and writes it to a file, so
|
|
||||||
the login changes every time the container restarts; the stack writes that
|
|
||||||
file itself from this variable instead. `AIRFLOW_ADMIN_USER` defaults to
|
|
||||||
`admin`.
|
|
||||||
- Open the staging Airflow UI (`http://<host>:8081`) and trigger, in order:
|
- Open the staging Airflow UI (`http://<host>:8081`) and trigger, in order:
|
||||||
`school_data_daily`, `school_data_monthly_ofsted`, then the manual-schedule
|
`school_data_daily`, `school_data_monthly_ofsted`, then the manual-schedule
|
||||||
`school_data_annual_ees` and `school_data_annual_idaci`.
|
`school_data_annual_ees` and `school_data_annual_idaci`.
|
||||||
|
|||||||
@@ -190,7 +190,7 @@ with DAG(
|
|||||||
|
|
||||||
dbt_build_ees = BashOperator(
|
dbt_build_ees = 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_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+",
|
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+",
|
||||||
)
|
)
|
||||||
|
|
||||||
sync_typesense_ees = BashOperator(
|
sync_typesense_ees = BashOperator(
|
||||||
|
|||||||
+18
-74
@@ -188,25 +188,21 @@ class DestinationsStream(Stream):
|
|||||||
_national_pinned: list[str] = []
|
_national_pinned: list[str] = []
|
||||||
_indicators: dict[str, str] = {}
|
_indicators: dict[str, str] = {}
|
||||||
|
|
||||||
# School rows and the England reference are DIFFERENT GRAINS, so they are
|
schema = th.PropertiesList(
|
||||||
# different streams. Carrying both in one table meant a null `urn` inside
|
th.Property("urn", th.StringType),
|
||||||
# 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("time_period", th.StringType),
|
||||||
th.Property("pupil_group", th.StringType),
|
th.Property("pupil_group", th.StringType),
|
||||||
th.Property("destination_measure", th.StringType),
|
th.Property("destination_measure", th.StringType),
|
||||||
th.Property("cohort_pupils", th.StringType),
|
th.Property("cohort_pupils", th.StringType),
|
||||||
th.Property("pupils_raw", th.StringType),
|
th.Property("pupils_raw", th.StringType),
|
||||||
th.Property("percentage_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
|
replication_key = None
|
||||||
|
|
||||||
def _query(self, period: str, pinned: list[str], keep_level: str,
|
def _query(self, period: str, pinned: list[str], keep_level: str,
|
||||||
urn_by_location: dict[str, str]): # noqa: D401
|
urn_by_location: dict[str, str]):
|
||||||
"""Page one period of one geographic level, yielding Singer records."""
|
"""Page one period of one geographic level, yielding Singer records."""
|
||||||
criteria = [
|
criteria = [
|
||||||
{"filters": {"in": list(self._destination_slugs)}},
|
{"filters": {"in": list(self._destination_slugs)}},
|
||||||
@@ -252,51 +248,19 @@ class DestinationsStream(Stream):
|
|||||||
"%s: %d school locations, %d time periods",
|
"%s: %d school locations, %d time periods",
|
||||||
self.name, len(urn_by_location), len(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:
|
for period in periods:
|
||||||
yield from self._query(
|
yield from self._query(period, self._pinned, "SCH", urn_by_location)
|
||||||
period, self._pinned_for_level, self._level, urn_by_location,
|
yield from self._query(period, self._national_pinned, "NAT", urn_by_location)
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
class SchoolDestinationsStream(DestinationsStream):
|
class KS4DestinationsStream(DestinationsStream):
|
||||||
"""School-level rows. `urn` is part of the key and is never null."""
|
name = "ees_ks4_destinations"
|
||||||
|
|
||||||
_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
|
_dataset_id = KS4_DATASET
|
||||||
_destination_slugs = KS4_DESTINATION_SLUGS
|
_destination_slugs = KS4_DESTINATION_SLUGS
|
||||||
_pupil_group_slugs = KS4_PUPIL_GROUP_SLUGS
|
_pupil_group_slugs = KS4_PUPIL_GROUP_SLUGS
|
||||||
@@ -305,7 +269,8 @@ class _KS4Config(DestinationsStream):
|
|||||||
_indicators = KS4_INDICATORS
|
_indicators = KS4_INDICATORS
|
||||||
|
|
||||||
|
|
||||||
class _KS5Config(DestinationsStream):
|
class KS5DestinationsStream(DestinationsStream):
|
||||||
|
name = "ees_ks5_destinations"
|
||||||
_dataset_id = KS5_DATASET
|
_dataset_id = KS5_DATASET
|
||||||
_destination_slugs = KS5_DESTINATION_SLUGS
|
_destination_slugs = KS5_DESTINATION_SLUGS
|
||||||
_pupil_group_slugs = KS5_PUPIL_GROUP_SLUGS
|
_pupil_group_slugs = KS5_PUPIL_GROUP_SLUGS
|
||||||
@@ -314,33 +279,12 @@ class _KS5Config(DestinationsStream):
|
|||||||
_indicators = KS5_INDICATORS
|
_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):
|
class TapUKEESDestinations(Tap):
|
||||||
name = "tap-uk-ees-destinations"
|
name = "tap-uk-ees-destinations"
|
||||||
config_jsonschema = th.PropertiesList().to_dict()
|
config_jsonschema = th.PropertiesList().to_dict()
|
||||||
|
|
||||||
def discover_streams(self):
|
def discover_streams(self):
|
||||||
return [
|
return [KS4DestinationsStream(self), KS5DestinationsStream(self)]
|
||||||
KS4DestinationsStream(self),
|
|
||||||
KS5DestinationsStream(self),
|
|
||||||
KS4NationalDestinationsStream(self),
|
|
||||||
KS5NationalDestinationsStream(self),
|
|
||||||
]
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
@@ -132,68 +132,3 @@ def test_establishment_dimensions_are_not_pinned():
|
|||||||
establishment_totals = {"4369U", "rgHcN", "EfHQq", "S4ROV"}
|
establishment_totals = {"4369U", "rgHcN", "EfHQq", "S4ROV"}
|
||||||
assert not set(KS4_PINNED) & establishment_totals
|
assert not set(KS4_PINNED) & establishment_totals
|
||||||
assert not set(KS5_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
|
|
||||||
@@ -22,7 +22,8 @@ select
|
|||||||
pupils,
|
pupils,
|
||||||
percentage,
|
percentage,
|
||||||
status
|
status
|
||||||
from {{ ref('stg_ees_ks4_destinations_national') }}
|
from {{ ref('stg_ees_ks4_destinations') }}
|
||||||
|
where urn is null
|
||||||
|
|
||||||
union all
|
union all
|
||||||
|
|
||||||
@@ -35,4 +36,5 @@ select
|
|||||||
pupils,
|
pupils,
|
||||||
percentage,
|
percentage,
|
||||||
status
|
status
|
||||||
from {{ ref('stg_ees_ks5_destinations_national') }}
|
from {{ ref('stg_ees_ks5_destinations') }}
|
||||||
|
where urn is null
|
||||||
@@ -47,17 +47,6 @@ sources:
|
|||||||
16-18 study leavers destinations. Same grain, same suppression
|
16-18 study leavers destinations. Same grain, same suppression
|
||||||
caveat, and only institutions with post-16 provision appear.
|
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
|
- name: ees_ks4_performance
|
||||||
description: KS4 performance tables (long format — one row per school × breakdown × sex)
|
description: KS4 performance tables (long format — one row per school × breakdown × sex)
|
||||||
|
|
||||||
|
|||||||
@@ -1,9 +1,7 @@
|
|||||||
{{ config(materialized='table') }}
|
{{ config(materialized='table') }}
|
||||||
|
|
||||||
-- Staging model: KS4 leavers destinations, school level.
|
-- Staging model: KS4 leavers destinations, school level plus the England
|
||||||
--
|
-- reference (which carries a null urn).
|
||||||
-- 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
|
-- 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"
|
-- (z, c, x, q, u) to NULL, which is right for attainment — there, "suppressed"
|
||||||
@@ -16,12 +14,13 @@
|
|||||||
|
|
||||||
with source as (
|
with source as (
|
||||||
select * from {{ source('raw', 'ees_ks4_destinations') }}
|
select * from {{ source('raw', 'ees_ks4_destinations') }}
|
||||||
where urn ~ '^[0-9]+$'
|
-- National rows carry a null urn and feed fact_destination_national.
|
||||||
|
where (urn is null or urn = '' or urn ~ '^[0-9]+$')
|
||||||
and time_period ~ '^[0-9]+$'
|
and time_period ~ '^[0-9]+$'
|
||||||
)
|
)
|
||||||
|
|
||||||
select
|
select
|
||||||
cast(trim(urn) as integer) as urn,
|
case when urn ~ '^[0-9]+$' then cast(trim(urn) as integer) end as urn,
|
||||||
cast(trim(time_period) as integer) as year,
|
cast(trim(time_period) as integer) as year,
|
||||||
trim(pupil_group) as pupil_group,
|
trim(pupil_group) as pupil_group,
|
||||||
trim(destination_measure) as destination_measure,
|
trim(destination_measure) as destination_measure,
|
||||||
|
|||||||
@@ -1,38 +0,0 @@
|
|||||||
{{ 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
|
|
||||||
@@ -1,10 +1,8 @@
|
|||||||
{{ config(materialized='table') }}
|
{{ config(materialized='table') }}
|
||||||
|
|
||||||
-- Staging model: 16-18 study leavers destinations, institution level.
|
-- Staging model: 16-18 study leavers destinations, institution level plus the
|
||||||
--
|
-- England reference (which carries a null urn). Only sixth forms and colleges
|
||||||
-- Only sixth forms and colleges appear, so a secondary with no post-16
|
-- appear here, so a secondary with no post-16 provision has no rows at all.
|
||||||
-- 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
|
-- 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"
|
-- (z, c, x, q, u) to NULL, which is right for attainment — there, "suppressed"
|
||||||
@@ -17,12 +15,13 @@
|
|||||||
|
|
||||||
with source as (
|
with source as (
|
||||||
select * from {{ source('raw', 'ees_ks5_destinations') }}
|
select * from {{ source('raw', 'ees_ks5_destinations') }}
|
||||||
where urn ~ '^[0-9]+$'
|
-- National rows carry a null urn and feed fact_destination_national.
|
||||||
|
where (urn is null or urn = '' or urn ~ '^[0-9]+$')
|
||||||
and time_period ~ '^[0-9]+$'
|
and time_period ~ '^[0-9]+$'
|
||||||
)
|
)
|
||||||
|
|
||||||
select
|
select
|
||||||
cast(trim(urn) as integer) as urn,
|
case when urn ~ '^[0-9]+$' then cast(trim(urn) as integer) end as urn,
|
||||||
cast(trim(time_period) as integer) as year,
|
cast(trim(time_period) as integer) as year,
|
||||||
trim(pupil_group) as pupil_group,
|
trim(pupil_group) as pupil_group,
|
||||||
trim(destination_measure) as destination_measure,
|
trim(destination_measure) as destination_measure,
|
||||||
|
|||||||
@@ -1,38 +0,0 @@
|
|||||||
{{ 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
|
|
||||||
Reference in new issue
Block a user