From 82f6c4fa0ac21197def7e21c78eebaf3b07e83d5 Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Fri, 25 Sep 2026 13:09:31 +0000 Subject: [PATCH 01/10] Automate CDC WONDER single-race mortality import data acquisition - Add download.py to automate live data downloading from CDC WONDER (D158) with session agreement handling, state-level querying, and dynamic 2-year chunking to respect the 75k row export limit. - Add download_test.py with unit tests covering query construction, session initialization, 75k row cap partitioning, large state chunking, and 429 rate-limiting backoff. - Update download.sh to execute download.py. - Update README.md to reflect automated import pipeline and document download parameters and execution commands. --- statvar_imports/us_cdc/single_race/README.md | 48 +- .../us_cdc/single_race/download.py | 655 ++++++++++++++++++ .../us_cdc/single_race/download.sh | 6 +- .../us_cdc/single_race/download_test.py | 248 +++++++ 4 files changed, 927 insertions(+), 30 deletions(-) create mode 100755 statvar_imports/us_cdc/single_race/download.py create mode 100644 statvar_imports/us_cdc/single_race/download_test.py diff --git a/statvar_imports/us_cdc/single_race/README.md b/statvar_imports/us_cdc/single_race/README.md index 7b76c5985a..35f81fef0e 100644 --- a/statvar_imports/us_cdc/single_race/README.md +++ b/statvar_imports/us_cdc/single_race/README.md @@ -4,7 +4,7 @@ - Source URL: https://wonder.cdc.gov/ucd-icd10-expanded.html -- Import Type: Semi-Automated +- Import Type: Automated - Data Availability: 2018 onwards @@ -12,48 +12,40 @@ ### 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: +sh 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 : +After downloading, input files will be placed into the `input_files/` directory. The data is processed using the `stat_var_processor.py` script: ```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. - - -```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 to run automatically on the second Saturday of every month schedule. - Cron Expression: 30 08 8-14 * 6 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..50feba0b3e --- /dev/null +++ b/statvar_imports/us_cdc/single_race/download.py @@ -0,0 +1,655 @@ +# 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 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 retry 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'.", +) +flags.DEFINE_string( + "years", + "2018-2024", + "Year range ('2018-2024') or comma-separated years ('2018,2019,2020').", +) +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 __init__( + self, + landing_url: str = SOURCE_LANDING_URL, + timeout: int = 120, + delay: float = 2.0, + ): + self.landing_url = landing_url + self.timeout = timeout + self.delay = delay + self.session = requests.Session() + self.session.headers.update({ + "User-Agent": + "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" + }) + self.action_url: Optional[str] = None + self.base_post_data: List[Tuple[str, 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 = requests.Session() + self.session.headers.update({ + "User-Agent": + "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" + }) + self.action_url = None + self.base_post_data = [] + + 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 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: + 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: + res = self.session.post(self.action_url, + data=payload, + timeout=self.timeout) + if res.status_code == 429: + 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)...", + wait_time, + wait_time // 60, + attempt, + ) + time.sleep(wait_time) + # Re-initialize session to renew cookies and session ID + try: + self.init_session() + payload = self._build_post_data(state_fips, years) + except Exception as e: + logging.warning("Session re-initialization error: %s", + e) + continue + + if res.status_code == 400: + logging.warning( + "CDC WONDER returned HTTP 400 (likely query buffer overrun for large state). Returning for partitioning." + ) + return "CDC WONDER 400 Bad Request (query too large)" + + res.raise_for_status() + 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) + try: + self.init_session() + payload = self._build_post_data(state_fips, years) + except Exception as session_err: + logging.warning("Session re-initialization error: %s", + session_err) + + 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) + + need_partitioning = False + + 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) + need_partitioning = True + else: + # Try querying all years first + try: + 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)] + else: + logging.warning( + "%s response not TSV (likely exceeded 75k rows: %s). Partitioning into chunks...", + state_name, + response_text[:120].strip().replace("\n", " "), + ) + need_partitioning = True + except Exception as e: + logging.warning( + "Querying all %d years for %s encountered %s. Partitioning into year chunks...", + len(years), + state_name, + e, + ) + need_partitioning = True + + if need_partitioning: + # 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 {single_year}." + ) + chunk_results.append((single_year, sy_text)) + else: + raise ValueError( + f"Failed to query {state_name} for chunk {chunk_label}: {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 + + 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: + latest_year = years[-1] + has_chunk = any( + f.name.endswith(f"_{latest_year}.csv") + or f"_{latest_year}_" in f.name for f in matches) + if has_chunk: + return True + single_file = Path( + output_dir) / f"UnderlyingCauseofDeath_SingleRace_{state_fips}.csv" + if single_file.exists(): + content = single_file.read_text(encoding="utf-8", errors="replace") + return f",{latest_year}," in content + return False + return True + + +def download_single_race_data( + states: List[str], + years: List[str], + output_dir: str, + delay: float = 3.0, + timeout: int = 120, + skip_existing: bool = True, + batch_size: int = 10, + batch_cooldown: float = 20.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() + + 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, + ) + + 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 10 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(_): + 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..664c817311 100755 --- a/statvar_imports/us_cdc/single_race/download.sh +++ b/statvar_imports/us_cdc/single_race/download.sh @@ -2,5 +2,7 @@ set -e -o 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..8612bc9e7c --- /dev/null +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -0,0 +1,248 @@ +# 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): + + def test_parse_year_list_range(self): + years = download.parse_year_list("2018-2023") + self.assertEqual(years, + ["2018", "2019", "2020", "2021", "2022", "2023"]) + + def test_parse_year_list_comma(self): + years = download.parse_year_list("2018, 2020, 2022") + self.assertEqual(years, ["2018", "2020", "2022"]) + + def test_parse_year_list_single(self): + years = download.parse_year_list("2024") + self.assertEqual(years, ["2024"]) + + @mock.patch.object(download.requests, "Session") + def test_init_session_success(self, mock_session_cls): + 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(mock_session.post.call_count, 1) + + @mock.patch.object(download.requests, "Session") + def test_init_session_missing_form(self, mock_session_cls): + 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): + 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): + 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): + 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): + 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): + 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_is_state_downloaded(self): + 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 file > 100 bytes + f.write_text("Header,col1,col2,col3\n" + + "val1,val2,val3,val4\n" * 10) + self.assertTrue(download.is_state_downloaded(temp_dir, "02")) + + @mock.patch.object(download.time, "sleep") + @mock.patch.object(download.CdcWonderSingleRaceDownloader, "init_session") + def test_execute_query_429_backoff(self, mock_init, mock_sleep): + 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() + + +if __name__ == "__main__": + unittest.main() From bfdcf4391eaac1ac904175f5f6e8d71a6721f7e7 Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 29 Sep 2026 05:17:26 +0000 Subject: [PATCH 02/10] Address review comments for CDC Single Race import automation - Remove exception suppression during session re-initialization in 429 and general retry handlers. - Update is_state_downloaded to verify all requested years are covered across monolithic and chunked CSV files. - Expand unit test test_is_state_downloaded with multi-year and chunk coverage assertions. - Clarify Data Commons manifest cron schedule and standard POSIX crontab conditional wrapper in README.md. --- statvar_imports/us_cdc/single_race/README.md | 6 ++- .../us_cdc/single_race/download.py | 50 ++++++++++--------- .../us_cdc/single_race/download_test.py | 19 ++++++- 3 files changed, 47 insertions(+), 28 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/README.md b/statvar_imports/us_cdc/single_race/README.md index 35f81fef0e..19a986285e 100644 --- a/statvar_imports/us_cdc/single_race/README.md +++ b/statvar_imports/us_cdc/single_race/README.md @@ -45,7 +45,9 @@ python3 ../../../tools/statvar_importer/stat_var_processor.py --existing_statvar ### Automation -This import pipeline is configured to run automatically on the second Saturday of every month schedule. +This import pipeline is configured in `manifest.json` on the second Saturday of every month schedule: -- Cron Expression: 30 08 8-14 * 6 +- Data Commons Manifest Schedule: `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 ] && sh 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 index 50feba0b3e..eb428299c6 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -366,12 +366,8 @@ def execute_query( ) time.sleep(wait_time) # Re-initialize session to renew cookies and session ID - try: - self.init_session() - payload = self._build_post_data(state_fips, years) - except Exception as e: - logging.warning("Session re-initialization error: %s", - e) + self.init_session() + payload = self._build_post_data(state_fips, years) continue if res.status_code == 400: @@ -394,12 +390,8 @@ def execute_query( max_retries, ) time.sleep(wait_time) - try: - self.init_session() - payload = self._build_post_data(state_fips, years) - except Exception as session_err: - logging.warning("Session re-initialization error: %s", - session_err) + self.init_session() + payload = self._build_post_data(state_fips, years) raise RuntimeError( f"Failed to query {state_fips} after {max_retries} attempts.") @@ -532,18 +524,28 @@ def is_state_downloaded(output_dir: str, if not all(f.stat().st_size > 100 for f in matches): return False if years: - latest_year = years[-1] - has_chunk = any( - f.name.endswith(f"_{latest_year}.csv") - or f"_{latest_year}_" in f.name for f in matches) - if has_chunk: - return True - single_file = Path( - output_dir) / f"UnderlyingCauseofDeath_SingleRace_{state_fips}.csv" - if single_file.exists(): - content = single_file.read_text(encoding="utf-8", errors="replace") - return f",{latest_year}," in content - return False + covered_years = set() + for f in matches: + name = f.stem + parts = name.split("_") + if len(parts) == 3: + try: + content = f.read_text(encoding="utf-8", errors="replace") + for y in years: + if f",{y}," in content: + covered_years.add(y) + 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 diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 8612bc9e7c..3e88a3f6ce 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -213,10 +213,25 @@ def test_is_state_downloaded(self): f.write_text("") self.assertFalse(download.is_state_downloaded(temp_dir, "02")) - # Create valid file > 100 bytes + # Create valid monolithic file covering 2018 and 2019 f.write_text("Header,col1,col2,col3\n" + - "val1,val2,val3,val4\n" * 10) + ",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"])) + + # 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") From c20aade5a6a2d239254713bf23c7733c6ae82a5c Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 29 Sep 2026 11:16:32 +0000 Subject: [PATCH 03/10] Address review findings for CDC Single Race import automation - Parse Year column explicitly in is_state_downloaded() to eliminate false positives from numeric data in other columns. - Ensure consistent explicit return in download_state() to avoid implicit None return. - Add pre-dispatch outbound HTTP request logging and response telemetry in execute_query(). - Add data row validation in save_tsv_as_csv() to fail on empty/header-only responses. - Harmonize flag and function parameter defaults for delay, batch_size, and batch_cooldown. - Wrap long lines to adhere strictly to <= 100 char limit and add complete docstrings across download_test.py. --- .../us_cdc/single_race/download.py | 188 +++++++++++------- .../us_cdc/single_race/download_test.py | 33 ++- 2 files changed, 145 insertions(+), 76 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/download.py b/statvar_imports/us_cdc/single_race/download.py index eb428299c6..3567461b71 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -191,7 +191,7 @@ def __init__( self, landing_url: str = SOURCE_LANDING_URL, timeout: int = 120, - delay: float = 2.0, + delay: float = 5.0, ): self.landing_url = landing_url self.timeout = timeout @@ -199,7 +199,7 @@ def __init__( self.session = requests.Session() self.session.headers.update({ "User-Agent": - "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" + "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" }) self.action_url: Optional[str] = None self.base_post_data: List[Tuple[str, str]] = [] @@ -218,7 +218,7 @@ def init_session(self): self.session = requests.Session() self.session.headers.update({ "User-Agent": - "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" + "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" }) self.action_url = None self.base_post_data = [] @@ -236,7 +236,8 @@ def init_session(self): 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")] + 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)...") @@ -265,8 +266,8 @@ def init_session(self): continue if itype in ["checkbox", "radio"]: if el.has_attr("checked"): - self.base_post_data.append( - (name, el.get("value", "on"))) + self.base_post_data.append((name, el.get("value", + "on"))) else: self.base_post_data.append((name, el.get("value", ""))) elif el.name == "select": @@ -276,8 +277,7 @@ def init_session(self): ] if selected_opts: for opt in selected_opts: - self.base_post_data.append((name, opt.get("value", - ""))) + self.base_post_data.append((name, opt.get("value", ""))) else: if not el.has_attr("multiple"): first_opt = el.find("option") @@ -345,21 +345,32 @@ def execute_query( 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: 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." + # "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. + # 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)...", + "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)...", wait_time, wait_time // 60, attempt, @@ -372,18 +383,27 @@ def execute_query( if res.status_code == 400: logging.warning( - "CDC WONDER returned HTTP 400 (likely query buffer overrun for large state). Returning for partitioning." + "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)...", + "Request error: %s. Renewing session and retrying in %d seconds " + "(attempt %d/%d)...", e, wait_time, attempt, @@ -427,65 +447,67 @@ def download_state(self, state_fips: str, "Successfully fetched %s (all requested years in 1 query).", state_name) return [("all", response_text)] - else: - logging.warning( - "%s response not TSV (likely exceeded 75k rows: %s). Partitioning into chunks...", - state_name, - response_text[:120].strip().replace("\n", " "), - ) - need_partitioning = True + logging.warning( + "%s response not TSV (likely exceeded 75k rows: %s). " + "Partitioning into chunks...", + state_name, + response_text[:120].strip().replace("\n", " "), + ) + need_partitioning = True except Exception as e: logging.warning( - "Querying all %d years for %s encountered %s. Partitioning into year chunks...", + "Querying all %d years for %s encountered %s. " + "Partitioning into year chunks...", len(years), state_name, e, ) need_partitioning = True - if need_partitioning: - # 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 {single_year}." - ) - chunk_results.append((single_year, sy_text)) - else: - raise ValueError( - f"Failed to query {state_name} for chunk {chunk_label}: {chunk_text[:300]}" - ) + if not need_partitioning: + return [] + + # 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: - chunk_results.append((chunk_label, chunk_text)) + 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 + return chunk_results def save_tsv_as_csv(raw_tsv: str, output_filepath: str) -> int: @@ -501,12 +523,19 @@ def save_tsv_as_csv(raw_tsv: str, output_filepath: str) -> int: if not row: continue # Stop at metadata notes footer - if row[0].startswith("---") or (len(row) > 1 - and row[1].startswith("---")): + 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 @@ -530,10 +559,19 @@ def is_state_downloaded(output_dir: str, parts = name.split("_") if len(parts) == 3: try: - content = f.read_text(encoding="utf-8", errors="replace") - for y in years: - if f",{y}," in content: - covered_years.add(y) + 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: @@ -553,11 +591,11 @@ def download_single_race_data( states: List[str], years: List[str], output_dir: str, - delay: float = 3.0, + delay: float = 5.0, timeout: int = 120, skip_existing: bool = True, - batch_size: int = 10, - batch_cooldown: float = 20.0, + 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) @@ -611,10 +649,11 @@ def download_single_race_data( states_in_batch += 1 - # Automatically refresh session after each batch of 10 states + # 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...", + "Completed session batch of %d states. Cooling down for %.1fs and renewing " + "CDC session...", states_in_batch, batch_cooldown, ) @@ -629,6 +668,7 @@ def download_single_race_data( def main(_): + """Main entry point to execute the CDC WONDER downloader.""" years = parse_year_list(FLAGS.years) if FLAGS.states.lower() == "all": diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 3e88a3f6ce..163fb914df 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -31,22 +31,27 @@ 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.requests, "Session") def test_init_session_success(self, mock_session_cls): + """Tests successful session handshake and parameter extraction.""" mock_session = mock.MagicMock() mock_session_cls.return_value = mock_session @@ -91,6 +96,7 @@ def test_init_session_success(self, mock_session_cls): @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 @@ -105,6 +111,7 @@ def test_init_session_missing_form(self, mock_session_cls): 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"), @@ -134,6 +141,7 @@ def test_build_post_data(self): @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 @@ -148,6 +156,7 @@ def test_download_state_single_query(self, mock_query): @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.") @@ -170,6 +179,7 @@ def test_download_state_row_limit_partitioning(self, mock_query): @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" @@ -184,6 +194,7 @@ def test_download_state_large_state_preemptive(self, mock_query): 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" @@ -201,10 +212,18 @@ def test_save_tsv_as_csv(self): 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') + 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")) @@ -222,6 +241,15 @@ def test_is_state_downloaded(self): 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) / @@ -236,6 +264,7 @@ def test_is_state_downloaded(self): @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")] From 367d4c62cc6e706bff15f939efdbb5725ac92ce4 Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 04:10:04 +0000 Subject: [PATCH 04/10] Address CRA review findings for CDC single-race data acquisition - Purge preexisting state CSVs in download_single_race_data to avoid duplicate observations. - Enforce immediate fail-fast on HTTP 429 when attempt reaches max_retries. - Bubble network/runtime exceptions in download_state instead of masking as row limit partitions. - Set dynamic end-year default to current year for automated annual rollovers. - Add connection pooling and unified SSL via HTTPAdapter and urllib3 Retry. - Add comprehensive unit tests in download_test.py covering download, purge, and 429 backoff. - Update download.sh to set -euo pipefail and document bash download.sh in README.md. - Remove dead code in download_state and wrap lines exceeding 100 characters in README.md. --- statvar_imports/us_cdc/single_race/README.md | 32 +++-- .../us_cdc/single_race/download.py | 98 ++++++++------- .../us_cdc/single_race/download.sh | 2 +- .../us_cdc/single_race/download_test.py | 112 ++++++++++++++++++ 4 files changed, 195 insertions(+), 49 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/README.md b/statvar_imports/us_cdc/single_race/README.md index 19a986285e..3c0888ace8 100644 --- a/statvar_imports/us_cdc/single_race/README.md +++ b/statvar_imports/us_cdc/single_race/README.md @@ -1,6 +1,7 @@ ### 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 @@ -14,19 +15,23 @@ - Download: Automated live downloader (`download.py`) -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: +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 -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. +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. To run the live download: ```bash # Execute via shell wrapper: -sh download.sh +bash download.sh # Or directly with python: python3 download.py @@ -37,17 +42,28 @@ python3 download.py --states=02,48 --years=2018-2024 ### Data Processing -After downloading, input files will be placed into the `input_files/` directory. The data is processed using the `stat_var_processor.py` script: +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 in `manifest.json` 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` -*(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 ] && sh download.sh` to restrict execution strictly to the second Saturday).* +*(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 index 3567461b71..0e7fc8b017 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -29,6 +29,7 @@ """ import csv +import datetime import io import os from pathlib import Path @@ -41,7 +42,9 @@ 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") @@ -137,10 +140,13 @@ "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", - "2018-2024", - "Year range ('2018-2024') or comma-separated years ('2018,2019,2020').", + _DEFAULT_YEARS, + f"Year range (e.g. '2018-{_CURRENT_YEAR}') or comma-separated years.", ) flags.DEFINE_string( "output_dir", @@ -187,6 +193,28 @@ def parse_year_list(year_str: str) -> List[str]: 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, @@ -196,11 +224,7 @@ def __init__( self.landing_url = landing_url self.timeout = timeout self.delay = delay - self.session = requests.Session() - self.session.headers.update({ - "User-Agent": - "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" - }) + self.session = self._create_session() self.action_url: Optional[str] = None self.base_post_data: List[Tuple[str, str]] = [] @@ -215,11 +239,7 @@ def init_session(self): self.session.close() except Exception: pass - self.session = requests.Session() - self.session.headers.update({ - "User-Agent": - "Mozilla/5.0 (DataCommons CDC Importer; contact: support@datacommons.org)" - }) + self.session = self._create_session() self.action_url = None self.base_post_data = [] @@ -357,6 +377,8 @@ def execute_query( 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" @@ -370,10 +392,11 @@ def execute_query( 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)...", + "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 @@ -439,33 +462,20 @@ def download_state(self, state_fips: str, need_partitioning = True else: # Try querying all years first - try: - 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", " "), - ) - need_partitioning = True - except Exception as e: - logging.warning( - "Querying all %d years for %s encountered %s. " - "Partitioning into year chunks...", - len(years), - state_name, - e, - ) - need_partitioning = True - - if not need_partitioning: - return [] + 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", " "), + ) + need_partitioning = True # Partition years into 2-year chunks chunk_results = [] @@ -634,6 +644,14 @@ def download_single_race_data( 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: diff --git a/statvar_imports/us_cdc/single_race/download.sh b/statvar_imports/us_cdc/single_race/download.sh index 664c817311..cf855e664c 100755 --- a/statvar_imports/us_cdc/single_race/download.sh +++ b/statvar_imports/us_cdc/single_race/download.sh @@ -1,6 +1,6 @@ #!/bin/bash -set -e -o pipefail +set -euo pipefail SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" mkdir -p "${SCRIPT_DIR}/input_files" diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 163fb914df..77373d3862 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -287,6 +287,118 @@ def test_execute_query_429_backoff(self, mock_init, mock_sleep): 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.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() + if __name__ == "__main__": unittest.main() From 48c155b471d66daf822205028150e3989b3acf5d Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 05:06:04 +0000 Subject: [PATCH 05/10] Update validation_config.json with per-StatVar date freshness and MAX_DATE_CONSISTENT --- .../us_cdc/single_race/validation_config.json | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/validation_config.json b/statvar_imports/us_cdc/single_race/validation_config.json index 9ccffb978e..a67607deba 100644 --- a/statvar_imports/us_cdc/single_race/validation_config.json +++ b/statvar_imports/us_cdc/single_race/validation_config.json @@ -3,12 +3,17 @@ "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 each StatVar reaches at least 2024 for CDC WONDER single-race mortality data.", "validator": "SQL_VALIDATOR", "params": { - "query": "SELECT CAST(MAX(MaxDate) AS INTEGER) AS latest_year FROM stats", - "condition": "latest_year >= 2024" + "query": "SELECT StatVar, MaxDate, COUNT(*) OVER () AS total_svs FROM stats", + "condition": "MaxDate >= '2024' AND total_svs > 0" } + }, + { + "rule_id": "check_max_date_consistent", + "description": "Checks if the MaxDate is the same for all StatVars.", + "validator": "MAX_DATE_CONSISTENT" } ] } From 7aae6518688e3bd4570d173050ccb1fd50bf200a Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 05:59:38 +0000 Subject: [PATCH 06/10] Filter requested years against CDC WONDER published years to prevent HTTP 500 errors --- .../us_cdc/single_race/download.py | 35 ++++++++++++++++++- .../us_cdc/single_race/download_test.py | 23 ++++++++++++ 2 files changed, 57 insertions(+), 1 deletion(-) diff --git a/statvar_imports/us_cdc/single_race/download.py b/statvar_imports/us_cdc/single_race/download.py index 0e7fc8b017..336a182c79 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -227,6 +227,7 @@ def __init__( 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, @@ -242,6 +243,7 @@ def init_session(self): 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) @@ -274,6 +276,19 @@ def init_session(self): 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: + self.available_years = [] + # Extract pre-populated query parameters self.base_post_data = [] for el in form_req.find_all(["input", "select", "textarea"]): @@ -346,7 +361,8 @@ def _build_post_data( if years: for y in years: - query_data.append(("F_D158.V1", y)) + 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 @@ -612,6 +628,23 @@ def download_single_race_data( 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 diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 77373d3862..2a60eebc70 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -399,6 +399,29 @@ def test_download_single_race_data_skip_existing(self, mock_download_state, mock_download_state.assert_not_called() + @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() From 894147372ebc4ce86a432c84845c7cebefc88f3a Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 08:46:56 +0000 Subject: [PATCH 07/10] Update freshness check to aggregate MAX date and remove MAX_DATE_CONSISTENT for sparse mortality data --- .../us_cdc/single_race/validation_config.json | 11 +++-------- 1 file changed, 3 insertions(+), 8 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/validation_config.json b/statvar_imports/us_cdc/single_race/validation_config.json index a67607deba..ef155a81ab 100644 --- a/statvar_imports/us_cdc/single_race/validation_config.json +++ b/statvar_imports/us_cdc/single_race/validation_config.json @@ -3,17 +3,12 @@ "rules": [ { "rule_id": "check_latest_date_freshness", - "description": "Asserts each StatVar reaches at least 2024 for CDC WONDER single-race mortality data.", + "description": "Asserts latest observation date reaches at least 2024 across StatVars.", "validator": "SQL_VALIDATOR", "params": { - "query": "SELECT StatVar, MaxDate, COUNT(*) OVER () AS total_svs FROM stats", - "condition": "MaxDate >= '2024' AND total_svs > 0" + "query": "SELECT MAX(TRY_CAST(MaxDate AS INT)) AS max_year FROM stats", + "condition": "max_year >= 2024" } - }, - { - "rule_id": "check_max_date_consistent", - "description": "Checks if the MaxDate is the same for all StatVars.", - "validator": "MAX_DATE_CONSISTENT" } ] } From 5475674880d918e54c4272216fabec924f7a691f Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 11:40:45 +0000 Subject: [PATCH 08/10] Address remaining CRA report findings: 3-way freshness gate, remove unused need_partitioning and assert mock_init --- statvar_imports/us_cdc/single_race/download.py | 4 ---- statvar_imports/us_cdc/single_race/download_test.py | 2 ++ statvar_imports/us_cdc/single_race/validation_config.json | 6 +++--- 3 files changed, 5 insertions(+), 7 deletions(-) diff --git a/statvar_imports/us_cdc/single_race/download.py b/statvar_imports/us_cdc/single_race/download.py index 336a182c79..b3bc560b37 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -469,13 +469,10 @@ def download_state(self, state_fips: str, logging.info("Querying data for %s (FIPS %s) for years %s...", state_name, state_fips, years) - need_partitioning = False - 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) - need_partitioning = True else: # Try querying all years first response_text = self.execute_query(state_fips, years) @@ -491,7 +488,6 @@ def download_state(self, state_fips: str, state_name, response_text[:120].strip().replace("\n", " "), ) - need_partitioning = True # Partition years into 2-year chunks chunk_results = [] diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 2a60eebc70..8a4d3ad178 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -377,6 +377,7 @@ def test_download_single_race_data_purges_preexisting_files( "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") @@ -398,6 +399,7 @@ def test_download_single_race_data_skip_existing(self, mock_download_state, ) 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( diff --git a/statvar_imports/us_cdc/single_race/validation_config.json b/statvar_imports/us_cdc/single_race/validation_config.json index ef155a81ab..baceeb2586 100644 --- a/statvar_imports/us_cdc/single_race/validation_config.json +++ b/statvar_imports/us_cdc/single_race/validation_config.json @@ -3,11 +3,11 @@ "rules": [ { "rule_id": "check_latest_date_freshness", - "description": "Asserts latest observation date reaches at least 2024 across StatVars.", + "description": "Asserts data reaches 2024 with >= 690 StatVars reaching 2024.", "validator": "SQL_VALIDATOR", "params": { - "query": "SELECT MAX(TRY_CAST(MaxDate AS INT)) AS max_year FROM stats", - "condition": "max_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" } } ] From d0ecebfa6ed6f25c20177da163b5ed0ec9a5bc7f Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Tue, 6 Oct 2026 11:55:21 +0000 Subject: [PATCH 09/10] Add check_deleted_records_percent (0.1% threshold) to validation_config.json --- statvar_imports/us_cdc/single_race/validation_config.json | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/statvar_imports/us_cdc/single_race/validation_config.json b/statvar_imports/us_cdc/single_race/validation_config.json index baceeb2586..99f939f8f9 100644 --- a/statvar_imports/us_cdc/single_race/validation_config.json +++ b/statvar_imports/us_cdc/single_race/validation_config.json @@ -9,6 +9,14 @@ "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 + } } ] } From 94646c3fffd20fc43f136896c73bcf671086540b Mon Sep 17 00:00:00 2001 From: shvngisingh Date: Thu, 8 Oct 2026 09:08:59 +0000 Subject: [PATCH 10/10] Log warning when year select element F_D158.V1 is missing in init_session() --- statvar_imports/us_cdc/single_race/download.py | 3 +++ statvar_imports/us_cdc/single_race/download_test.py | 7 ++++++- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/statvar_imports/us_cdc/single_race/download.py b/statvar_imports/us_cdc/single_race/download.py index b3bc560b37..f1a5281ab3 100755 --- a/statvar_imports/us_cdc/single_race/download.py +++ b/statvar_imports/us_cdc/single_race/download.py @@ -287,6 +287,9 @@ def init_session(self): 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 diff --git a/statvar_imports/us_cdc/single_race/download_test.py b/statvar_imports/us_cdc/single_race/download_test.py index 8a4d3ad178..28ee0f3294 100644 --- a/statvar_imports/us_cdc/single_race/download_test.py +++ b/statvar_imports/us_cdc/single_race/download_test.py @@ -49,8 +49,9 @@ def test_parse_year_list_single(self): 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): + 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 @@ -92,6 +93,10 @@ def test_init_session_success(self, mock_session_cls): 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")