Skip to content

Commit 5026289

Browse files
committed
hospitalization data source
1 parent cd391c1 commit 5026289

13 files changed

Lines changed: 1195 additions & 750 deletions

File tree

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
# This dataset is static, and thus only needs to be extracted in it's raw form
2+
# once and saved to blob storage. For a human viewable we format, see:
3+
# https://data.cdc.gov/resource/aemt-mg7g.csv
4+
# Raw (extracted) and transformed (loaded) data will be stored in azure blob storage
5+
# which requires the account, container and path for access
6+
7+
[properties]
8+
9+
name = "hospitalization"
10+
automate = false
11+
transform_template = "hospitalization.sql"
12+
schema = "hospitalization.py"
13+
14+
[source]
15+
16+
url = "https://data.cdc.gov/resource/aemt-mg7g.csv"
17+
pagination = {limit = 1000}
18+
19+
[extract]
20+
21+
account = "cfadatalakeprd"
22+
container = "cfapredict"
23+
prefix = "dataops/scenarios/raw/hospitalization"
24+
25+
[load]
26+
27+
account = "cfadatalakeprd"
28+
container = "cfapredict"
29+
prefix = "dataops/scenarios/transformed/hospitalization"
30+
31+
# TODO: add some data schema validation
Lines changed: 120 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,120 @@
1+
import random
2+
3+
import pandas as pd
4+
import pandera.pandas as pa
5+
from faker import Faker
6+
7+
fake = Faker()
8+
df_len = 1000
9+
10+
extract_schema = pa.DataFrameSchema(
11+
{
12+
"week_end_date": pa.Column(str),
13+
"jurisdiction": pa.Column(str),
14+
"weekly_actual_days_reporting_any_data": pa.Column(float, nullable=True, coerce = True),
15+
"weekly_percent_days_reporting_any_data": pa.Column(float, nullable=True, coerce = True),
16+
"num_hospitals_previous_day_admission_adult_covid_confirmed": pa.Column('int', nullable=True),
17+
"num_hospitals_previous_day_admission_pediatric_covid_confirmed": pa.Column('int', nullable=True),
18+
"num_hospitals_previous_day_admission_influenza_confirmed": pa.Column('int', nullable=True),
19+
"num_hospitals_total_patients_hospitalized_confirmed_influenza": pa.Column('int', nullable=True),
20+
"num_hospitals_icu_patients_confirmed_influenza": pa.Column('int', nullable=True),
21+
"num_hospitals_inpatient_beds": pa.Column('int', nullable=True),
22+
"num_hospitals_total_icu_beds": pa.Column('int', nullable=True),
23+
"num_hospitals_inpatient_beds_used": pa.Column('int', nullable=True),
24+
"num_hospitals_icu_beds_used": pa.Column('int', nullable=True),
25+
"num_hospitals_percent_inpatient_beds_occupied": pa.Column('int', nullable=True),
26+
"num_hospitals_percent_staff_icu_beds_occupied": pa.Column('int', nullable=True),
27+
"num_hospitals_percent_inpatient_beds_covid": pa.Column('int', nullable=True),
28+
"num_hospitals_percent_inpatient_beds_influenza": pa.Column('int', nullable=True),
29+
"num_hospitals_percent_staff_icu_beds_covid": pa.Column('int', nullable=True),
30+
"num_hospitals_percent_icu_beds_influenza": pa.Column('int', nullable=True),
31+
"num_hospitals_admissions_all_covid_confirmed": pa.Column('int', nullable=True),
32+
"num_hospitals_total_patients_hospitalized_covid_confirmed": pa.Column('int', nullable=True),
33+
"num_hospitals_staff_icu_patients_covid_confirmed": pa.Column('int', nullable=True),
34+
"avg_admissions_adult_covid_confirmed": pa.Column('float', nullable=True),
35+
"total_admissions_adult_covid_confirmed": pa.Column('float', nullable=True),
36+
"avg_admissions_pediatric_covid_confirmed": pa.Column('float', nullable=True),
37+
"total_admissions_pediatric_covid_confirmed": pa.Column('float', nullable=True),
38+
"avg_admissions_all_covid_confirmed": pa.Column('float', nullable=True),
39+
"total_admissions_all_covid_confirmed": pa.Column('float', nullable=True),
40+
"avg_admissions_all_influenza_confirmed": pa.Column('float', nullable=True),
41+
"total_admissions_all_influenza_confirmed": pa.Column('float', nullable=True),
42+
"avg_total_patients_hospitalized_covid_confirmed": pa.Column('float', nullable=True),
43+
"avg_total_patients_hospitalized_influenza_confirmed": pa.Column('float', nullable=True),
44+
"avg_staff_icu_patients_covid_confirmed": pa.Column('float', nullable=True),
45+
"avg_icu_patients_influenza_confirmed": pa.Column('float', nullable=True),
46+
"avg_inpatient_beds": pa.Column('float', nullable=True),
47+
"avg_total_icu_beds": pa.Column('float', nullable=True),
48+
"avg_inpatient_beds_used": pa.Column('float', nullable=True),
49+
"avg_icu_beds_used": pa.Column('float', nullable=True),
50+
"avg_percent_inpatient_beds_occupied": pa.Column('float', nullable=True),
51+
"avg_percent_staff_icu_beds_occupied": pa.Column('float', nullable=True),
52+
"avg_percent_inpatient_beds_covid": pa.Column('float', nullable=True),
53+
"avg_percent_inpatient_beds_influenza": pa.Column('float', nullable=True),
54+
"avg_percent_staff_icu_beds_covid": pa.Column('float', nullable=True),
55+
"avg_percent_icu_beds_influenza": pa.Column('float', nullable=True),
56+
"percent_adult_covid_admissions": pa.Column('float', nullable=True),
57+
"percent_pediatric_covid_admissions": pa.Column('float', nullable=True),
58+
"percent_hospitals_previous_day_admission_adult_covid_confirmed": pa.Column('float', nullable=True),
59+
"percent_hospitals_previous_day_admission_pediatric_covid_confirmed": pa.Column('float', nullable=True),
60+
"percent_hospitals_previous_day_admission_influenza_confirmed": pa.Column('float', nullable=True),
61+
"percent_hospitals_total_patients_hospitalized_confirmed_influenza": pa.Column('float', nullable=True),
62+
"percent_hospitals_icu_patients_confirmed_influenza": pa.Column('float', nullable=True),
63+
"percent_hospitals_inpatient_beds": pa.Column('float', nullable=True),
64+
"percent_hospitals_total_icu_beds": pa.Column('float', nullable=True),
65+
"percent_hospitals_inpatient_beds_used": pa.Column('float', nullable=True),
66+
"percent_hospitals_icu_beds_used": pa.Column('float', nullable=True),
67+
"percent_hospitals_percent_inpatient_beds_occupied": pa.Column('float', nullable=True),
68+
"percent_hospitals_percent_staff_icu_beds_occupied": pa.Column('float', nullable=True),
69+
"percent_hospitals_percent_inpatient_beds_covid": pa.Column('float', nullable=True),
70+
"percent_hospitals_percent_inpatient_beds_influenza": pa.Column('float', nullable=True),
71+
"percent_hospitals_percent_staff_icu_beds_covid": pa.Column('float', nullable=True),
72+
"percent_hospitals_percent_icu_beds_influenza": pa.Column('float', nullable=True),
73+
"percent_hospitals_admissions_all_covid_confirmed": pa.Column('float', nullable=True),
74+
"percent_hospitals_total_patients_hospitalized_covid_confirmed": pa.Column('float', nullable=True),
75+
"percent_hospitals_staff_icu_patients_covid_confirmed": pa.Column('float', nullable=True),
76+
"abs_chg_percent_hospitals_previous_day_admission_adult_covid_confirmed": pa.Column('float', nullable=True),
77+
"abs_chg_percent_hospitals_previous_day_admission_pediatric_covid_confirmed": pa.Column('float', nullable=True),
78+
"abs_chg_percent_hospitals_previous_day_admission_influenza_confirmed": pa.Column('float', nullable=True),
79+
"abs_chg_percent_hospitals_total_patients_hospitalized_confirmed_influenza": pa.Column('float', nullable=True),
80+
"abs_chg_percent_hospitals_icu_patients_confirmed_influenza": pa.Column('float', nullable=True),
81+
"abs_chg_percent_hospitals_inpatient_beds": pa.Column('float', nullable=True),
82+
"abs_chg_percent_hospitals_total_icu_beds": pa.Column('float', nullable=True),
83+
"abs_chg_percent_hospitals_inpatient_beds_used": pa.Column('float', nullable=True),
84+
"abs_chg_percent_hospitals_icu_beds_used": pa.Column('float', nullable=True),
85+
"abs_chg_percent_hospitals_percent_inpatient_beds_occupied": pa.Column('float', nullable=True),
86+
"abs_chg_percent_hospitals_percent_staff_icu_beds_occupied": pa.Column('float', nullable=True),
87+
"abs_chg_percent_hospitals_percent_inpatient_beds_covid": pa.Column('float', nullable=True),
88+
"abs_chg_percent_hospitals_percent_inpatient_beds_influenza": pa.Column('float', nullable=True),
89+
"abs_chg_percent_hospitals_percent_staff_icu_beds_covid": pa.Column('float', nullable=True),
90+
"abs_chg_percent_hospitals_percent_icu_beds_influenza": pa.Column('float', nullable=True),
91+
"abs_chg_percent_hospitals_admissions_all_covid_confirmed": pa.Column('float', nullable=True),
92+
"abs_chg_percent_hospitals_total_patients_hospitalized_covid_confirmed": pa.Column('float', nullable=True),
93+
"abs_chg_percent_hospitals_staff_icu_patients_covid_confirmed": pa.Column('float', nullable=True)
94+
}
95+
)
96+
97+
load_schema = pa.DataFrameSchema({
98+
"date": pa.Column(str, coerce=True),
99+
"state": pa.Column(str, coerce=True),
100+
"total": pa.Column(str, coerce = True, nullable = True),
101+
"stname" : pa.Column(str, coerce=True),
102+
})
103+
104+
105+
raw_synth_data = pd.DataFrame({})
106+
107+
stname_tf = {
108+
"CA": "california",
109+
"TX": "texas",
110+
"NY": "new_york",
111+
"FL": "florida",
112+
"IL": "illinois"
113+
}
114+
tf_synth_data = pd.DataFrame({
115+
"date": [fake.date_this_year() for _ in range(df_len)],
116+
"state": [random.choice(["CA", "TX", "NY", "FL", "IL"]) for _ in range(df_len)],
117+
"total": [random.randint(0, 1000) for _ in range(df_len)],
118+
"stname": [fake.state() for _ in range(df_len)],
119+
})
120+
tf_synth_data["stname"] = tf_synth_data["state"].map(stname_tf)
Lines changed: 153 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,153 @@
1+
"""For ETL of the COVID 19 vaccination trends"""
2+
3+
from io import StringIO
4+
5+
import duckdb
6+
import httpx
7+
import pandas as pd
8+
import pandera.pandas as pa
9+
from tqdm import tqdm
10+
from cfa.scenarios.dataops.datasets.catalog import get_data
11+
12+
from ..datasets import datasets
13+
14+
# import schema
15+
from ..datasets.schemas.hospitalization import extract_schema, load_schema
16+
from .utils import get_timestamp, transform_template_lookup
17+
18+
config = datasets.hospitalization
19+
20+
21+
def extract() -> pd.DataFrame:
22+
"""Get data and return raw values as a DataFrame
23+
24+
Returns:
25+
pd.DataFrame: the extracted data
26+
"""
27+
28+
r_count = httpx.get(config.source.url, params={"$select": "count(*)"})
29+
dataset_len = int("".join(char for char in r_count.text if char.isdigit()))
30+
offsets = range(0, dataset_len, config.source.pagination.limit)
31+
32+
extract_blob_version = f"{get_timestamp()}"
33+
34+
get_dfs = []
35+
36+
for idx, offset_i in enumerate(tqdm(offsets)):
37+
params = {
38+
"$limit": config.source.pagination.limit,
39+
"$offset": offset_i,
40+
}
41+
42+
r = httpx.get(config.source.url, params=params)
43+
44+
config.extract.write_blob(
45+
file_buffer=bytes(r.text, "utf-8"),
46+
path_after_prefix=f"{extract_blob_version}/part_{idx}.csv",
47+
)
48+
49+
get_dfs.append(pd.read_csv(StringIO(r.text)))
50+
51+
return pd.concat(get_dfs)
52+
53+
54+
def transform(df: pd.DataFrame) -> pd.DataFrame: # noqa: W0613
55+
"""Run a SQL transform on the raw DataFrame
56+
57+
Args:
58+
df (pd.DataFrame): the extracted data
59+
60+
Returns:
61+
pd.DataFrame: the transformed data
62+
"""
63+
region_id = get_data("fips_to_name_improved")
64+
template = transform_template_lookup.get_template(
65+
config.properties.transform_template
66+
)
67+
query = template.render(data_source="df", region_id = "region_id")
68+
transformed_db = duckdb.sql(query).df()
69+
return transformed_db
70+
71+
72+
def load(df: pd.DataFrame) -> None:
73+
"""Loads transformed data to blob storage
74+
75+
Args:
76+
df (pd.DataFrame): the transformed data to be loaded as parquet
77+
"""
78+
79+
load_blob_version = f"{get_timestamp()}"
80+
81+
config.load.write_blob(
82+
file_buffer=df.to_parquet(),
83+
path_after_prefix=f"{load_blob_version}/data.parquet",
84+
)
85+
86+
87+
def main(
88+
run_extract: bool = False, val_raw: bool = False, val_tf: bool = False
89+
) -> None:
90+
"""ETL main runner
91+
92+
Args:
93+
extract (bool, optional): should the extraction run. Useful in the case
94+
that the raw data is static and doesn't require update while iteration
95+
on the transformation of the raw data. Defaults to False.
96+
val_raw (bool, optional): whether to validate the raw data schema during extraction.
97+
val_tf (bool, optional): whether to validate the transformed data schema before loading to Blob Storage.
98+
"""
99+
100+
if run_extract:
101+
raw_df = extract()
102+
if val_raw:
103+
# check raw_df against schema
104+
try:
105+
extract_schema.validate(raw_df)
106+
except pa.errors.SchemaError as exc:
107+
raise (exc)
108+
109+
else:
110+
try:
111+
buffers = config.extract.read_blobs()
112+
raw_df = pd.concat([pd.read_csv(i) for i in buffers])
113+
except IndexError as e:
114+
raise AttributeError(
115+
"Run extract set to False, but no latest version of extract "
116+
"data to use. Run with extraction step to fetch latest raw "
117+
"version."
118+
) from e
119+
# transform dataframe
120+
transformed_df = transform(raw_df)
121+
# validate transformed df
122+
if val_tf:
123+
try:
124+
load_schema.validate(transformed_df)
125+
except pa.errors.SchemaError as exc:
126+
print(
127+
"Validation of tranformed dataframe failed. Data not loaded to Blob Storage."
128+
)
129+
print("Fix the pipeline and try again.")
130+
raise exc
131+
# load data to blob storage
132+
load(transformed_df)
133+
134+
135+
if __name__ == "__main__":
136+
import argparse
137+
138+
parser = argparse.ArgumentParser(
139+
prog="hospitalization trends data etl pipeline",
140+
)
141+
parser.add_argument("--extract", "-e", action="store_true", default=False)
142+
parser.add_argument(
143+
"--validate_raw", "-v", action="store_true", default=False
144+
)
145+
parser.add_argument(
146+
"--validate_transf", "-t", action="store_true", default=False
147+
)
148+
args = parser.parse_args()
149+
main(
150+
run_extract=args.extract,
151+
val_raw=args.validate_raw,
152+
val_tf=args.validate_transf,
153+
)
Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,14 @@
1-
WITH data AS (
2-
SELECT "state_name" AS stname,
3-
"state_fipcode" as st,
4-
"state_code" as stusps
5-
FROM ${data_source}
6-
)
7-
SELECT 'United States' AS stname,
8-
'US' AS st,
9-
'US' AS stusps
10-
UNION
11-
SELECT DISTINCT stname, st, stusps
12-
FROM data
13-
WHERE stname <> 'Puerto Rico'
14-
ORDER BY stname
1+
WITH data AS (
2+
SELECT "state_name" AS stname,
3+
"state_fipcode" as st,
4+
"state_code" as stusps
5+
FROM ${data_source}
6+
)
7+
SELECT 'United States' AS stname,
8+
'US' AS st,
9+
'US' AS stusps
10+
UNION
11+
SELECT DISTINCT stname, st, stusps
12+
FROM data
13+
WHERE stname <> 'Puerto Rico'
14+
ORDER BY stname
Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,14 @@
1-
WITH region_data AS (
2-
SELECT 'USA' AS region, 'US' AS states
3-
UNION SELECT region, states FROM ${data_source}
4-
5-
),
6-
fips_data AS (
7-
SELECT * FROM ${fips}
8-
)
9-
SELECT fips_data.stname, fips_data.st, fips_data.stusps, region_data.region FROM fips_data
10-
LEFT JOIN
11-
region_data
12-
ON
13-
fips_data.stusps = region_data.states
14-
ORDER BY fips_data.stname
1+
WITH region_data AS (
2+
SELECT 'USA' AS region, 'US' AS states
3+
UNION SELECT region, states FROM ${data_source}
4+
5+
),
6+
fips_data AS (
7+
SELECT * FROM ${fips}
8+
)
9+
SELECT fips_data.stname, fips_data.st, fips_data.stusps, region_data.region FROM fips_data
10+
LEFT JOIN
11+
region_data
12+
ON
13+
fips_data.stusps = region_data.states
14+
ORDER BY fips_data.stname
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
WITH tmp as (
2+
SELECT LEFT(week_end_date, 10) as date,
3+
REPLACE(jurisdiction, 'USA', 'US') as state,
4+
total_admissions_all_covid_confirmed as total
5+
FROM ${data_source}
6+
)
7+
SELECT tmp.date as "date", tmp.state, CAST(tmp.total as int) as total, LOWER(REPLACE(ri.stname, ' ', '_')) as stname
8+
FROM tmp
9+
RIGHT JOIN ${region_id} as ri
10+
ON tmp.state = ri.stusps
Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,3 @@
1-
from . import covid
2-
3-
all = [covid]
1+
from . import covid
2+
3+
all = [covid]

0 commit comments

Comments
 (0)