Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 81 additions & 8 deletions experimental/dm/marimo_stuff/edm_recipes.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import marimo

__generated_with = "0.17.0"
__generated_with = "0.23.3"
app = marimo.App(width="medium")

with app.setup(hide_code=True):
Expand All @@ -11,25 +11,27 @@

@app.cell(hide_code=True)
def _():
mo.md(r"""# `edm-recipes` Data Catalog""")
mo.md(r"""
# `edm-recipes` Data Catalog
""")
return


@app.cell(hide_code=True)
def _():
mo.md(
r"""
mo.md(r"""
This notebook is for exploring Data Engineering's source data stored in the `edm-recipes` S3 bucket in Digital Ocean.

Every day, a duckdb database in `edm-reicpes` is refreshed to have up-to-date views of all versions of all datasets. This is limited to the versions that have a parquet files in them.
"""
)
""")
return


@app.cell(hide_code=True)
def _():
mo.md(r"""## Setup""")
mo.md(r"""
## Setup
""")
return


Expand Down Expand Up @@ -83,7 +85,78 @@ def _(conn):

@app.cell(hide_code=True)
def _():
mo.md(r"""## Explore an example dataset""")
mo.md(r"""
## Explore a dataset
""")
return


@app.cell
def _(conn):
dataset_names = (
conn.sql("select distinct schema from (show all tables) order by schema")
.df()["schema"]
.to_list()
)
return (dataset_names,)


@app.cell
def _(dataset_names):
dropdown_dataset_names = mo.ui.dropdown(options=dataset_names)
dropdown_dataset_names
return (dropdown_dataset_names,)


@app.cell(hide_code=True)
def _(conn, dropdown_dataset_names):
_df = mo.sql(
f"""
select
*
from
(show all tables)
where
schema = '{dropdown_dataset_names.value}'
""",
engine=conn,
)
return


@app.cell
def _(conn, dropdown_dataset_names):
dataset_versions = (
conn.sql(
f"select name from (show all tables) where schema = '{dropdown_dataset_names.value}' order by name"
)
.df()["name"]
.to_list()
)
dropdown_dataset_versions = mo.ui.dropdown(options=dataset_versions)
dropdown_dataset_versions
return (dropdown_dataset_versions,)


@app.cell(hide_code=True)
def _(conn, dropdown_dataset_names, dropdown_dataset_versions):
_df = mo.sql(
f"""
select
*
from
"{dropdown_dataset_names.value}"."{dropdown_dataset_versions.value}"
""",
engine=conn,
)
return


@app.cell(hide_code=True)
def _():
mo.md(r"""
## Explore an example dataset
""")
return


Expand Down
31 changes: 20 additions & 11 deletions products/facilities/facdb/cli.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import concurrent.futures
import importlib
import shlex
import shutil
import subprocess
import zipfile
Expand Down Expand Up @@ -113,16 +114,14 @@ def _export_fgdb(pg: postgres.PostgresClient, output_dir: Path) -> None:


def _dbt(args: list[str]) -> None:
subprocess.check_call(
[
"dbt",
*args,
"--quiet",
"--warn-error-options",
'{"error": ["NoNodesForSelectionCriteria"]}',
],
cwd=PRODUCT_PATH,
)
cmd = [
"dbt",
*args,
"--warn-error-options",
'{"error": ["NoNodesForSelectionCriteria"]}',
]
typer.echo(typer.style(f"$ {shlex.join(cmd)}", fg=typer.colors.CYAN))
subprocess.check_call(cmd, cwd=PRODUCT_PATH)


@app.command("init")
Expand All @@ -136,7 +135,15 @@ def _cli_init():
)
postgres.execute_file_via_shell(BUILD_ENGINE, SQL_PATH / "_procedures.sql")
_dbt(["deps"])
_dbt(["seed"])
_dbt(
[
"build",
"--select",
"config.materialized:seed",
"--indirect-selection=cautious",
"--full-refresh",
]
)


@app.command("build")
Expand All @@ -159,6 +166,8 @@ def _cli_build():
build_schema=BUILD_NAME,
)
postgres.execute_file_via_shell(BUILD_ENGINE, SQL_PATH / "_deduplication.sql")
postgres.execute_file_via_shell(BUILD_ENGINE, SQL_PATH / "_create_rectype.sql")
postgres.execute_file_via_shell(BUILD_ENGINE, SQL_PATH / "_create_sgr_mock.sql")


