diff --git a/statvar_imports/us_cdc/single_race/README.md b/statvar_imports/us_cdc/single_race/README.md index 7b76c5985a..3c0888ace8 100644 --- a/statvar_imports/us_cdc/single_race/README.md +++ b/statvar_imports/us_cdc/single_race/README.md @@ -1,10 +1,11 @@ ### This import process handles data from wonder.cdc platform. -- Description: Mortality statistics, categorized by demographic factors and specific causes of death, location, race at county level. +- Description: Mortality statistics, categorized by demographic factors and specific + causes of death, location, race at county level. - Source URL: https://wonder.cdc.gov/ucd-icd10-expanded.html -- Import Type: Semi-Automated +- Import Type: Automated - Data Availability: 2018 onwards @@ -12,48 +13,57 @@ ### Preprocessing and Data Acquisition --Download: Manual +- Download: Automated live downloader (`download.py`) -To obtain the raw input files, data must be manually downloaded from the source. The download process involves selecting specific criteria from the dropdown menus: +The script connects directly to the CDC WONDER platform +(`https://wonder.cdc.gov/ucd-icd10-expanded.html`), automates the session agreement, +and downloads county-level mortality datasets across: + * Year (2018 onwards) + * County + * Sex (Male, Female) + * Single Race (6 categories) + * ICD-10-113 Cause List - *Year - *County - *Sex - *Single Race (6 categories) - *ICD-10-113 Cause List +The script automatically partitions queries state by state, dynamically splits +high-population states into 2-year chunks to respect CDC's 75,000 row export limit, +and batches downloads in sessions with automatic renewal and cooldown to avoid rate limits. -For each download, a specific state must be selected. Critical form options: - * **Show Totals**: Disabled (must be unchecked to avoid subtotal pollution) - * **Show Zero Values**: Disabled - * **Show Suppressed Values**: False +To run the live download: +```bash +# Execute via shell wrapper: +bash download.sh -After making the selections, click the "Send" button at the bottom to initiate the download. +# Or directly with python: +python3 download.py -Once all state files are downloaded, stage them to GCS: -```bash -gsutil -m cp *.csv gs://unresolved_mcf/cdc/UnderlyingCause/Single_Race/latest/input_files/ +# Download specific states or years: +python3 download.py --states=02,48 --years=2018-2024 ``` - ### Data Processing -To get the input files, run the following command. The `download.sh` script will create an input_files folder and copy all the necessary files into it from the GCS : - -```bash - - sh download.sh -``` -After the files are downloaded, the data is processed using the stat_var_processor.py script. The script uses various command-line arguments to specify the input data, pvmap, configuration file, and output path. - +After downloading, input files will be placed into the `input_files/` directory. +The data is processed using the `stat_var_processor.py` script: ```bash - - python3 ../../../tools/statvar_importer/stat_var_processor.py --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf --input_data=input_files/*.csv --pv_map=single_race_pvmap.csv --config_file=single_race_metadata.csv --output_path=output/underlyingcauseofdeath_singlerace --output_counters=counters/underlyingcauseofdeath_singlerace.csv +python3 ../../../tools/statvar_importer/stat_var_processor.py \ + --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf \ + --input_data=input_files/*.csv \ + --pv_map=single_race_pvmap.csv \ + --config_file=single_race_metadata.csv \ + --output_path=output/underlyingcauseofdeath_singlerace \ + --output_counters=counters/underlyingcauseofdeath_singlerace.csv ``` ### Automation -This import pipeline is configured to run Semi-automatic on the second Saturday of every month schedule. +This import pipeline is configured in `manifest.json` on the second Saturday of +every month schedule: + +- Data Commons Manifest Schedule: `30 08 8-14 * 6` -- Cron Expression: 30 08 8-14 * 6 +*(Note: In standard POSIX crontabs where day-of-month and day-of-week evaluate as an +`OR` condition, use +`30 08 * * 6 [ $(date +\%d) -ge 8 ] && [ $(date +\%d) -le 14 ] && bash download.sh` +to restrict execution strictly to the second Saturday).* diff --git a/statvar_imports/us_cdc/single_race/download.py b/statvar_imports/us_cdc/single_race/download.py new file mode 100755 index 0000000000..f1a5281ab3 --- /dev/null +++ b/statvar_imports/us_cdc/single_race/download.py @@ -0,0 +1,747 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Automated live downloader for CDC WONDER Single Race Mortality Data (D158). + +This script automates downloading county-level mortality statistics from CDC +WONDER (Database D158: Underlying Cause of Death, Single Race). + +Data is broken down by: +- Year (2018 to 2024, or specified range) +- County +- Sex (Male, Female) +- Single Race 6 (6 categories) +- ICD-10 113 Cause List + +CDC WONDER imposes a hard limit of 75,000 rows per export query. This script +queries state by state, automatically detects when a state query exceeds the +75,000 row cap, and splits into smaller year chunks to download complete data. +""" + +import csv +import datetime +import io +import os +from pathlib import Path +import time +from typing import Dict, List, Optional, Tuple +from urllib.parse import urljoin + +from absl import app +from absl import flags +from absl import logging +from bs4 import BeautifulSoup +import requests +from requests.adapters import HTTPAdapter +from retry import retry +from urllib3.util import Retry + +_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__)) +DEFAULT_INPUT_DIR = os.path.join(_SCRIPT_DIR, "input_files") +SOURCE_LANDING_URL = "https://wonder.cdc.gov/ucd-icd10-expanded.html" + +# 50 States + District of Columbia FIPS codes +US_STATES: Dict[str, str] = { + "01": "Alabama", + "02": "Alaska", + "04": "Arizona", + "05": "Arkansas", + "06": "California", + "08": "Colorado", + "09": "Connecticut", + "10": "Delaware", + "11": "District of Columbia", + "12": "Florida", + "13": "Georgia", + "15": "Hawaii", + "16": "Idaho", + "17": "Illinois", + "18": "Indiana", + "19": "Iowa", + "20": "Kansas", + "21": "Kentucky", + "22": "Louisiana", + "23": "Maine", + "24": "Maryland", + "25": "Massachusetts", + "26": "Michigan", + "27": "Minnesota", + "28": "Mississippi", + "29": "Missouri", + "30": "Montana", + "31": "Nebraska", + "32": "Nevada", + "33": "New Hampshire", + "34": "New Jersey", + "35": "New Mexico", + "36": "New York", + "37": "North Carolina", + "38": "North Dakota", + "39": "Ohio", + "40": "Oklahoma", + "41": "Oregon", + "42": "Pennsylvania", + "44": "Rhode Island", + "45": "South Carolina", + "46": "South Dakota", + "47": "Tennessee", + "48": "Texas", + "49": "Utah", + "50": "Vermont", + "51": "Virginia", + "53": "Washington", + "54": "West Virginia", + "55": "Wisconsin", + "56": "Wyoming", +} + +# Populous states that exceed CDC WONDER's 75,000 row cap across 6 years. +# Querying directly in 2-year chunks prevents query buffer overruns and HTTP 400 errors. +LARGE_STATES: set[str] = { + "01", + "06", + "12", + "13", + "17", + "18", + "21", + "22", + "26", + "27", + "28", + "29", + "34", + "36", + "37", + "39", + "40", + "42", + "45", + "47", + "48", + "51", + "53", + "55", +} + +FLAGS = flags.FLAGS +flags.DEFINE_string( + "states", + "all", + "Comma-separated 2-digit FIPS codes of states to download (e.g. '02,48'), or 'all'.", +) +_CURRENT_YEAR = datetime.date.today().year +_DEFAULT_YEARS = f"2018-{_CURRENT_YEAR}" + +flags.DEFINE_string( + "years", + _DEFAULT_YEARS, + f"Year range (e.g. '2018-{_CURRENT_YEAR}') or comma-separated years.", +) +flags.DEFINE_string( + "output_dir", + DEFAULT_INPUT_DIR, + "Directory where downloaded CSV files will be saved.", +) +flags.DEFINE_float( + "delay", + 5.0, + "Politeness delay in seconds between successive HTTP queries.", +) +flags.DEFINE_integer( + "timeout", + 120, + "HTTP request timeout in seconds.", +) +flags.DEFINE_bool( + "skip_existing", + True, + "Skip downloading states that already have existing non-empty CSV files in output_dir.", +) +flags.DEFINE_integer( + "batch_size", + 8, + "Number of states to process per session before automatically refreshing session.", +) +flags.DEFINE_float( + "batch_cooldown", + 60.0, + "Cooldown delay in seconds between session batches to prevent rate limits.", +) + + +def parse_year_list(year_str: str) -> List[str]: + """Parses a year string like '2018-2024' or '2018,2019' into a list of year strings.""" + year_str = year_str.strip() + if "-" in year_str and not year_str.startswith("-"): + parts = year_str.split("-") + start, end = int(parts[0]), int(parts[1]) + return [str(y) for y in range(start, end + 1)] + return [y.strip() for y in year_str.split(",") if y.strip()] + + +class CdcWonderSingleRaceDownloader: + """Automates CDC WONDER sessions and queries for Single Race mortality data.""" + + def _create_session(self) -> requests.Session: + """Creates a requests.Session with connection pooling and unified SSL.""" + session = requests.Session() + session.verify = True + session.headers.update({ + "User-Agent": + "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" + }) + adapter = HTTPAdapter( + pool_connections=10, + pool_maxsize=10, + max_retries=Retry( + total=3, + backoff_factor=1, + status_forcelist=[500, 502, 503, 504], + raise_on_status=False, + ), + ) + session.mount("https://", adapter) + session.mount("http://", adapter) + return session + + def __init__( + self, + landing_url: str = SOURCE_LANDING_URL, + timeout: int = 120, + delay: float = 5.0, + ): + self.landing_url = landing_url + self.timeout = timeout + self.delay = delay + self.session = self._create_session() + self.action_url: Optional[str] = None + self.base_post_data: List[Tuple[str, str]] = [] + self.available_years: List[str] = [] + + @retry(tries=3, + delay=5, + backoff=2, + exceptions=(requests.RequestException, ValueError)) + def init_session(self): + """Accesses landing page, submits agreement, and extracts base query form parameters.""" + if hasattr(self, "session") and self.session: + try: + self.session.close() + except Exception: + pass + self.session = self._create_session() + self.action_url = None + self.base_post_data = [] + self.available_years = [] + + logging.info("Connecting to CDC WONDER landing page: %s", + self.landing_url) + res = self.session.get(self.landing_url, timeout=self.timeout) + res.raise_for_status() + + soup = BeautifulSoup(res.text, "lxml") + form = soup.find("form", id="wonderform") + if not form: + raise ValueError( + "Could not find initial wonderform on CDC WONDER page.") + + action = urljoin(self.landing_url, form.get("action")) + agree_inputs = [(inp.get("name"), inp.get("value", "")) + for inp in form.find_all("input") + if inp.get("name")] + agree_inputs.append(("action-I Agree", "I Agree")) + + logging.info("Submitting Data Use Agreement (I Agree)...") + res_agree = self.session.post(action, + data=agree_inputs, + timeout=self.timeout) + res_agree.raise_for_status() + + soup_req = BeautifulSoup(res_agree.text, "lxml") + form_req = soup_req.find("form", id="wonderform") + if not form_req: + raise ValueError( + "Could not find request form after agreeing to terms.") + + self.action_url = urljoin(self.landing_url, form_req.get("action")) + + # Extract available years published in the CDC WONDER form + year_select = form_req.find("select", {"name": "F_D158.V1"}) + if year_select: + self.available_years = [ + opt.get("value") + for opt in year_select.find_all("option") + if opt.get("value") and opt.get("value") != "*All*" + ] + logging.info("Detected available years on CDC WONDER: %s", + self.available_years) + else: + logging.warning( + "Year select element 'F_D158.V1' not found in CDC WONDER form. " + "Dynamic year filtering disabled.") + self.available_years = [] + + # Extract pre-populated query parameters + self.base_post_data = [] + for el in form_req.find_all(["input", "select", "textarea"]): + name = el.get("name") + if not name: + continue + if el.name == "input": + itype = el.get("type", "text").lower() + if itype in ["submit", "button", "reset", "image"]: + continue + if itype in ["checkbox", "radio"]: + if el.has_attr("checked"): + self.base_post_data.append((name, el.get("value", + "on"))) + else: + self.base_post_data.append((name, el.get("value", ""))) + elif el.name == "select": + selected_opts = [ + opt for opt in el.find_all("option") + if opt.has_attr("selected") + ] + if selected_opts: + for opt in selected_opts: + self.base_post_data.append((name, opt.get("value", ""))) + else: + if not el.has_attr("multiple"): + first_opt = el.find("option") + if first_opt: + self.base_post_data.append( + (name, first_opt.get("value", ""))) + elif el.name == "textarea": + self.base_post_data.append((name, el.text or "")) + + logging.info( + "Successfully established CDC WONDER session with action: %s", + self.action_url) + + def _build_post_data( + self, + state_fips: str, + years: Optional[List[str]] = None) -> List[Tuple[str, str]]: + """Constructs query payload for single race mortality with groupings and filters.""" + query_data: List[Tuple[str, str]] = [] + for k, v in self.base_post_data: + # Grouping fields: + # B_1: Year + # B_2: County + # B_3: Sex + # B_4: Single Race 6 + # B_5: ICD-10 113 Cause List + if k == "B_1": + query_data.append((k, "D158.V1-level1")) + elif k == "B_2": + query_data.append((k, "D158.V9-level2")) + elif k == "B_3": + query_data.append((k, "D158.V7")) + elif k == "B_4": + query_data.append((k, "D158.V42")) + elif k == "B_5": + query_data.append((k, "D158.V4")) + elif k == "F_D158.V9": + # Filter by state FIPS code + query_data.append((k, state_fips)) + elif k == "F_D158.V1": + # Year filter - handle separately below if specific years requested + if not years: + query_data.append((k, v)) + else: + query_data.append((k, v)) + + if years: + for y in years: + if not self.available_years or y in self.available_years: + query_data.append(("F_D158.V1", y)) + + query_data.append(("action-Export Results", "Export Results")) + return query_data + + def execute_query( + self, + state_fips: str, + years: Optional[List[str]] = None, + max_retries: int = 5, + ) -> str: + """Executes query with automatic 429 rate-limit backoff and session renewal.""" + if not self.action_url or not self.base_post_data: + self.init_session() + + payload = self._build_post_data(state_fips, years) + + for attempt in range(1, max_retries + 1): + try: + logging.info( + "Dispatching POST %s for state %s, years %s (attempt %d/%d)...", + self.action_url, + state_fips, + years or "all", + attempt, + max_retries, + ) + res = self.session.post(self.action_url, + data=payload, + timeout=self.timeout) + if res.status_code == 429: + if attempt == max_retries: + res.raise_for_status() + retry_after = res.headers.get("Retry-After") + # CDC WONDER WAF explicitly states: + # "Your IP address has been temporarily blocked... Please wait 30 minutes" + # " before trying again." + # Any probe before 30 minutes resets the firewall penalty timer. + # Therefore, on 429 we must pause for the full 30 minutes (+ 1 min buffer) + # in complete silence. + wait_time = int( + retry_after) if retry_after and retry_after.isdigit( + ) else 1860 # 31 minutes + logging.warning( + "Encountered HTTP 429 (Too Many Requests). CDC WONDER enforces a " + "30-minute IP block. Waiting %d seconds (%d min) in complete silence " + "for block to clear (attempt %d/%d)...", + wait_time, + wait_time // 60, + attempt, + max_retries, + ) + time.sleep(wait_time) + # Re-initialize session to renew cookies and session ID + self.init_session() + payload = self._build_post_data(state_fips, years) + continue + + if res.status_code == 400: + logging.warning( + "CDC WONDER returned HTTP 400 for state %s (likely query buffer " + "overrun for large state). Returning for partitioning.", + state_fips, + ) + return "CDC WONDER 400 Bad Request (query too large)" + + res.raise_for_status() + logging.info( + "Successfully received HTTP %d for state %s (%d bytes)", + res.status_code, + state_fips, + len(res.content), + ) + return res.text + except (requests.RequestException, ValueError) as e: + if attempt == max_retries: + raise + wait_time = 15 * attempt + logging.warning( + "Request error: %s. Renewing session and retrying in %d seconds " + "(attempt %d/%d)...", + e, + wait_time, + attempt, + max_retries, + ) + time.sleep(wait_time) + self.init_session() + payload = self._build_post_data(state_fips, years) + + raise RuntimeError( + f"Failed to query {state_fips} after {max_retries} attempts.") + + def download_state(self, state_fips: str, + years: List[str]) -> List[Tuple[str, str]]: + """Downloads data for a given state. + + Automatically detects if the state exceeds the 75,000 row limit, times out, + or encounters server errors across all years, and partitions into smaller chunks. + + Returns: + List of tuples: (chunk_label, tsv_content) + """ + state_name = US_STATES.get(state_fips, f"FIPS-{state_fips}") + logging.info("Querying data for %s (FIPS %s) for years %s...", + state_name, state_fips, years) + + if state_fips in LARGE_STATES and len(years) > 2: + logging.info( + "%s is a high-volume state (>75k rows). Querying directly in 2-year chunks...", + state_name) + else: + # Try querying all years first + response_text = self.execute_query(state_fips, years) + first_line = response_text.split("\n", 1)[0] + if "County Code" in first_line: + logging.info( + "Successfully fetched %s (all requested years in 1 query).", + state_name) + return [("all", response_text)] + logging.warning( + "%s response not TSV (likely exceeded 75k rows: %s). " + "Partitioning into chunks...", + state_name, + response_text[:120].strip().replace("\n", " "), + ) + + # Partition years into 2-year chunks + chunk_results = [] + chunk_size = 2 if len(years) > 2 else 1 + for i in range(0, len(years), chunk_size): + year_chunk = years[i:i + chunk_size] + chunk_label = f"{year_chunk[0]}_{year_chunk[-1]}" if len( + year_chunk) > 1 else year_chunk[0] + logging.info( + "Querying %s for chunk %s (%s)...", + state_name, + chunk_label, + year_chunk, + ) + time.sleep(self.delay) + chunk_text = self.execute_query(state_fips, year_chunk) + chunk_first_line = chunk_text.split("\n", 1)[0] + + if "County Code" not in chunk_first_line: + # If 2-year chunk is still too big, try 1-year chunks + if len(year_chunk) > 1: + logging.warning( + "Chunk %s still too large for %s. Splitting into 1-year chunks...", + chunk_label, + state_name, + ) + for single_year in year_chunk: + time.sleep(self.delay) + sy_text = self.execute_query(state_fips, [single_year]) + if "County Code" not in sy_text.split("\n", 1)[0]: + raise ValueError( + f"Failed to query {state_name} even for single year " + f"{single_year}.") + chunk_results.append((single_year, sy_text)) + else: + raise ValueError( + f"Failed to query {state_name} for chunk {chunk_label}: " + f"{chunk_text[:300]}") + else: + chunk_results.append((chunk_label, chunk_text)) + + return chunk_results + + +def save_tsv_as_csv(raw_tsv: str, output_filepath: str) -> int: + """Converts TSV text from CDC WONDER into CSV format, writing atomically.""" + Path(os.path.dirname(output_filepath)).mkdir(parents=True, exist_ok=True) + tsv_reader = csv.reader(io.StringIO(raw_tsv), delimiter="\t") + + temp_filepath = f"{output_filepath}.tmp" + row_count = 0 + with open(temp_filepath, "w", newline="", encoding="utf-8") as f: + csv_writer = csv.writer(f) + for row in tsv_reader: + if not row: + continue + # Stop at metadata notes footer + if row[0].startswith("---") or (len(row) > 1 and + row[1].startswith("---")): + break + csv_writer.writerow(row) + row_count += 1 + + if row_count <= 1: + if os.path.exists(temp_filepath): + os.remove(temp_filepath) + raise ValueError( + f"No data rows found in TSV response for {output_filepath} " + f"(row_count={row_count}).") + + os.replace(temp_filepath, output_filepath) + logging.info("Saved %d rows to %s", row_count, output_filepath) + return row_count + + +def is_state_downloaded(output_dir: str, + state_fips: str, + years: Optional[List[str]] = None) -> bool: + """Checks if valid non-empty CSV files for this state already exist in output_dir.""" + pattern = f"UnderlyingCauseofDeath_SingleRace_{state_fips}*.csv" + matches = list(Path(output_dir).glob(pattern)) + if not matches: + return False + # Verify that all found files are non-empty (> 100 bytes) + if not all(f.stat().st_size > 100 for f in matches): + return False + if years: + covered_years = set() + for f in matches: + name = f.stem + parts = name.split("_") + if len(parts) == 3: + try: + with open(f, mode="r", encoding="utf-8", + errors="replace") as fh: + reader = csv.reader(fh) + header = next(reader, None) + if header: + year_idx = header.index( + "Year") if "Year" in header else 1 + for row in reader: + if len(row + ) > year_idx and row[year_idx] in years: + covered_years.add(row[year_idx]) + if all(y in covered_years for y in years): + break + except Exception: + pass + elif len(parts) == 4: + covered_years.add(parts[3]) + elif len(parts) == 5: + try: + start, end = int(parts[3]), int(parts[4]) + for y in range(start, end + 1): + covered_years.add(str(y)) + except ValueError: + pass + return all(y in covered_years for y in years) + return True + + +def download_single_race_data( + states: List[str], + years: List[str], + output_dir: str, + delay: float = 5.0, + timeout: int = 120, + skip_existing: bool = True, + batch_size: int = 8, + batch_cooldown: float = 60.0, +): + """Downloads CDC Single Race mortality data for specified states and years.""" + os.makedirs(output_dir, exist_ok=True) + downloader = CdcWonderSingleRaceDownloader(timeout=timeout, delay=delay) + downloader.init_session() + + # Intersect requested years with published years on CDC WONDER to prevent HTTP 500 + if downloader.available_years: + valid_years = [y for y in years if y in downloader.available_years] + if valid_years: + if len(valid_years) < len(years): + skipped = [y for y in years if y not in downloader.available_years] + logging.info( + "Skipping years not yet published in CDC WONDER: %s (querying: %s)", + skipped, valid_years) + years = valid_years + else: + logging.warning( + "None of requested years %s found in CDC WONDER available years %s. " + "Defaulting to available years: %s", + years, downloader.available_years, downloader.available_years) + years = downloader.available_years + + total_files = 0 + total_rows = 0 + states_in_batch = 0 + + for idx, state_fips in enumerate(states, start=1): + state_name = US_STATES.get(state_fips, f"FIPS-{state_fips}") + + if skip_existing and is_state_downloaded( + output_dir, state_fips, years=years): + existing_files = list( + Path(output_dir).glob( + f"UnderlyingCauseofDeath_SingleRace_{state_fips}*.csv")) + logging.info( + "[%d/%d] Skipping %s (FIPS %s): %d existing file(s) found.", + idx, + len(states), + state_name, + state_fips, + len(existing_files), + ) + continue + + logging.info( + "[%d/%d] Processing %s (FIPS %s) (Session batch item %d/%d)...", + idx, + len(states), + state_name, + state_fips, + states_in_batch + 1, + batch_size, + ) + + # Purge preexisting files for this state to prevent duplicate data ingestion + # if earlier runs created monolithic vs chunked files (CHK-5.2). + for old_file in Path(output_dir).glob( + f"UnderlyingCauseofDeath_SingleRace_{state_fips}*.csv"): + if old_file.is_file(): + logging.info("Purging preexisting state file: %s", old_file) + old_file.unlink() + + chunks = downloader.download_state(state_fips, years) + + for chunk_label, tsv_data in chunks: + if chunk_label == "all": + filename = f"UnderlyingCauseofDeath_SingleRace_{state_fips}.csv" + else: + filename = f"UnderlyingCauseofDeath_SingleRace_{state_fips}_{chunk_label}.csv" + + output_file = os.path.join(output_dir, filename) + rows = save_tsv_as_csv(tsv_data, output_file) + total_files += 1 + total_rows += rows + + states_in_batch += 1 + + # Automatically refresh session after each batch of states + if states_in_batch >= batch_size and idx < len(states): + logging.info( + "Completed session batch of %d states. Cooling down for %.1fs and renewing " + "CDC session...", + states_in_batch, + batch_cooldown, + ) + time.sleep(batch_cooldown) + downloader.init_session() + states_in_batch = 0 + elif idx < len(states): + time.sleep(delay) + + logging.info("Download complete: Saved %d files with %d total rows in %s.", + total_files, total_rows, output_dir) + + +def main(_): + """Main entry point to execute the CDC WONDER downloader.""" + years = parse_year_list(FLAGS.years) + + if FLAGS.states.lower() == "all": + states = sorted(list(US_STATES.keys())) + else: + states = [ + s.strip().zfill(2) for s in FLAGS.states.split(",") if s.strip() + ] + + logging.info( + "Starting CDC Single Race live download for %d states, years: %s", + len(states), years) + download_single_race_data( + states=states, + years=years, + output_dir=FLAGS.output_dir, + delay=FLAGS.delay, + timeout=FLAGS.timeout, + skip_existing=FLAGS.skip_existing, + batch_size=FLAGS.batch_size, + batch_cooldown=FLAGS.batch_cooldown, + ) + + +if __name__ == "__main__": + app.run(main) diff --git a/statvar_imports/us_cdc/single_race/download.sh b/statvar_imports/us_cdc/single_race/download.sh index 2a712aa2b1..cf855e664c 100755 --- a/statvar_imports/us_cdc/single_race/download.sh +++ b/statvar_imports/us_cdc/single_race/download.sh @@ -1,6 +1,8 @@ #!/bin/bash -set -e -o pipefail +set -euo pipefail -mkdir -p input_files -gcloud storage cp gs://unresolved_mcf/cdc/UnderlyingCause/Single_Race/latest/input_files/*.csv input_files/ +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +mkdir -p "${SCRIPT_DIR}/input_files" + +python3 "${SCRIPT_DIR}/download.py" "$@" diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py new file mode 100644 index 0000000000..28ee0f3294 --- /dev/null +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -0,0 +1,434 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Unit tests for CDC WONDER Single Race Downloader.""" + +import os +from pathlib import Path +import tempfile +import unittest +from unittest import mock + +import sys + +_MODULE_DIR = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, _MODULE_DIR) + +try: + from statvar_imports.us_cdc.single_race import download +except ModuleNotFoundError: + import download + + +class DownloadTest(unittest.TestCase): + """Unit test suite for CDC WONDER Single Race Downloader.""" + + def test_parse_year_list_range(self): + """Tests parsing a range of years string.""" + years = download.parse_year_list("2018-2023") + self.assertEqual(years, + ["2018", "2019", "2020", "2021", "2022", "2023"]) + + def test_parse_year_list_comma(self): + """Tests parsing comma-separated years string.""" + years = download.parse_year_list("2018, 2020, 2022") + self.assertEqual(years, ["2018", "2020", "2022"]) + + def test_parse_year_list_single(self): + """Tests parsing a single year string.""" + years = download.parse_year_list("2024") + self.assertEqual(years, ["2024"]) + + @mock.patch.object(download.logging, "warning") + @mock.patch.object(download.requests, "Session") + def test_init_session_success(self, mock_session_cls, mock_warning): + """Tests successful session handshake and parameter extraction.""" + mock_session = mock.MagicMock() + mock_session_cls.return_value = mock_session + + # 1. Landing page response + mock_res1 = mock.MagicMock() + mock_res1.text = """ + + +
+ +
+ + + """ + mock_res1.raise_for_status.return_value = None + + # 2. Agreement POST response + mock_res2 = mock.MagicMock() + mock_res2.text = """ + + +
+ + + +
+ + + """ + mock_res2.raise_for_status.return_value = None + + mock_session.get.return_value = mock_res1 + mock_session.post.return_value = mock_res2 + + downloader = download.CdcWonderSingleRaceDownloader() + downloader.init_session() + + self.assertIn("controller/datarequest/D158;jsessionid=TEST1234", + downloader.action_url) + self.assertTrue(len(downloader.base_post_data) > 0) + self.assertEqual(downloader.available_years, []) + mock_warning.assert_called_once_with( + "Year select element 'F_D158.V1' not found in CDC WONDER form. " + "Dynamic year filtering disabled.") + self.assertEqual(mock_session.post.call_count, 1) + + @mock.patch.object(download.requests, "Session") + def test_init_session_missing_form(self, mock_session_cls): + """Tests error handling when the initial landing form is missing.""" + mock_session = mock.MagicMock() + mock_session_cls.return_value = mock_session + + mock_res = mock.MagicMock() + mock_res.text = "No form here" + mock_res.raise_for_status.return_value = None + mock_session.get.return_value = mock_res + + downloader = download.CdcWonderSingleRaceDownloader() + with self.assertRaisesRegex(ValueError, + "Could not find initial wonderform"): + downloader.init_session.__wrapped__(downloader) + + def test_build_post_data(self): + """Tests constructing query payload with proper groupings and filters.""" + downloader = download.CdcWonderSingleRaceDownloader() + downloader.base_post_data = [ + ("B_1", "old_val"), + ("B_2", "old_val"), + ("B_3", "old_val"), + ("B_4", "old_val"), + ("B_5", "old_val"), + ("F_D158.V9", "*All*"), + ("F_D158.V1", "*All*"), + ("other_key", "other_val"), + ] + + payload = downloader._build_post_data("02", ["2018", "2019"]) + payload_dict = dict(payload) + + self.assertEqual(payload_dict["B_1"], "D158.V1-level1") + self.assertEqual(payload_dict["B_2"], "D158.V9-level2") + self.assertEqual(payload_dict["B_3"], "D158.V7") + self.assertEqual(payload_dict["B_4"], "D158.V42") + self.assertEqual(payload_dict["B_5"], "D158.V4") + self.assertEqual(payload_dict["F_D158.V9"], "02") + self.assertEqual(payload_dict["action-Export Results"], + "Export Results") + + year_params = [v for k, v in payload if k == "F_D158.V1"] + self.assertEqual(year_params, ["2018", "2019"]) + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "execute_query") + def test_download_state_single_query(self, mock_query): + """Tests downloading state data that fits within a single query.""" + tsv_output = ("Notes\tYear\tCounty\tCounty Code\tDeaths\n" + "\t2018\tAnchorage Borough, AK\t02020\t20\n") + mock_query.return_value = tsv_output + + downloader = download.CdcWonderSingleRaceDownloader() + results = downloader.download_state("02", ["2018", "2019"]) + + self.assertEqual(len(results), 1) + self.assertEqual(results[0][0], "all") + self.assertEqual(results[0][1], tsv_output) + mock_query.assert_called_once_with("02", ["2018", "2019"]) + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "execute_query") + def test_download_state_row_limit_partitioning(self, mock_query): + """Tests dynamic partitioning when CDC WONDER row limit is exceeded.""" + err_msg = ( + "This request produces 168,754 rows, but 75,000 is the maximum allowed. " + "Simplify this request, or send a series of smaller ones.") + tsv_chunk1 = "Notes\tYear\tCounty Code\n\t2018\t06001\n" + tsv_chunk2 = "Notes\tYear\tCounty Code\n\t2020\t06001\n" + tsv_chunk3 = "Notes\tYear\tCounty Code\n\t2022\t06001\n" + + # Test dynamic partitioning when 75k limit is hit for a non-preemptive state + mock_query.side_effect = [err_msg, tsv_chunk1, tsv_chunk2, tsv_chunk3] + + downloader = download.CdcWonderSingleRaceDownloader(delay=0.0) + results = downloader.download_state( + "99", ["2018", "2019", "2020", "2021", "2022", "2023"]) + + self.assertEqual(len(results), 3) + self.assertEqual(results[0][0], "2018_2019") + self.assertEqual(results[1][0], "2020_2021") + self.assertEqual(results[2][0], "2022_2023") + self.assertEqual(mock_query.call_count, 4) + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "execute_query") + def test_download_state_large_state_preemptive(self, mock_query): + """Tests preemptive chunking for high-volume states.""" + tsv_chunk1 = "Notes\tYear\tCounty Code\n\t2018\t06001\n" + tsv_chunk2 = "Notes\tYear\tCounty Code\n\t2020\t06001\n" + tsv_chunk3 = "Notes\tYear\tCounty Code\n\t2022\t06001\n" + mock_query.side_effect = [tsv_chunk1, tsv_chunk2, tsv_chunk3] + + downloader = download.CdcWonderSingleRaceDownloader(delay=0.0) + # California ("06") is in LARGE_STATES, should directly query 3 chunks (no 6-year attempt) + results = downloader.download_state( + "06", ["2018", "2019", "2020", "2021", "2022", "2023"]) + + self.assertEqual(len(results), 3) + self.assertEqual(mock_query.call_count, 3) + + def test_save_tsv_as_csv(self): + """Tests converting CDC WONDER TSV output to CSV format.""" + raw_tsv = ("Notes\tYear\tCounty\tCounty Code\tDeaths\n" + "\t2018\tAnchorage Borough, AK\t02020\t20\n" + "\t2018\tFairbanks North Star Borough, AK\t02090\t15\n" + "---\n" + "Query Parameters:\n" + "Caveats:\n") + + with tempfile.TemporaryDirectory() as temp_dir: + output_csv = os.path.join(temp_dir, "test_output.csv") + rows = download.save_tsv_as_csv(raw_tsv, output_csv) + + self.assertEqual(rows, 3) + self.assertTrue(os.path.exists(output_csv)) + lines = Path(output_csv).read_text(encoding="utf-8").splitlines() + + self.assertEqual(len(lines), 3) + self.assertEqual(lines[0], "Notes,Year,County,County Code,Deaths") + self.assertEqual(lines[1], ',2018,"Anchorage Borough, AK",02020,20') + + def test_save_tsv_as_csv_empty_raises(self): + """Tests that saving an empty or header-only TSV raises ValueError.""" + raw_tsv = "Notes\tYear\tCounty\tCounty Code\tDeaths\n---\n" + with tempfile.TemporaryDirectory() as temp_dir: + output_csv = os.path.join(temp_dir, "empty_output.csv") + with self.assertRaisesRegex(ValueError, "No data rows found"): + download.save_tsv_as_csv(raw_tsv, output_csv) + + def test_is_state_downloaded(self): + """Tests state download detection across monolithic and chunked files.""" + with tempfile.TemporaryDirectory() as temp_dir: + self.assertFalse(download.is_state_downloaded(temp_dir, "02")) + + # Create empty file + f = Path(temp_dir) / "UnderlyingCauseofDeath_SingleRace_02.csv" + f.write_text("") + self.assertFalse(download.is_state_downloaded(temp_dir, "02")) + + # Create valid monolithic file covering 2018 and 2019 + f.write_text("Header,col1,col2,col3\n" + + ",2018,val\n,2019,val\n" * 10) + self.assertTrue(download.is_state_downloaded(temp_dir, "02")) + self.assertTrue( + download.is_state_downloaded(temp_dir, "02", ["2018", "2019"])) + self.assertFalse( + download.is_state_downloaded(temp_dir, "02", ["2018", "2020"])) + + # Verify numeric value in another column does not trigger false positive + f_other_col = (Path(temp_dir) / + "UnderlyingCauseofDeath_SingleRace_03.csv") + f_other_col.write_text("Notes,Year,Deaths\n" + ",2020,2018\n" * 10) + self.assertFalse( + download.is_state_downloaded(temp_dir, "03", ["2018"])) + self.assertTrue( + download.is_state_downloaded(temp_dir, "03", ["2020"])) + + # Test chunk files + f.unlink() + f_chunk = (Path(temp_dir) / + "UnderlyingCauseofDeath_SingleRace_02_2018_2019.csv") + f_chunk.write_text("Header,col1,col2,col3\n" + + "val1,val2,val3,val4\n" * 10) + self.assertTrue( + download.is_state_downloaded(temp_dir, "02", ["2018", "2019"])) + self.assertFalse( + download.is_state_downloaded(temp_dir, "02", ["2018", "2020"])) + + @mock.patch.object(download.time, "sleep") + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "init_session") + def test_execute_query_429_backoff(self, mock_init, mock_sleep): + """Tests rate-limit handling and backoff on HTTP 429.""" + downloader = download.CdcWonderSingleRaceDownloader() + downloader.action_url = "https://wonder.cdc.gov/test" + downloader.base_post_data = [("B_1", "test")] + + mock_res_429 = mock.MagicMock() + mock_res_429.status_code = 429 + mock_res_429.headers = {"Retry-After": "1"} + + mock_res_200 = mock.MagicMock() + mock_res_200.status_code = 200 + mock_res_200.text = "Notes\tCounty Code\n" + mock_res_200.raise_for_status.return_value = None + + downloader.session.post = mock.MagicMock( + side_effect=[mock_res_429, mock_res_200]) + + result = downloader.execute_query("02", ["2018"], max_retries=2) + self.assertEqual(result, "Notes\tCounty Code\n") + self.assertEqual(downloader.session.post.call_count, 2) + mock_sleep.assert_called_with(1) + mock_init.assert_called_once() + + @mock.patch.object(download.time, "sleep") + def test_execute_query_429_final_attempt_fails_fast(self, mock_sleep): + """Tests that HTTP 429 on the final attempt fails immediately without sleeping.""" + downloader = download.CdcWonderSingleRaceDownloader() + downloader.action_url = "https://wonder.cdc.gov/test" + downloader.base_post_data = [("B_1", "test")] + + mock_res_429 = mock.MagicMock() + mock_res_429.status_code = 429 + mock_res_429.raise_for_status.side_effect = ( + download.requests.HTTPError("429 Client Error")) + + downloader.session.post = mock.MagicMock(return_value=mock_res_429) + + with self.assertRaises(download.requests.HTTPError): + downloader.execute_query("02", ["2018"], max_retries=1) + + mock_sleep.assert_not_called() + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "execute_query") + def test_download_state_network_exception_bubbles(self, mock_query): + """Tests that network or runtime exceptions bubble up without entering chunking.""" + mock_query.side_effect = download.requests.ConnectionError("Connection aborted") + + downloader = download.CdcWonderSingleRaceDownloader() + with self.assertRaises(download.requests.ConnectionError): + downloader.download_state("02", ["2018", "2019"]) + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "init_session") + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "download_state") + def test_download_single_race_data_success(self, mock_download_state, + mock_init): + """Tests orchestration and file writing in download_single_race_data.""" + tsv_content = ("Notes\tYear\tCounty\tCounty Code\tDeaths\n" + "\t2018\tAnchorage Borough, AK\t02020\t20\n" + "---\n") + mock_download_state.return_value = [("all", tsv_content)] + + with tempfile.TemporaryDirectory() as temp_dir: + download.download_single_race_data( + states=["02"], + years=["2018"], + output_dir=temp_dir, + delay=0.0, + skip_existing=False, + ) + expected_file = ( + Path(temp_dir) / "UnderlyingCauseofDeath_SingleRace_02.csv") + self.assertTrue(expected_file.exists()) + lines = expected_file.read_text(encoding="utf-8").splitlines() + self.assertEqual(len(lines), 2) + mock_init.assert_called_once() + mock_download_state.assert_called_once_with("02", ["2018"]) + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "init_session") + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "download_state") + def test_download_single_race_data_purges_preexisting_files( + self, mock_download_state, mock_init + ): + """Tests that stale monolithic or chunked files are purged before new downloads.""" + with tempfile.TemporaryDirectory() as temp_dir: + # Create old monolithic file and unrelated state file + old_monolith = ( + Path(temp_dir) / "UnderlyingCauseofDeath_SingleRace_02.csv") + old_monolith.write_text("old data") + other_state = ( + Path(temp_dir) / "UnderlyingCauseofDeath_SingleRace_04.csv") + other_state.write_text("other state data") + + # Mock new download returning a 2-year chunk + tsv_chunk = ("Notes\tYear\tCounty\tCounty Code\tDeaths\n" + "\t2018\tAnchorage Borough, AK\t02020\t20\n" + "---\n") + mock_download_state.return_value = [("2018_2019", tsv_chunk)] + + download.download_single_race_data( + states=["02"], + years=["2018", "2019"], + output_dir=temp_dir, + delay=0.0, + skip_existing=False, + ) + + # Monolith for 02 should be purged, new chunk should exist, and other state intact + self.assertFalse(old_monolith.exists()) + new_chunk = ( + Path(temp_dir) / + "UnderlyingCauseofDeath_SingleRace_02_2018_2019.csv") + self.assertTrue(new_chunk.exists()) + self.assertTrue(other_state.exists()) + mock_init.assert_called_once() + + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "init_session") + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "download_state") + def test_download_single_race_data_skip_existing(self, mock_download_state, + mock_init): + """Tests skipping states that are already downloaded.""" + with tempfile.TemporaryDirectory() as temp_dir: + # Create valid file covering 2018 + f = Path(temp_dir) / "UnderlyingCauseofDeath_SingleRace_02.csv" + f.write_text("Notes,Year,County,County Code,Deaths\n" + + ",2018,Anchorage,02020,20\n" * 5) + + download.download_single_race_data( + states=["02"], + years=["2018"], + output_dir=temp_dir, + delay=0.0, + skip_existing=True, + ) + + mock_download_state.assert_not_called() + mock_init.assert_called_once() + + @mock.patch.object(download, "CdcWonderSingleRaceDownloader") + def test_download_single_race_data_filters_unavailable_years( + self, mock_downloader_cls + ): + """Tests that unpublished future years are filtered out from requested years.""" + mock_instance = mock.MagicMock() + mock_instance.available_years = ["2018", "2019", "2020", "2024"] + mock_instance.download_state.return_value = [] + mock_downloader_cls.return_value = mock_instance + + with tempfile.TemporaryDirectory() as temp_dir: + download.download_single_race_data( + states=["02"], + years=["2018", "2024", "2025", "2026"], + output_dir=temp_dir, + delay=0.0, + skip_existing=False, + ) + + # download_state should only be called with published years (2018 and 2024) + mock_instance.download_state.assert_called_once_with( + "02", ["2018", "2024"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/statvar_imports/us_cdc/single_race/validation_config.json b/statvar_imports/us_cdc/single_race/validation_config.json index 9ccffb978e..99f939f8f9 100644 --- a/statvar_imports/us_cdc/single_race/validation_config.json +++ b/statvar_imports/us_cdc/single_race/validation_config.json @@ -3,11 +3,19 @@ "rules": [ { "rule_id": "check_latest_date_freshness", - "description": "Asserts the latest observation date reaches at least 2024 for CDC WONDER single-race mortality data.", + "description": "Asserts data reaches 2024 with >= 690 StatVars reaching 2024.", "validator": "SQL_VALIDATOR", "params": { - "query": "SELECT CAST(MAX(MaxDate) AS INTEGER) AS latest_year FROM stats", - "condition": "latest_year >= 2024" + "query": "SELECT CAST(MAX(MaxDate) AS INTEGER) AS latest_year, COUNT(DISTINCT StatVar) AS total_svs, SUM(CASE WHEN MaxDate >= '2024' THEN 1 ELSE 0 END) AS fresh_svs FROM stats", + "condition": "latest_year >= 2024 AND total_svs >= 805 AND fresh_svs >= 690" + } + }, + { + "rule_id": "check_deleted_records_percent", + "description": "Checks that the percentage of deleted records is within 0.1% threshold.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 0.1 } } ]