@app.command("qaqc")
Expand Down
13 changes: 13 additions & 0 deletions products/facilities/facdb/sql/_create_rectype.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
-- Assign RECTYPE (facility / program) from the factype_rectype_mapping seed.
-- The seed is the source of truth for valid values; a blank rectype or an
-- unmapped factype yields NULL.

ALTER TABLE facdb ADD COLUMN IF NOT EXISTS rectype text;

-- Reset first so a standalone re-run is deterministic (unmapped factypes -> NULL).
UPDATE facdb SET rectype = NULL;

UPDATE facdb
SET rectype = nullif(lower(m.rectype), '')
FROM factype_rectype_mapping AS m
WHERE facdb.factype = m.factype;
96 changes: 96 additions & 0 deletions products/facilities/facdb/sql/_create_sgr_mock.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
-- Add mock SGR (State of Good Repair) columns to the facdb table.
-- The asset id and SGR values are derived deterministically from md5(bin::text)
-- so that (a) every record for the same building (BIN) gets identical values and
-- (b) values are stable across nightly builds. No random() calls; no setseed.
--
-- TEMPORARY: remove this file when the real AIMS SGR source lands (scores/grades
-- arrive pre-computed).

-- Helper function: letter grade from 0–100 score.
-- TEMPORARY/mock — not canonical grade logic. Bands live upstream (omb-aims)
-- where the real SGR source is produced; this is only for generating mock data.
DROP FUNCTION IF EXISTS sgr_grade(smallint);
CREATE FUNCTION sgr_grade(score smallint) RETURNS text AS $$
BEGIN
IF score >= 90 THEN RETURN 'A';
ELSIF score >= 75 THEN RETURN 'B';
ELSIF score >= 60 THEN RETURN 'C';
ELSIF score >= 50 THEN RETURN 'D';
ELSE RETURN 'F';
END IF;
END;
$$ LANGUAGE plpgsql IMMUTABLE STRICT;

-- Add columns (IF NOT EXISTS for idempotency on re-runs).
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS asset_id integer;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_score_arch smallint;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_grade_arch text;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_score_syst smallint;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_grade_syst text;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_score_tot smallint;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_grade_tot text;
ALTER TABLE facdb ADD COLUMN IF NOT EXISTS sgr_assmnt_year smallint;

-- ASSET_ID: mock AIMS asset number. An AIMS asset is a building, so the id is
-- keyed off BIN (not uid): every record sharing a BIN gets the same asset_id and
-- the same SGR values below. Eligible buildings are CITY facility records with a
-- real BIN: AIMS only surveys city facilities, so State/Federal facilities are
-- excluded here even though their rectype is 'facility'. The placeholder
-- "million BINs" are already NULL in facdb (see _create_facdb_spatial.sql).
-- Selection: ~50% of eligible buildings are "in AIMS", via a deterministic hash
-- gate on BIN. Each selected building gets a UNIQUE id from a hash-ordered ranking
-- (1..N), capped at 16000 to match the real AIMS asset-number range. asset_id is
-- 1:1 with BIN, mirroring the AIMS source — no two BINs share an id.
-- Only city facility records receive an asset_id/SGR (the final UPDATE filters
-- rectype = 'facility' and overlevel = 'City'). A program operating inside a
-- scored building does NOT inherit its values — SGR shows only on the facility.
WITH selected_bins AS (
SELECT DISTINCT bin
FROM facdb
WHERE
rectype = 'facility'
AND overlevel = 'City'
AND bin IS NOT NULL
AND abs(('x' || substr(md5(bin::text), 1, 8))::bit(32)::bigint) % 100 < 50
),

numbered_bins AS (
SELECT
bin,
row_number() OVER (ORDER BY md5(bin::text)) AS asset_id
FROM selected_bins
)

UPDATE facdb
SET asset_id = numbered_bins.asset_id
FROM numbered_bins
WHERE
facdb.bin = numbered_bins.bin
AND facdb.rectype = 'facility'
AND facdb.overlevel = 'City'
AND numbered_bins.asset_id <= 16000;

-- Scores and assessment year are also keyed off BIN, so all records for the same
-- building share identical SGR. Non-overlapping md5(bin) substrings keep the three
-- values independent. ('x'||8hex)::bit(32)::bigint can be negative; abs() then mod.
UPDATE facdb SET
sgr_score_arch = (abs(('x' || substr(md5(bin::text), 9, 8))::bit(32)::bigint) % 101)::smallint,
sgr_score_syst = (abs(('x' || substr(md5(bin::text), 17, 8))::bit(32)::bigint) % 101)::smallint,
-- Year: chars 25–32; map to [2010, current_year] inclusive.
sgr_assmnt_year = (
2010 + abs(('x' || substr(md5(bin::text), 25, 8))::bit(32)::bigint)
% (extract(YEAR FROM current_date)::int - 2010 + 1)
)::smallint
WHERE asset_id IS NOT NULL;

-- Total score: weighted average of arch and syst
UPDATE facdb SET
sgr_score_tot = round(0.6 * sgr_score_arch + 0.4 * sgr_score_syst)::smallint
WHERE asset_id IS NOT NULL;

-- Letter grades derived from scores
UPDATE facdb SET
sgr_grade_arch = sgr_grade(sgr_score_arch),
sgr_grade_syst = sgr_grade(sgr_score_syst),
sgr_grade_tot = sgr_grade(sgr_score_tot)
WHERE asset_id IS NOT NULL;
20 changes: 20 additions & 0 deletions products/facilities/facdb/sql/_reformat_facdb.sql
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
-- mapping of data types between postgres and geodabase reference: https://pro.arcgis.com/en/pro-app/3.1/help/data/geodatabases/manage-postgresql/data-types-postgresql.htm

-- create facdb table with expected column names and data types
DROP TABLE IF EXISTS facdb_export;
CREATE TABLE facdb_export AS
SELECT
facname::VARCHAR(250) AS "FACNAME",
Expand Down Expand Up @@ -40,12 +41,22 @@ SELECT
schooldist::VARCHAR(3) AS "SCHOOLDIST",
policeprct::SMALLINT AS "POLICEPRCT",
datasource::VARCHAR(150) AS "DATASOURCE",
rectype::VARCHAR(12) AS "RECTYPE",
asset_id::INTEGER AS "ASSET_ID",
sgr_score_tot::SMALLINT AS "SGR_SCORE_TOT",
sgr_grade_tot::VARCHAR(1) AS "SGR_GRADE_TOT",
sgr_score_arch::SMALLINT AS "SGR_SCORE_ARCH",
sgr_grade_arch::VARCHAR(1) AS "SGR_GRADE_ARCH",
sgr_score_syst::SMALLINT AS "SGR_SCORE_SYST",
sgr_grade_syst::VARCHAR(1) AS "SGR_GRADE_SYST",
sgr_assmnt_year::SMALLINT AS "SGR_ASSMNT_YEAR",
uid::VARCHAR AS "UID",
geom
FROM facdb;


-- create facdb table without geometry column
DROP TABLE IF EXISTS facdb_export_csv;
CREATE TABLE facdb_export_csv AS
SELECT
facname::VARCHAR(250) AS "FACNAME",
Expand Down Expand Up @@ -84,6 +95,15 @@ SELECT
schooldist::VARCHAR(3) AS "SCHOOLDIST",
policeprct::SMALLINT AS "POLICEPRCT",
datasource::VARCHAR(150) AS "DATASOURCE",
rectype::VARCHAR(12) AS "RECTYPE",
asset_id::INTEGER AS "ASSET_ID",
sgr_score_tot::SMALLINT AS "SGR_SCORE_TOT",
sgr_grade_tot::VARCHAR(1) AS "SGR_GRADE_TOT",
sgr_score_arch::SMALLINT AS "SGR_SCORE_ARCH",
sgr_grade_arch::VARCHAR(1) AS "SGR_GRADE_ARCH",
sgr_score_syst::SMALLINT AS "SGR_SCORE_SYST",
sgr_grade_syst::VARCHAR(1) AS "SGR_GRADE_SYST",
sgr_assmnt_year::SMALLINT AS "SGR_ASSMNT_YEAR",
uid::VARCHAR AS "UID"
FROM facdb;

Expand Down
Loading
Loading