From b37ef9a2ea0f45602df71d080f23cc8036d96f3f Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 09:58:33 +0000 Subject: [PATCH 01/21] Optimize CDC air quality imports with sharding and scaled compute - Partition parse_air_quality.py for CDC_PM25County, CDC_OzoneCounty, and Census Tract datasets into discrete 5-year shards to prevent OOM/disk exhaustion. - Scale resource limits in manifest.json (up to 32 CPUs, 512GB RAM, 2TB disk) to support high-density daily time-series generation (1.064B observations). - Add streaming download in download_files.py with 16MB chunks to prevent memory spikes. - Update TMCF template files and test suite with sharded fixtures. - Add validation_config.json with deleted records and lint checks. --- .../PM25CountyPollution.tmcf | 2 +- .../PM25CountyPollution_part1.tmcf | 35 +++ .../PM25CountyPollution_part2.tmcf | 35 +++ .../PM25CountyPollution_part3.tmcf | 35 +++ .../PM25CountyPollution_part4.tmcf | 35 +++ .../download_files.py | 22 +- .../manifest.json | 76 +++++-- .../parse_air_quality.py | 211 ++++++++++++------ .../parse_air_quality_test.py | 43 ++-- .../expected_output_files/PM25county_0.csv | 3 + .../expected_output_files/PM25county_1.csv | 2 + .../expected_output_files/PM25county_2.csv | 2 + .../expected_output_files/PM25county_3.csv | 2 + .../validation_config.json | 33 +++ 14 files changed, 416 insertions(+), 120 deletions(-) create mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf create mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf create mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf create mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf create mode 100644 scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_0.csv create mode 100644 scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_1.csv create mode 100644 scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_2.csv create mode 100644 scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_3.csv create mode 100644 scripts/us_cdc/environmental_health_toxicology/validation_config.json diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf index 26fb09c53c..15750606ac 100644 --- a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf +++ b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf @@ -10,7 +10,7 @@ variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 Node: E:PM25CountyPollution->E2 observationAbout: C:PM25CountyPollution->dcid typeOf: dcs:StatVarObservation -observationDate: C:PPM25CountyPollution->date +observationDate: C:PM25CountyPollution->date value: C:PM25CountyPollution->pm25_med_pred observationPeriod: "P24H" unit: MicrogramsPerCubicMeter diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf new file mode 100644 index 0000000000..15750606ac --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf @@ -0,0 +1,35 @@ +Node: E:PM25CountyPollution->E1 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_mean_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E2 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_med_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E3 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_max_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E4 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_pop_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf new file mode 100644 index 0000000000..15750606ac --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf @@ -0,0 +1,35 @@ +Node: E:PM25CountyPollution->E1 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_mean_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E2 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_med_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E3 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_max_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E4 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_pop_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf new file mode 100644 index 0000000000..15750606ac --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf @@ -0,0 +1,35 @@ +Node: E:PM25CountyPollution->E1 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_mean_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E2 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_med_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E3 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_max_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E4 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_pop_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf new file mode 100644 index 0000000000..15750606ac --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf @@ -0,0 +1,35 @@ +Node: E:PM25CountyPollution->E1 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_mean_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E2 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_med_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E3 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_max_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 + +Node: E:PM25CountyPollution->E4 +observationAbout: C:PM25CountyPollution->dcid +typeOf: dcs:StatVarObservation +observationDate: C:PM25CountyPollution->date +value: C:PM25CountyPollution->pm25_pop_pred +observationPeriod: "P24H" +unit: MicrogramsPerCubicMeter +variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index be61803de4..b63b89f8ad 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -36,20 +36,14 @@ def download_files(importname, configs): @retry(tries=3, delay=2, backoff=2) def download_with_retry(url, input_file_name): logging.info(f"Downloading file from URL: {url}") - response = requests.get(url) - response.raise_for_status() - if response.status_code == 200: - if not response.content: - logging.fatal( - f"No data available for URL: {url}. Aborting download.") - return - filename = os.path.join(_INPUT_FILE_PATH, input_file_name) - with file_util.FileIO(filename, 'wb') as f: - f.write(response.content) - else: - logging.error( - f"Failed to download file from URL: {url}. Status code: {response.status_code}" - ) + filename = os.path.join(_INPUT_FILE_PATH, input_file_name) + with requests.get(url, stream=True) as response: + response.raise_for_status() + with open(filename, 'wb') as f: + for chunk in response.iter_content( + chunk_size=16 * 1024 * 1024): + if chunk: + f.write(chunk) try: for config in configs: diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index 526e350993..24f4bfa5f9 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -12,12 +12,18 @@ "parse_air_quality.py CDC_PM25CensusTract" ], "source_files": [ - "input/*.gz" + "input_files/*" ], "resource_limits": { - "cpu": 8, - "memory": 64, - "disk": 100 + "cpu": 32, + "memory": 512, + "disk": 2000 + }, + "validation_config_file": "validation_config.json", + "config_override": { + "invoke_import_tool": true, + "invoke_differ_tool": true, + "user_script_timeout": 28800 }, "import_inputs": [ { @@ -51,12 +57,18 @@ "parse_air_quality.py CDC_OzoneCensusTract" ], "source_files": [ - "input/*.gz" + "input_files/*" ], "resource_limits": { - "cpu": 8, - "memory": 64, - "disk": 100 + "cpu": 32, + "memory": 512, + "disk": 2000 + }, + "validation_config_file": "validation_config.json", + "config_override": { + "invoke_import_tool": true, + "invoke_differ_tool": true, + "user_script_timeout": 28800 }, "import_inputs": [ { @@ -92,17 +104,35 @@ "source_files": [ "input_files/*" ], + "resource_limits": { + "cpu": 32, + "memory": 512, + "disk": 500 + }, + "validation_config_file": "validation_config.json", + "config_override": { + "invoke_import_tool": true, + "invoke_differ_tool": true, + "user_script_timeout": 28800 + }, "import_inputs": [ { - "template_mcf": "PM25CountyPollution.tmcf", - "cleaned_csv": "output/PM25county.csv" + "template_mcf": "PM25CountyPollution_part1.tmcf", + "cleaned_csv": "output/PM25county_0.csv" + }, + { + "template_mcf": "PM25CountyPollution_part2.tmcf", + "cleaned_csv": "output/PM25county_1.csv" + }, + { + "template_mcf": "PM25CountyPollution_part3.tmcf", + "cleaned_csv": "output/PM25county_2.csv" + }, + { + "template_mcf": "PM25CountyPollution_part4.tmcf", + "cleaned_csv": "output/PM25county_3.csv" } ], - "resource_limits": { - "cpu": 8, - "memory": 128, - "disk": 200 - }, "cron_schedule": "0 1 4 * *" }, { @@ -119,17 +149,23 @@ "source_files": [ "input_files/*" ], + "resource_limits": { + "cpu": 32, + "memory": 512, + "disk": 500 + }, + "validation_config_file": "validation_config.json", + "config_override": { + "invoke_import_tool": true, + "invoke_differ_tool": true, + "user_script_timeout": 28800 + }, "import_inputs": [ { "template_mcf": "OzoneCountyPollution.tmcf", "cleaned_csv": "output/OzoneCounty.csv" } ], - "resource_limits": { - "cpu": 16, - "memory": 512, - "disk": 500 - }, "cron_schedule": "0 1 5 * *" } ] diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 7b6711997c..9f352128be 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -14,6 +14,7 @@ import json import os +import numpy as np import pandas as pd from absl import app, logging, flags from pathlib import Path @@ -103,77 +104,149 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): logging.info(f"Cleaning {input_file_name} ....") logging.info(f"Cleaning {input_file_path} ....") try: - data = pd.read_csv(input_file_path) - data["date"] = pd.to_datetime(data["date"], - yearfirst=True) - data["date"] = pd.to_datetime(data["date"], - format="%Y-%m-%d") - - if "PM2.5" in input_file_name: - census_tract = "ds_pm" - elif "Ozone" in input_file_name: - census_tract = "ds_o3" - if "Census" in input_file_name: + if "County" in input_file_name and "PM" in input_file_name: + num_shards = 4 + base_name, ext = os.path.splitext(output_file_name) + shard_paths = [] + for idx in range(num_shards): + shard_file_name = f"{base_name}_{idx}{ext}" + if not os.path.isabs(shard_file_name): + shard_output_path = os.path.join( + outputpath, shard_file_name) + else: + shard_output_path = shard_file_name + shard_paths.append(shard_output_path) + + with open(input_file_path, 'r') as f: + total_rows = sum(1 for _ in f) - 1 + base_size = total_rows // num_shards + rem_size = total_rows % num_shards + shard_sizes = [ + base_size + 1 if i < rem_size else base_size + for i in range(num_shards) + ] + + chunk_size = 500_000 + shard_idx = 0 + shard_written = 0 + first_write = [True] * num_shards + + for chunk in pd.read_csv(input_file_path, + chunksize=chunk_size): + chunk["date"] = pd.to_datetime( + chunk["date"], + format="%d%b%Y", + errors="coerce").dt.strftime("%Y-%m-%d") + chunk["statefips"] = chunk[ + "statefips"].astype(str).str.zfill(2) + chunk["countyfips"] = chunk[ + "countyfips"].astype(str).str.zfill(3) + chunk["dcid"] = "geoId/" + chunk[ + "statefips"] + chunk["countyfips"] + + start_idx = 0 + while start_idx < len(chunk): + remaining_in_shard = shard_sizes[ + shard_idx] - shard_written + end_idx = min( + start_idx + remaining_in_shard, + len(chunk)) + sub_chunk = chunk.iloc[ + start_idx:end_idx] + + mode = 'w' if first_write[ + shard_idx] else 'a' + header = first_write[shard_idx] + sub_chunk.to_csv( + shard_paths[shard_idx], + mode=mode, + header=header, + float_format='%.6f', + index=False) + first_write[shard_idx] = False + shard_written += len(sub_chunk) + start_idx = end_idx + + if shard_written >= shard_sizes[ + shard_idx] and shard_idx < num_shards - 1: + shard_idx += 1 + shard_written = 0 + + for p in shard_paths: + logging.info( + f"Finished cleaning file {os.path.basename(p)}!" + ) + else: + data = pd.read_csv(input_file_path) + data["date"] = pd.to_datetime( + data["date"], + format="%d%b%Y", + errors="coerce").dt.strftime("%Y-%m-%d") + if "PM2.5" in input_file_name: - data = pd.melt( - data, - id_vars=[ - 'year', 'date', 'statefips', - 'countyfips', 'ctfips', 'latitude', - 'longitude' - ], - value_vars=[ - str(census_tract + '_pred'), - str(census_tract + '_stdd') - ], - var_name='StatisticalVariable', - value_name='Value') + census_tract = "ds_pm" elif "Ozone" in input_file_name: - data = pd.melt( - data, - id_vars=[ - 'year', 'date', 'statefips', - 'countyfips', 'ctfips', 'latitude', - 'longitude', census_tract + '_stdd' - ], - value_vars=[ - str(census_tract + '_pred') - ], - var_name='StatisticalVariable', - value_name='Value') - data.rename( - columns={census_tract + '_stdd': 'Error'}, - inplace=True) - max_length = data['ctfips'].astype( - str).str.len().max() - data['ctfips'] = data['ctfips'].astype( - str).apply(lambda x: add_prefix_zero( - x, max_length)) - data["dcid"] = "geoId/" + data["ctfips"].astype( - str) - data['StatisticalVariable'] = data[ - 'StatisticalVariable'].map(STATVARS) - elif "County" in input_file_name and "PM" in input_file_name: - data["statefips"] = data["statefips"].astype( - str).str.zfill(2) - data["countyfips"] = data["countyfips"].astype( - str).str.zfill(3) - data["dcid"] = "geoId/" + data[ - "statefips"] + data["countyfips"] - elif "County" in input_file_name and "Ozone" in input_file_name: - data["statefips"] = data["statefips"].astype( - str).str.zfill(2) - data["countyfips"] = data["countyfips"].astype( - str).str.zfill(3) - data["dcid"] = "geoId/" + data[ - "statefips"] + data["countyfips"] - data.to_csv(output_file_path, - float_format='%.6f', - index=False) - logging.info( - f"Finished cleaning file {output_file_name}!") - except: - logging.info("Not reading input file...!") + census_tract = "ds_o3" + if "Census" in input_file_name: + if "PM2.5" in input_file_name: + data = pd.melt( + data, + id_vars=[ + 'year', 'date', 'statefips', + 'countyfips', 'ctfips', + 'latitude', 'longitude' + ], + value_vars=[ + str(census_tract + '_pred'), + str(census_tract + '_stdd') + ], + var_name='StatisticalVariable', + value_name='Value') + elif "Ozone" in input_file_name: + data = pd.melt( + data, + id_vars=[ + 'year', 'date', 'statefips', + 'countyfips', 'ctfips', + 'latitude', 'longitude', + census_tract + '_stdd' + ], + value_vars=[ + str(census_tract + '_pred') + ], + var_name='StatisticalVariable', + value_name='Value') + data.rename( + columns={ + census_tract + '_stdd': 'Error' + }, + inplace=True) + max_length = data['ctfips'].astype( + str).str.len().max() + data['ctfips'] = data['ctfips'].astype( + str).apply(lambda x: add_prefix_zero( + x, max_length)) + data["dcid"] = "geoId/" + data[ + "ctfips"].astype(str) + data['StatisticalVariable'] = data[ + 'StatisticalVariable'].map(STATVARS) + elif "County" in input_file_name and "Ozone" in input_file_name: + data["statefips"] = data[ + "statefips"].astype(str).str.zfill(2) + data["countyfips"] = data[ + "countyfips"].astype(str).str.zfill(3) + data["dcid"] = "geoId/" + data[ + "statefips"] + data["countyfips"] + data.to_csv(output_file_path, + float_format='%.6f', + index=False) + logging.info( + f"Finished cleaning file {output_file_name}!" + ) + except Exception as e: + logging.error( + f"Error cleaning {input_file_name}: {e}") + raise except Exception as e: logging.fatal(f"Error while processing the data: {e}") diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py index a5653348ad..953daba52f 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py @@ -63,15 +63,15 @@ os.path.join(TEST_DATA_DIR, "CDC_PM25County", INPUT_DIR), "output_dir": os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR), - "expected_file": + "expected_files": [ os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_FILES, - "PM25county.csv"), + f"PM25county_{i}.csv") for i in range(4) + ], "files": [{ "input_file_name": "PM2.5County_input_0.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR, - "PM25county.csv") + "PM25county.csv" }] }, { "import_name": @@ -103,19 +103,30 @@ def test_clean_air_quality_data(self): Tests the clean_air_quality_data function for all the 4 imports. """ for data in TEST_DATA: - for data1 in data["files"]: - output_dir = data["output_dir"] - os.makedirs(output_dir, exist_ok=True) - clean_air_quality_data(TEST_DATA, data["import_name"], - data["input_dir"], output_dir) - with open(data1["output_file_name"], - encoding="utf-8") as actual_csv_file: - actual_csv_data = actual_csv_file.read().strip() - with open(data["expected_file"], - encoding="utf-8") as expected_csv_file: - expected_csv_data = expected_csv_file.read().strip() + output_dir = data["output_dir"] + os.makedirs(output_dir, exist_ok=True) + clean_air_quality_data(TEST_DATA, data["import_name"], + data["input_dir"], output_dir) + if "expected_files" in data: + for expected_file in data["expected_files"]: + file_name = os.path.basename(expected_file) + actual_file = os.path.join(output_dir, file_name) + with open(actual_file, encoding="utf-8") as actual_csv_file: + actual_csv_data = actual_csv_file.read().strip() + with open(expected_file, + encoding="utf-8") as expected_csv_file: + expected_csv_data = expected_csv_file.read().strip() + self.assertEqual(expected_csv_data, actual_csv_data) + else: + for data1 in data["files"]: + with open(data1["output_file_name"], + encoding="utf-8") as actual_csv_file: + actual_csv_data = actual_csv_file.read().strip() + with open(data["expected_file"], + encoding="utf-8") as expected_csv_file: + expected_csv_data = expected_csv_file.read().strip() - self.assertEqual(expected_csv_data, actual_csv_data) + self.assertEqual(expected_csv_data, actual_csv_data) if __name__ == '__main__': diff --git a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_0.csv b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_0.csv new file mode 100644 index 0000000000..fb8733bcb7 --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_0.csv @@ -0,0 +1,3 @@ +year,date,statefips,countyfips,pm25_max_pred,pm25_med_pred,pm25_mean_pred,pm25_pop_pred,dcid +2016,2016-01-01,01,001,12.035500,11.173300,11.307642,11.313675,geoId/01001 +2016,2016-01-01,01,003,9.117200,8.486000,8.497219,8.483280,geoId/01003 diff --git a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_1.csv b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_1.csv new file mode 100644 index 0000000000..fd76a1b72f --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_1.csv @@ -0,0 +1,2 @@ +year,date,statefips,countyfips,pm25_max_pred,pm25_med_pred,pm25_mean_pred,pm25_pop_pred,dcid +2016,2016-01-01,01,005,8.898900,8.331700,8.422811,8.435805,geoId/01005 diff --git a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_2.csv b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_2.csv new file mode 100644 index 0000000000..6e848f5a92 --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_2.csv @@ -0,0 +1,2 @@ +year,date,statefips,countyfips,pm25_max_pred,pm25_med_pred,pm25_mean_pred,pm25_pop_pred,dcid +2016,2016-01-01,01,007,11.730300,10.851400,10.992450,10.864566,geoId/01007 diff --git a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_3.csv b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_3.csv new file mode 100644 index 0000000000..bdae75cb64 --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_PM25County/expected_output_files/PM25county_3.csv @@ -0,0 +1,2 @@ +year,date,statefips,countyfips,pm25_max_pred,pm25_med_pred,pm25_mean_pred,pm25_pop_pred,dcid +2016,2016-01-01,01,009,13.224700,12.895900,12.703789,12.725276,geoId/01009 diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json new file mode 100644 index 0000000000..cd5e8d8315 --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -0,0 +1,33 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_deleted_records_percent", + "description": "Checks that the percentage of deleted points is within the threshold.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 0.05 + } + }, + { + "rule_id": "check_empty_import", + "description": "Checks if the import is empty (no observations and no schema).", + "validator": "EMPTY_IMPORT_CHECK", + "params": {} + }, + { + "rule_id": "check_missing_refs_count", + "validator": "MISSING_REFS_COUNT", + "params": { + "threshold": 0 + } + }, + { + "rule_id": "check_lint_error_count", + "validator": "LINT_ERROR_COUNT", + "params": { + "threshold": 0 + } + } + ] +} From 4ebc67ec83d624dfa7711446bfacba2c8f6c7fcd Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 10:09:39 +0000 Subject: [PATCH 02/21] Fix county sharding logic: initialize headers for all shards and prevent infinite loop on last shard --- .../parse_air_quality.py | 33 +++++++++++-------- 1 file changed, 19 insertions(+), 14 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 9f352128be..4b5e355361 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -129,7 +129,7 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): chunk_size = 500_000 shard_idx = 0 shard_written = 0 - first_write = [True] * num_shards + first_chunk = True for chunk in pd.read_csv(input_file_path, chunksize=chunk_size): @@ -144,31 +144,36 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): chunk["dcid"] = "geoId/" + chunk[ "statefips"] + chunk["countyfips"] + if first_chunk: + for p in shard_paths: + pd.DataFrame(columns=chunk.columns).to_csv( + p, index=False) + first_chunk = False + start_idx = 0 while start_idx < len(chunk): - remaining_in_shard = shard_sizes[ - shard_idx] - shard_written - end_idx = min( - start_idx + remaining_in_shard, - len(chunk)) + if shard_idx < num_shards - 1: + remaining_in_shard = shard_sizes[ + shard_idx] - shard_written + end_idx = min( + start_idx + remaining_in_shard, + len(chunk)) + else: + end_idx = len(chunk) + sub_chunk = chunk.iloc[ start_idx:end_idx] - mode = 'w' if first_write[ - shard_idx] else 'a' - header = first_write[shard_idx] sub_chunk.to_csv( shard_paths[shard_idx], - mode=mode, - header=header, + mode='a', + header=False, float_format='%.6f', index=False) - first_write[shard_idx] = False shard_written += len(sub_chunk) start_idx = end_idx - if shard_written >= shard_sizes[ - shard_idx] and shard_idx < num_shards - 1: + if shard_idx < num_shards - 1 and shard_written >= shard_sizes[shard_idx]: shard_idx += 1 shard_written = 0 From 5a7f855ad16a02ea0e016eb6d6577a8ecd278e38 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 10:19:53 +0000 Subject: [PATCH 03/21] Apply yapf Google code formatting to scripts/us_cdc/environmental_health_toxicology --- .../download_files.py | 3 +- .../parse_air_quality.py | 31 ++++++++++--------- .../parse_air_quality_test.py | 6 ++-- 3 files changed, 19 insertions(+), 21 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index b63b89f8ad..f1183a1138 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -40,8 +40,7 @@ def download_with_retry(url, input_file_name): with requests.get(url, stream=True) as response: response.raise_for_status() with open(filename, 'wb') as f: - for chunk in response.iter_content( - chunk_size=16 * 1024 * 1024): + for chunk in response.iter_content(chunk_size=16 * 1024 * 1024): if chunk: f.write(chunk) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 4b5e355361..b28a686987 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -106,7 +106,8 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): try: if "County" in input_file_name and "PM" in input_file_name: num_shards = 4 - base_name, ext = os.path.splitext(output_file_name) + base_name, ext = os.path.splitext( + output_file_name) shard_paths = [] for idx in range(num_shards): shard_file_name = f"{base_name}_{idx}{ext}" @@ -146,8 +147,9 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): if first_chunk: for p in shard_paths: - pd.DataFrame(columns=chunk.columns).to_csv( - p, index=False) + pd.DataFrame( + columns=chunk.columns).to_csv( + p, index=False) first_chunk = False start_idx = 0 @@ -164,16 +166,16 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): sub_chunk = chunk.iloc[ start_idx:end_idx] - sub_chunk.to_csv( - shard_paths[shard_idx], - mode='a', - header=False, - float_format='%.6f', - index=False) + sub_chunk.to_csv(shard_paths[shard_idx], + mode='a', + header=False, + float_format='%.6f', + index=False) shard_written += len(sub_chunk) start_idx = end_idx - if shard_idx < num_shards - 1 and shard_written >= shard_sizes[shard_idx]: + if shard_idx < num_shards - 1 and shard_written >= shard_sizes[ + shard_idx]: shard_idx += 1 shard_written = 0 @@ -221,11 +223,10 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): ], var_name='StatisticalVariable', value_name='Value') - data.rename( - columns={ - census_tract + '_stdd': 'Error' - }, - inplace=True) + data.rename(columns={ + census_tract + '_stdd': 'Error' + }, + inplace=True) max_length = data['ctfips'].astype( str).str.len().max() data['ctfips'] = data['ctfips'].astype( diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py index 953daba52f..2147622914 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py @@ -68,10 +68,8 @@ f"PM25county_{i}.csv") for i in range(4) ], "files": [{ - "input_file_name": - "PM2.5County_input_0.csv", - "output_file_name": - "PM25county.csv" + "input_file_name": "PM2.5County_input_0.csv", + "output_file_name": "PM25county.csv" }] }, { "import_name": From 30e7938fb8d7d808525199f87e3727b927cec99d Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 10:29:13 +0000 Subject: [PATCH 04/21] Add automatic test artifact cleanup in parse_air_quality_test tearDown --- .../parse_air_quality_test.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py index 2147622914..9d62d0de7e 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py @@ -11,8 +11,9 @@ # 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. -import unittest import os +import shutil +import unittest from .parse_air_quality import clean_air_quality_data _MODULE_DIR = os.path.dirname(__file__) @@ -126,6 +127,10 @@ def test_clean_air_quality_data(self): self.assertEqual(expected_csv_data, actual_csv_data) + def tearDown(self): + for data in TEST_DATA: + shutil.rmtree(data["output_dir"], ignore_errors=True) + if __name__ == '__main__': unittest.main() From 4e0cd49885909ea004942673a4ce0cf242fe1909 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 10:48:02 +0000 Subject: [PATCH 05/21] Remove redundant config_override block from manifest.json --- .../manifest.json | 20 ------------------- 1 file changed, 20 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index 24f4bfa5f9..f15b66f913 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -20,11 +20,6 @@ "disk": 2000 }, "validation_config_file": "validation_config.json", - "config_override": { - "invoke_import_tool": true, - "invoke_differ_tool": true, - "user_script_timeout": 28800 - }, "import_inputs": [ { "template_mcf": "PM25CensusTractPollution_part1.tmcf", @@ -65,11 +60,6 @@ "disk": 2000 }, "validation_config_file": "validation_config.json", - "config_override": { - "invoke_import_tool": true, - "invoke_differ_tool": true, - "user_script_timeout": 28800 - }, "import_inputs": [ { "template_mcf": "OzoneCensusTractPollution_part1.tmcf", @@ -110,11 +100,6 @@ "disk": 500 }, "validation_config_file": "validation_config.json", - "config_override": { - "invoke_import_tool": true, - "invoke_differ_tool": true, - "user_script_timeout": 28800 - }, "import_inputs": [ { "template_mcf": "PM25CountyPollution_part1.tmcf", @@ -155,11 +140,6 @@ "disk": 500 }, "validation_config_file": "validation_config.json", - "config_override": { - "invoke_import_tool": true, - "invoke_differ_tool": true, - "user_script_timeout": 28800 - }, "import_inputs": [ { "template_mcf": "OzoneCountyPollution.tmcf", From 1363a5f2fd3dfa95f764b43004c6deb7026598f4 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 11:44:11 +0000 Subject: [PATCH 06/21] Update README.md with sharded tMCF files and architecture documentation --- .../environmental_health_toxicology/README.md | 20 +++++++------------ 1 file changed, 7 insertions(+), 13 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/README.md b/scripts/us_cdc/environmental_health_toxicology/README.md index 0092777035..f82a1190da 100644 --- a/scripts/us_cdc/environmental_health_toxicology/README.md +++ b/scripts/us_cdc/environmental_health_toxicology/README.md @@ -69,19 +69,13 @@ These data were collected as part of the [CDC National Environment Public Health [`small_Palmer_expected.csv`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/test_data/small_Palmer_expected.csv) #### tMCFs -[`OzoneCensusTractPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution.tmcf) - -[`OzoneCountyPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCountyPollution.tmcf) - -[`PalmerDroughtSeverityIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PalmerDroughtSeverityIndex.tmcf) - -[`PM25CensusTractPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution.tmcf) - -[`PM25CountyPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf) - -[`StandardizedPrecipitationEvapotranspirationIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/StandardizedPrecipitationEvapotranspirationIndex.tmcf) - -[`StandardizedPrecipitationIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/StandardizedPrecipitationIndex.tmcf) +* [`OzoneCountyPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCountyPollution.tmcf) +* **Sharded PM2.5 County tMCFs:** [`PM25CountyPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf) +* **Sharded PM2.5 Census Tract tMCFs:** [`PM25CensusTractPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part4.tmcf) +* **Sharded Ozone Census Tract tMCFs:** [`OzoneCensusTractPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part4.tmcf) +* [`PalmerDroughtSeverityIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PalmerDroughtSeverityIndex.tmcf) +* [`StandardizedPrecipitationEvapotranspirationIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/StandardizedPrecipitationEvapotranspirationIndex.tmcf) +* [`StandardizedPrecipitationIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/StandardizedPrecipitationIndex.tmcf) ### Import Procedure From c43b1d19777bcc27fdda86cf0e59d2c4289aa841 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Wed, 2 Sep 2026 05:32:39 +0000 Subject: [PATCH 07/21] Disable in-memory differ for census tract imports to prevent OOM on single VM --- .../us_cdc/environmental_health_toxicology/manifest.json | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index f15b66f913..3df2b223ef 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -19,6 +19,9 @@ "memory": 512, "disk": 2000 }, + "config_override": { + "invoke_differ_tool": false + }, "validation_config_file": "validation_config.json", "import_inputs": [ { @@ -59,6 +62,9 @@ "memory": 512, "disk": 2000 }, + "config_override": { + "invoke_differ_tool": false + }, "validation_config_file": "validation_config.json", "import_inputs": [ { From 46e4e57a5f86dc20aa6bafa512f5eb6625997eed Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Wed, 2 Sep 2026 13:24:14 +0000 Subject: [PATCH 08/21] Consolidate PM25County TMCFs to single template and streamline validation_config.json - Consolidate all 4 sharded CDC_PM25County import_inputs in manifest.json to reference single PM25CountyPollution.tmcf. - Remove redundant duplicate PM25CountyPollution_part[1-4].tmcf files. - Streamline validation_config.json by removing redundant rules inherited from base system config and adding description justifying the 0.05 deletion threshold. - Update README.md to reflect single PM25CountyPollution.tmcf. --- .../PM25CountyPollution_part1.tmcf | 35 ------------------- .../PM25CountyPollution_part2.tmcf | 35 ------------------- .../PM25CountyPollution_part3.tmcf | 35 ------------------- .../PM25CountyPollution_part4.tmcf | 35 ------------------- .../environmental_health_toxicology/README.md | 2 +- .../manifest.json | 8 ++--- .../validation_config.json | 22 +----------- 7 files changed, 6 insertions(+), 166 deletions(-) delete mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf delete mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf delete mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf delete mode 100644 scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf deleted file mode 100644 index 15750606ac..0000000000 --- a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf +++ /dev/null @@ -1,35 +0,0 @@ -Node: E:PM25CountyPollution->E1 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_mean_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E2 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_med_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E3 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_max_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E4 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_pop_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf deleted file mode 100644 index 15750606ac..0000000000 --- a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf +++ /dev/null @@ -1,35 +0,0 @@ -Node: E:PM25CountyPollution->E1 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_mean_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E2 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_med_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E3 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_max_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E4 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_pop_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf deleted file mode 100644 index 15750606ac..0000000000 --- a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf +++ /dev/null @@ -1,35 +0,0 @@ -Node: E:PM25CountyPollution->E1 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_mean_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E2 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_med_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E3 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_max_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E4 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_pop_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf b/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf deleted file mode 100644 index 15750606ac..0000000000 --- a/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf +++ /dev/null @@ -1,35 +0,0 @@ -Node: E:PM25CountyPollution->E1 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_mean_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Mean_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E2 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_med_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Median_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E3 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_max_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:Max_Concentration_AirPollutant_PM2.5 - -Node: E:PM25CountyPollution->E4 -observationAbout: C:PM25CountyPollution->dcid -typeOf: dcs:StatVarObservation -observationDate: C:PM25CountyPollution->date -value: C:PM25CountyPollution->pm25_pop_pred -observationPeriod: "P24H" -unit: MicrogramsPerCubicMeter -variableMeasured: dcs:PopulationWeighted_Concentration_AirPollutant_PM2.5 diff --git a/scripts/us_cdc/environmental_health_toxicology/README.md b/scripts/us_cdc/environmental_health_toxicology/README.md index f82a1190da..59c8280498 100644 --- a/scripts/us_cdc/environmental_health_toxicology/README.md +++ b/scripts/us_cdc/environmental_health_toxicology/README.md @@ -70,7 +70,7 @@ These data were collected as part of the [CDC National Environment Public Health #### tMCFs * [`OzoneCountyPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCountyPollution.tmcf) -* **Sharded PM2.5 County tMCFs:** [`PM25CountyPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution_part4.tmcf) +* [`PM25CountyPollution.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CountyPollution.tmcf) * **Sharded PM2.5 Census Tract tMCFs:** [`PM25CensusTractPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PM25CensusTractPollution_part4.tmcf) * **Sharded Ozone Census Tract tMCFs:** [`OzoneCensusTractPollution_part1.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part1.tmcf), [`part2`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part2.tmcf), [`part3`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part3.tmcf), [`part4`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/OzoneCensusTractPollution_part4.tmcf) * [`PalmerDroughtSeverityIndex.tmcf`](https://github.com/datacommonsorg/data/blob/master/scripts/us_cdc/environmental_health_toxicology/PalmerDroughtSeverityIndex.tmcf) diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index 3df2b223ef..afafdcdea0 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -108,19 +108,19 @@ "validation_config_file": "validation_config.json", "import_inputs": [ { - "template_mcf": "PM25CountyPollution_part1.tmcf", + "template_mcf": "PM25CountyPollution.tmcf", "cleaned_csv": "output/PM25county_0.csv" }, { - "template_mcf": "PM25CountyPollution_part2.tmcf", + "template_mcf": "PM25CountyPollution.tmcf", "cleaned_csv": "output/PM25county_1.csv" }, { - "template_mcf": "PM25CountyPollution_part3.tmcf", + "template_mcf": "PM25CountyPollution.tmcf", "cleaned_csv": "output/PM25county_2.csv" }, { - "template_mcf": "PM25CountyPollution_part4.tmcf", + "template_mcf": "PM25CountyPollution.tmcf", "cleaned_csv": "output/PM25county_3.csv" } ], diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json index cd5e8d8315..027c3724ca 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -3,31 +3,11 @@ "rules": [ { "rule_id": "check_deleted_records_percent", - "description": "Checks that the percentage of deleted points is within the threshold.", + "description": "Checks that the percentage of deleted points is within the threshold (allows up to 5% to accommodate 0.02% historical CDC county boundary recalibrations in OzoneCounty while maintaining strict guardrails against data loss).", "validator": "DELETED_RECORDS_PERCENT", "params": { "threshold": 0.05 } - }, - { - "rule_id": "check_empty_import", - "description": "Checks if the import is empty (no observations and no schema).", - "validator": "EMPTY_IMPORT_CHECK", - "params": {} - }, - { - "rule_id": "check_missing_refs_count", - "validator": "MISSING_REFS_COUNT", - "params": { - "threshold": 0 - } - }, - { - "rule_id": "check_lint_error_count", - "validator": "LINT_ERROR_COUNT", - "params": { - "threshold": 0 - } } ] } From 4740ef3a2b07e401b37665c21948308777709218 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Wed, 2 Sep 2026 18:36:56 +0000 Subject: [PATCH 09/21] Revert description in validation_config.json to standard form (details documented in PR description) --- .../environmental_health_toxicology/validation_config.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json index 027c3724ca..071a1c5179 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -3,7 +3,7 @@ "rules": [ { "rule_id": "check_deleted_records_percent", - "description": "Checks that the percentage of deleted points is within the threshold (allows up to 5% to accommodate 0.02% historical CDC county boundary recalibrations in OzoneCounty while maintaining strict guardrails against data loss).", + "description": "Checks that the percentage of deleted points is within the threshold.", "validator": "DELETED_RECORDS_PERCENT", "params": { "threshold": 0.05 From 681ad34c2fc273d247fa940d058561a094aa9904 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Thu, 3 Sep 2026 05:57:58 +0000 Subject: [PATCH 10/21] fix(cdc): remove validation_config_file from census tract imports where differ is decoupled --- scripts/us_cdc/environmental_health_toxicology/manifest.json | 2 -- 1 file changed, 2 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index afafdcdea0..e4e36a9fb3 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -22,7 +22,6 @@ "config_override": { "invoke_differ_tool": false }, - "validation_config_file": "validation_config.json", "import_inputs": [ { "template_mcf": "PM25CensusTractPollution_part1.tmcf", @@ -65,7 +64,6 @@ "config_override": { "invoke_differ_tool": false }, - "validation_config_file": "validation_config.json", "import_inputs": [ { "template_mcf": "OzoneCensusTractPollution_part1.tmcf", From ab963d1bf7e276b898e235013a99ca598b191dab Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Thu, 3 Sep 2026 09:14:44 +0000 Subject: [PATCH 11/21] fix(cdc): address code review findings for timeout, date parsing, and validation consistency --- .../us_cdc/environmental_health_toxicology/download_files.py | 5 +++-- .../environmental_health_toxicology/parse_air_quality.py | 4 ++-- .../environmental_health_toxicology/validation_config.json | 5 +++++ 3 files changed, 10 insertions(+), 4 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index f1183a1138..aad90a356a 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -37,7 +37,7 @@ def download_files(importname, configs): def download_with_retry(url, input_file_name): logging.info(f"Downloading file from URL: {url}") filename = os.path.join(_INPUT_FILE_PATH, input_file_name) - with requests.get(url, stream=True) as response: + with requests.get(url, stream=True, timeout=(30, 300)) as response: response.raise_for_status() with open(filename, 'wb') as f: for chunk in response.iter_content(chunk_size=16 * 1024 * 1024): @@ -55,7 +55,8 @@ def download_with_retry(url, input_file_name): logging.info(f"Input File Name {input_file_name}") get_record_count = requests.get( - url_new.replace('.csv', record_count_query)) + url_new.replace('.csv', record_count_query), + timeout=60) if get_record_count.status_code == 200: record_count = json.loads( get_record_count.text diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index b28a686987..20cf0a97cf 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -137,7 +137,7 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): chunk["date"] = pd.to_datetime( chunk["date"], format="%d%b%Y", - errors="coerce").dt.strftime("%Y-%m-%d") + errors="raise").dt.strftime("%Y-%m-%d") chunk["statefips"] = chunk[ "statefips"].astype(str).str.zfill(2) chunk["countyfips"] = chunk[ @@ -188,7 +188,7 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): data["date"] = pd.to_datetime( data["date"], format="%d%b%Y", - errors="coerce").dt.strftime("%Y-%m-%d") + errors="raise").dt.strftime("%Y-%m-%d") if "PM2.5" in input_file_name: census_tract = "ds_pm" diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json index 071a1c5179..bac5d68f65 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -1,6 +1,11 @@ { "schema_version": "1.0", "rules": [ + { + "rule_id": "check_max_date_consistent", + "description": "Checks if the MaxDate is the same for all StatVars.", + "validator": "MAX_DATE_CONSISTENT" + }, { "rule_id": "check_deleted_records_percent", "description": "Checks that the percentage of deleted points is within the threshold.", From 53d7457c81e166b00cba117015c97377e5421d41 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Thu, 3 Sep 2026 09:37:40 +0000 Subject: [PATCH 12/21] fix(cdc): remove check_max_date_consistent rule from validation_config.json --- .../environmental_health_toxicology/validation_config.json | 5 ----- 1 file changed, 5 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json index bac5d68f65..071a1c5179 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -1,11 +1,6 @@ { "schema_version": "1.0", "rules": [ - { - "rule_id": "check_max_date_consistent", - "description": "Checks if the MaxDate is the same for all StatVars.", - "validator": "MAX_DATE_CONSISTENT" - }, { "rule_id": "check_deleted_records_percent", "description": "Checks that the percentage of deleted points is within the threshold.", From dd29405fb8d74abfea429b1a137ae74dc516d11a Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Thu, 3 Sep 2026 09:56:20 +0000 Subject: [PATCH 13/21] style(cdc): format download_files.py with yapf Google style --- .../environmental_health_toxicology/download_files.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index aad90a356a..2d008fb3d0 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -54,9 +54,9 @@ def download_with_retry(url, input_file_name): input_file_name = file_info["input_file_name"] logging.info(f"Input File Name {input_file_name}") - get_record_count = requests.get( - url_new.replace('.csv', record_count_query), - timeout=60) + get_record_count = requests.get(url_new.replace( + '.csv', record_count_query), + timeout=60) if get_record_count.status_code == 200: record_count = json.loads( get_record_count.text From 06808e9d6765382f3035f970b5eae315aae7c367 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 06:08:24 +0000 Subject: [PATCH 14/21] fix(cdc): resolve adversarial review findings for download error handling and ozone streaming - Add raise_for_status() to download_files.py to eliminate silent exit on HTTP error - Prevent UnboundLocalError and re-raise exception in download_files.py error handling - Validate import_name and raise ValueError if unrecognized - Refactor CDC_OzoneCounty parsing to stream in 500k row chunks to eliminate ~45GB RAM spikes - Enforce fixed 11-digit zfill for Census Tract ctfips/dcid and update test fixtures - Add robust module import fallbacks in unittests --- .../download_files.py | 37 +++++----- .../parse_air_quality.py | 59 +++++++++++----- .../parse_air_quality_test.py | 69 ++++++++++--------- .../parse_precipitation_index_test.py | 10 ++- ...sus_Tract_Level_Ozone_Concentrations_0.csv | 10 +-- 5 files changed, 110 insertions(+), 75 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index 2d008fb3d0..4fb88ee60c 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -40,13 +40,17 @@ def download_with_retry(url, input_file_name): with requests.get(url, stream=True, timeout=(30, 300)) as response: response.raise_for_status() with open(filename, 'wb') as f: - for chunk in response.iter_content(chunk_size=16 * 1024 * 1024): + for chunk in response.iter_content(chunk_size=16 * 1024 * + 1024): if chunk: f.write(chunk) + url_new = None + import_found = False try: for config in configs: if config["import_name"] == importname: + import_found = True files = config["files"] for file_info in files: url_new = file_info["url"] @@ -57,24 +61,23 @@ def download_with_retry(url, input_file_name): get_record_count = requests.get(url_new.replace( '.csv', record_count_query), timeout=60) - if get_record_count.status_code == 200: - record_count = json.loads( - get_record_count.text - )[0]['COLUMN_ALIAS_GUARD__count'] - logging.info( - f"Numbers of records found for the URL {url_new} is {record_count}" - ) - url_new = f"{url_new}?$limit={record_count}&$offset=0" - download_with_retry(url_new, input_file_name) - logging.info( - "Successfully downloaded the source data...!!!!") - else: - logging.error( - f"Failed to download files, Status code: {get_record_count.status_code}" - ) + get_record_count.raise_for_status() + record_count = json.loads( + get_record_count.text)[0]['COLUMN_ALIAS_GUARD__count'] + logging.info( + f"Numbers of records found for the URL {url_new} is {record_count}" + ) + url_new = f"{url_new}?$limit={record_count}&$offset=0" + download_with_retry(url_new, input_file_name) + logging.info( + "Successfully downloaded the source data...!!!!") + if not import_found: + raise ValueError( + f"Import name '{importname}' not found in configuration") except Exception as e: - logging.fatal(f"Error downloading URL {url_new} - {e}") + logging.fatal(f"Error downloading URL {url_new or 'unknown'} - {e}") + raise def main(_): diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 20cf0a97cf..771638c3a5 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -70,8 +70,8 @@ # this method is applicable only for "census tract PM25" -def add_prefix_zero(value, length): - return value.zfill(length) +def add_prefix_zero(value, length=11): + return str(value).zfill(length) def clean_air_quality_data(configs, importname, inputpath, outputpath): @@ -123,7 +123,8 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): base_size = total_rows // num_shards rem_size = total_rows % num_shards shard_sizes = [ - base_size + 1 if i < rem_size else base_size + base_size + + 1 if i < rem_size else base_size for i in range(num_shards) ] @@ -166,11 +167,12 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): sub_chunk = chunk.iloc[ start_idx:end_idx] - sub_chunk.to_csv(shard_paths[shard_idx], - mode='a', - header=False, - float_format='%.6f', - index=False) + sub_chunk.to_csv( + shard_paths[shard_idx], + mode='a', + header=False, + float_format='%.6f', + index=False) shard_written += len(sub_chunk) start_idx = end_idx @@ -183,6 +185,35 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): logging.info( f"Finished cleaning file {os.path.basename(p)}!" ) + elif "County" in input_file_name and "Ozone" in input_file_name: + chunk_size = 500_000 + first_chunk = True + for chunk in pd.read_csv(input_file_path, + chunksize=chunk_size): + chunk["date"] = pd.to_datetime( + chunk["date"], + format="%d%b%Y", + errors="raise").dt.strftime("%Y-%m-%d") + chunk["statefips"] = chunk[ + "statefips"].astype(str).str.zfill(2) + chunk["countyfips"] = chunk[ + "countyfips"].astype(str).str.zfill(3) + chunk["dcid"] = "geoId/" + chunk[ + "statefips"] + chunk["countyfips"] + if first_chunk: + chunk.to_csv(output_file_path, + float_format='%.6f', + index=False) + first_chunk = False + else: + chunk.to_csv(output_file_path, + mode='a', + header=False, + float_format='%.6f', + index=False) + logging.info( + f"Finished cleaning file {output_file_name}!" + ) else: data = pd.read_csv(input_file_path) data["date"] = pd.to_datetime( @@ -227,22 +258,12 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): census_tract + '_stdd': 'Error' }, inplace=True) - max_length = data['ctfips'].astype( - str).str.len().max() data['ctfips'] = data['ctfips'].astype( - str).apply(lambda x: add_prefix_zero( - x, max_length)) + str).str.zfill(11) data["dcid"] = "geoId/" + data[ "ctfips"].astype(str) data['StatisticalVariable'] = data[ 'StatisticalVariable'].map(STATVARS) - elif "County" in input_file_name and "Ozone" in input_file_name: - data["statefips"] = data[ - "statefips"].astype(str).str.zfill(2) - data["countyfips"] = data[ - "countyfips"].astype(str).str.zfill(3) - data["dcid"] = "geoId/" + data[ - "statefips"] + data["countyfips"] data.to_csv(output_file_path, float_format='%.6f', index=False) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py index 9d62d0de7e..88f6ca3ee0 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py @@ -8,15 +8,19 @@ # # 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. import os import shutil +import sys import unittest -from .parse_air_quality import clean_air_quality_data _MODULE_DIR = os.path.dirname(__file__) +sys.path.insert(0, _MODULE_DIR) + +try: + from .parse_air_quality import clean_air_quality_data +except ImportError: + from parse_air_quality import clean_air_quality_data + TEST_DATA_DIR = os.path.join(_MODULE_DIR, 'test_data') INPUT_DIR = 'input_files' OUTPUT_DIR = 'actual_output_files' @@ -25,45 +29,45 @@ # test data for each import type TEST_DATA = [{ "import_name": - "CDC_PM25CensusTract", + "CDC_PM25CensusTract", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_FILES, - "PM2.5CensusTract_0.csv"), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_FILES, + "PM2.5CensusTract_0.csv"), "files": [{ "input_file_name": - "PM2.5CensusTractPollution_input_0.csv", + "PM2.5CensusTractPollution_input_0.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR, - "PM2.5CensusTract_0.csv") + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR, + "PM2.5CensusTract_0.csv") }] }, { "import_name": - "CDC_OzoneCensusTract", + "CDC_OzoneCensusTract", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_FILES, - "Census_Tract_Level_Ozone_Concentrations_0.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_FILES, + "Census_Tract_Level_Ozone_Concentrations_0.csv"), "files": [{ "input_file_name": - "Census_Tract_Level_Ozone_Concentrations_input_0.csv", + "Census_Tract_Level_Ozone_Concentrations_input_0.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR, - "Census_Tract_Level_Ozone_Concentrations_0.csv") + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR, + "Census_Tract_Level_Ozone_Concentrations_0.csv") }] }, { "import_name": - "CDC_PM25County", + "CDC_PM25County", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25County", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25County", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR), "expected_files": [ os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_FILES, f"PM25county_{i}.csv") for i in range(4) @@ -74,20 +78,20 @@ }] }, { "import_name": - "CDC_OzoneCounty", + "CDC_OzoneCounty", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_FILES, - "OzoneCounty.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_FILES, + "OzoneCounty.csv"), "files": [{ "input_file_name": - "OzoneCounty_input.csv", + "OzoneCounty_input.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR, - "OzoneCounty.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR, + "OzoneCounty.csv"), }] }] @@ -110,7 +114,8 @@ def test_clean_air_quality_data(self): for expected_file in data["expected_files"]: file_name = os.path.basename(expected_file) actual_file = os.path.join(output_dir, file_name) - with open(actual_file, encoding="utf-8") as actual_csv_file: + with open(actual_file, + encoding="utf-8") as actual_csv_file: actual_csv_data = actual_csv_file.read().strip() with open(expected_file, encoding="utf-8") as expected_csv_file: diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_precipitation_index_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_precipitation_index_test.py index b69ac7c8b1..8f5e1f4183 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_precipitation_index_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_precipitation_index_test.py @@ -20,11 +20,17 @@ python3 parse_precipitation_index_test.py input_file output_file ''' -import unittest import os -from .parse_precipitation_index import clean_precipitation_data +import sys +import unittest module_dir_ = os.path.dirname(__file__) +sys.path.insert(0, module_dir_) + +try: + from .parse_precipitation_index import clean_precipitation_data +except ImportError: + from parse_precipitation_index import clean_precipitation_data class TestParsePrecipitationData(unittest.TestCase): diff --git a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_OzoneCensusTract/expected_output_files/Census_Tract_Level_Ozone_Concentrations_0.csv b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_OzoneCensusTract/expected_output_files/Census_Tract_Level_Ozone_Concentrations_0.csv index 9debb88cb2..99ac99313b 100644 --- a/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_OzoneCensusTract/expected_output_files/Census_Tract_Level_Ozone_Concentrations_0.csv +++ b/scripts/us_cdc/environmental_health_toxicology/test_data/CDC_OzoneCensusTract/expected_output_files/Census_Tract_Level_Ozone_Concentrations_0.csv @@ -1,6 +1,6 @@ year,date,statefips,countyfips,ctfips,latitude,longitude,Error,StatisticalVariable,Value,dcid -2001,2001-01-01,1,1,1001020100,32.477180,-86.490010,5.931458,Mean_Concentration_AirPollutant_Ozone,31.659713,geoId/1001020100 -2001,2001-01-01,1,1,1001020200,32.474250,-86.473390,5.703084,Mean_Concentration_AirPollutant_Ozone,31.939058,geoId/1001020200 -2001,2001-01-01,1,1,1001020300,32.475440,-86.460200,5.779281,Mean_Concentration_AirPollutant_Ozone,31.859155,geoId/1001020300 -2001,2001-01-01,1,1,1001020400,32.472040,-86.443700,5.784730,Mean_Concentration_AirPollutant_Ozone,31.681802,geoId/1001020400 -2001,2001-01-01,1,1,1001020500,32.458920,-86.422710,5.797620,Mean_Concentration_AirPollutant_Ozone,31.752820,geoId/1001020500 +2001,2001-01-01,1,1,01001020100,32.477180,-86.490010,5.931458,Mean_Concentration_AirPollutant_Ozone,31.659713,geoId/01001020100 +2001,2001-01-01,1,1,01001020200,32.474250,-86.473390,5.703084,Mean_Concentration_AirPollutant_Ozone,31.939058,geoId/01001020200 +2001,2001-01-01,1,1,01001020300,32.475440,-86.460200,5.779281,Mean_Concentration_AirPollutant_Ozone,31.859155,geoId/01001020300 +2001,2001-01-01,1,1,01001020400,32.472040,-86.443700,5.784730,Mean_Concentration_AirPollutant_Ozone,31.681802,geoId/01001020400 +2001,2001-01-01,1,1,01001020500,32.458920,-86.422710,5.797620,Mean_Concentration_AirPollutant_Ozone,31.752820,geoId/01001020500 From 87ac66042f825c78380bf9bfedf638765f4d1c04 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 06:22:32 +0000 Subject: [PATCH 15/21] fix(cdc): restore check_max_date_consistent rule in validation_config.json --- .../environmental_health_toxicology/validation_config.json | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json index 071a1c5179..bac5d68f65 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config.json @@ -1,6 +1,11 @@ { "schema_version": "1.0", "rules": [ + { + "rule_id": "check_max_date_consistent", + "description": "Checks if the MaxDate is the same for all StatVars.", + "validator": "MAX_DATE_CONSISTENT" + }, { "rule_id": "check_deleted_records_percent", "description": "Checks that the percentage of deleted points is within the threshold.", From 288973ad9c589472451918bf3b293079c8cdee53 Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 06:36:34 +0000 Subject: [PATCH 16/21] style(cdc): format with yapf --style=google to satisfy CI lint check --- .../download_files.py | 3 +- .../parse_air_quality.py | 14 ++--- .../parse_air_quality_test.py | 57 +++++++++---------- 3 files changed, 35 insertions(+), 39 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index 4fb88ee60c..1e673eb564 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -40,8 +40,7 @@ def download_with_retry(url, input_file_name): with requests.get(url, stream=True, timeout=(30, 300)) as response: response.raise_for_status() with open(filename, 'wb') as f: - for chunk in response.iter_content(chunk_size=16 * 1024 * - 1024): + for chunk in response.iter_content(chunk_size=16 * 1024 * 1024): if chunk: f.write(chunk) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 771638c3a5..5adb0adf4b 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -123,8 +123,7 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): base_size = total_rows // num_shards rem_size = total_rows % num_shards shard_sizes = [ - base_size + - 1 if i < rem_size else base_size + base_size + 1 if i < rem_size else base_size for i in range(num_shards) ] @@ -167,12 +166,11 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): sub_chunk = chunk.iloc[ start_idx:end_idx] - sub_chunk.to_csv( - shard_paths[shard_idx], - mode='a', - header=False, - float_format='%.6f', - index=False) + sub_chunk.to_csv(shard_paths[shard_idx], + mode='a', + header=False, + float_format='%.6f', + index=False) shard_written += len(sub_chunk) start_idx = end_idx diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py index 88f6ca3ee0..c0b17923f6 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality_test.py @@ -29,45 +29,45 @@ # test data for each import type TEST_DATA = [{ "import_name": - "CDC_PM25CensusTract", + "CDC_PM25CensusTract", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_FILES, - "PM2.5CensusTract_0.csv"), + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_FILES, + "PM2.5CensusTract_0.csv"), "files": [{ "input_file_name": - "PM2.5CensusTractPollution_input_0.csv", + "PM2.5CensusTractPollution_input_0.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR, - "PM2.5CensusTract_0.csv") + os.path.join(TEST_DATA_DIR, "CDC_PM25CensusTract", OUTPUT_DIR, + "PM2.5CensusTract_0.csv") }] }, { "import_name": - "CDC_OzoneCensusTract", + "CDC_OzoneCensusTract", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_FILES, - "Census_Tract_Level_Ozone_Concentrations_0.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_FILES, + "Census_Tract_Level_Ozone_Concentrations_0.csv"), "files": [{ "input_file_name": - "Census_Tract_Level_Ozone_Concentrations_input_0.csv", + "Census_Tract_Level_Ozone_Concentrations_input_0.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR, - "Census_Tract_Level_Ozone_Concentrations_0.csv") + os.path.join(TEST_DATA_DIR, "CDC_OzoneCensusTract", OUTPUT_DIR, + "Census_Tract_Level_Ozone_Concentrations_0.csv") }] }, { "import_name": - "CDC_PM25County", + "CDC_PM25County", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25County", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25County", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_DIR), "expected_files": [ os.path.join(TEST_DATA_DIR, "CDC_PM25County", OUTPUT_FILES, f"PM25county_{i}.csv") for i in range(4) @@ -78,20 +78,20 @@ }] }, { "import_name": - "CDC_OzoneCounty", + "CDC_OzoneCounty", "input_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", INPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", INPUT_DIR), "output_dir": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR), "expected_file": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_FILES, - "OzoneCounty.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_FILES, + "OzoneCounty.csv"), "files": [{ "input_file_name": - "OzoneCounty_input.csv", + "OzoneCounty_input.csv", "output_file_name": - os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR, - "OzoneCounty.csv"), + os.path.join(TEST_DATA_DIR, "CDC_OzoneCounty", OUTPUT_DIR, + "OzoneCounty.csv"), }] }] @@ -114,8 +114,7 @@ def test_clean_air_quality_data(self): for expected_file in data["expected_files"]: file_name = os.path.basename(expected_file) actual_file = os.path.join(output_dir, file_name) - with open(actual_file, - encoding="utf-8") as actual_csv_file: + with open(actual_file, encoding="utf-8") as actual_csv_file: actual_csv_data = actual_csv_file.read().strip() with open(expected_file, encoding="utf-8") as expected_csv_file: From 44b3c598c9d48c5d591064e42c569850802272ef Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 06:47:46 +0000 Subject: [PATCH 17/21] feat(cdc): add separate validation configs for Census Tract and County imports - Add validation_config_census_tract.json disabling check_deleted_records_percent (since differ is decoupled) and configuring check_lint_error_count tolerance for remote network drops - Add validation_config_county.json with check_max_date_consistent and 5% check_deleted_records_percent threshold - Update manifest.json to explicitly bind Census Tract and County imports to their dedicated validation configurations --- .../manifest.json | 6 ++++-- .../validation_config_census_tract.json | 17 +++++++++++++++++ .../validation_config_county.json | 18 ++++++++++++++++++ 3 files changed, 39 insertions(+), 2 deletions(-) create mode 100644 scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json create mode 100644 scripts/us_cdc/environmental_health_toxicology/validation_config_county.json diff --git a/scripts/us_cdc/environmental_health_toxicology/manifest.json b/scripts/us_cdc/environmental_health_toxicology/manifest.json index e4e36a9fb3..265c11c316 100644 --- a/scripts/us_cdc/environmental_health_toxicology/manifest.json +++ b/scripts/us_cdc/environmental_health_toxicology/manifest.json @@ -19,6 +19,7 @@ "memory": 512, "disk": 2000 }, + "validation_config_file": "validation_config_census_tract.json", "config_override": { "invoke_differ_tool": false }, @@ -61,6 +62,7 @@ "memory": 512, "disk": 2000 }, + "validation_config_file": "validation_config_census_tract.json", "config_override": { "invoke_differ_tool": false }, @@ -103,7 +105,7 @@ "memory": 512, "disk": 500 }, - "validation_config_file": "validation_config.json", + "validation_config_file": "validation_config_county.json", "import_inputs": [ { "template_mcf": "PM25CountyPollution.tmcf", @@ -143,7 +145,7 @@ "memory": 512, "disk": 500 }, - "validation_config_file": "validation_config.json", + "validation_config_file": "validation_config_county.json", "import_inputs": [ { "template_mcf": "OzoneCountyPollution.tmcf", diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json new file mode 100644 index 0000000000..fe9e41637e --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json @@ -0,0 +1,17 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_deleted_records_percent", + "description": "Disable deleted records check because differ is decoupled for high-scale Census Tracts.", + "enabled": false + }, + { + "rule_id": "check_lint_error_count", + "description": "Tolerate transient remote network RPC drops during multi-hour existence checks.", + "params": { + "threshold": 500 + } + } + ] +} diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json new file mode 100644 index 0000000000..bac5d68f65 --- /dev/null +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json @@ -0,0 +1,18 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_max_date_consistent", + "description": "Checks if the MaxDate is the same for all StatVars.", + "validator": "MAX_DATE_CONSISTENT" + }, + { + "rule_id": "check_deleted_records_percent", + "description": "Checks that the percentage of deleted points is within the threshold.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 0.05 + } + } + ] +} From 42a34d0f1967d6908d74fa66f7697eac9b4293fb Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 07:01:59 +0000 Subject: [PATCH 18/21] chore(cdc): clean up redundant assignment and remove superseded validation_config.json - Remove duplicate _INPUT_FILE_PATH assignment in download_files.py - Remove superseded validation_config.json (now replaced by validation_config_county.json and validation_config_census_tract.json) --- .../download_files.py | 1 - .../validation_config.json | 18 ------------------ 2 files changed, 19 deletions(-) delete mode 100644 scripts/us_cdc/environmental_health_toxicology/validation_config.json diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index 1e673eb564..c52984af73 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -82,7 +82,6 @@ def download_with_retry(url, input_file_name): def main(_): """Main function to download the csv files.""" global _INPUT_FILE_PATH - _INPUT_FILE_PATH = os.path.join(_FLAGS.input_file_path) _INPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.input_file_path) Path(_INPUT_FILE_PATH).mkdir(parents=True, exist_ok=True) importname = sys.argv[1] diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config.json b/scripts/us_cdc/environmental_health_toxicology/validation_config.json deleted file mode 100644 index bac5d68f65..0000000000 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config.json +++ /dev/null @@ -1,18 +0,0 @@ -{ - "schema_version": "1.0", - "rules": [ - { - "rule_id": "check_max_date_consistent", - "description": "Checks if the MaxDate is the same for all StatVars.", - "validator": "MAX_DATE_CONSISTENT" - }, - { - "rule_id": "check_deleted_records_percent", - "description": "Checks that the percentage of deleted points is within the threshold.", - "validator": "DELETED_RECORDS_PERCENT", - "params": { - "threshold": 0.05 - } - } - ] -} From aaf3f00d7d3dfc43b90aa1bc60d809e75d2d537c Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Fri, 4 Sep 2026 07:41:07 +0000 Subject: [PATCH 19/21] fix(cdc): address code review findings for robustness, memory streaming, and validation - download_files.py: wrap Socrata record count metadata query in @retry matching chunk downloads - parse_air_quality.py: validate total_rows > 0 to prevent zero-row edge case - parse_air_quality.py: stream Census Tract processing in 500k-row chunks to bound RAM under 2.5GB - parse_air_quality.py: clean up unused numpy import, duplicate MODULE_DIR, unused query string, and dead assignments - validation_config_census_tract.json: add check_max_date_consistent rule --- .../download_files.py | 15 ++- .../parse_air_quality.py | 119 ++++++++++-------- .../validation_config_census_tract.json | 5 + 3 files changed, 80 insertions(+), 59 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index c52984af73..2c01f08967 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -44,6 +44,13 @@ def download_with_retry(url, input_file_name): if chunk: f.write(chunk) + @retry(tries=3, delay=2, backoff=2) + def get_record_count_with_retry(count_url): + logging.info(f"Querying record count from URL: {count_url}") + resp = requests.get(count_url, timeout=60) + resp.raise_for_status() + return json.loads(resp.text)[0]['COLUMN_ALIAS_GUARD__count'] + url_new = None import_found = False try: @@ -57,12 +64,8 @@ def download_with_retry(url, input_file_name): input_file_name = file_info["input_file_name"] logging.info(f"Input File Name {input_file_name}") - get_record_count = requests.get(url_new.replace( - '.csv', record_count_query), - timeout=60) - get_record_count.raise_for_status() - record_count = json.loads( - get_record_count.text)[0]['COLUMN_ALIAS_GUARD__count'] + count_url = url_new.replace('.csv', record_count_query) + record_count = get_record_count_with_retry(count_url) logging.info( f"Numbers of records found for the URL {url_new} is {record_count}" ) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 5adb0adf4b..3932b832a9 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -14,7 +14,6 @@ import json import os -import numpy as np import pandas as pd from absl import app, logging, flags from pathlib import Path @@ -30,10 +29,8 @@ 'config_file', 'gs://unresolved_mcf/cdc/environmental/import_configs.json', 'Config file path') flags.DEFINE_string('output_file_path', 'output', 'Output files path') -_MODULE_DIR = os.path.dirname(os.path.abspath(__file__)) _INPUT_FILE_PATH = None _OUTPUT_FILE_PATH = None -record_count_query = '?$query=select%20count(*)%20as%20COLUMN_ALIAS_GUARD__count' # Mapping of column names in file to StatVar names. STATVARS = { @@ -87,7 +84,6 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): a cleaned csv file """ try: - global output_file_name logging.info(f"import name from command line {importname}") for config in configs: if config["import_name"] == importname: @@ -120,6 +116,10 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): with open(input_file_path, 'r') as f: total_rows = sum(1 for _ in f) - 1 + if total_rows <= 0: + raise ValueError( + f"Input file {input_file_path} contains no data rows (total_rows={total_rows})." + ) base_size = total_rows // num_shards rem_size = total_rows % num_shards shard_sizes = [ @@ -213,58 +213,72 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): f"Finished cleaning file {output_file_name}!" ) else: - data = pd.read_csv(input_file_path) - data["date"] = pd.to_datetime( - data["date"], - format="%d%b%Y", - errors="raise").dt.strftime("%Y-%m-%d") - + chunk_size = 500_000 + first_chunk = True if "PM2.5" in input_file_name: census_tract = "ds_pm" elif "Ozone" in input_file_name: census_tract = "ds_o3" - if "Census" in input_file_name: - if "PM2.5" in input_file_name: - data = pd.melt( - data, - id_vars=[ - 'year', 'date', 'statefips', - 'countyfips', 'ctfips', - 'latitude', 'longitude' - ], - value_vars=[ - str(census_tract + '_pred'), - str(census_tract + '_stdd') - ], - var_name='StatisticalVariable', - value_name='Value') - elif "Ozone" in input_file_name: - data = pd.melt( - data, - id_vars=[ - 'year', 'date', 'statefips', - 'countyfips', 'ctfips', - 'latitude', 'longitude', - census_tract + '_stdd' - ], - value_vars=[ - str(census_tract + '_pred') - ], - var_name='StatisticalVariable', - value_name='Value') - data.rename(columns={ - census_tract + '_stdd': 'Error' - }, - inplace=True) - data['ctfips'] = data['ctfips'].astype( - str).str.zfill(11) - data["dcid"] = "geoId/" + data[ - "ctfips"].astype(str) - data['StatisticalVariable'] = data[ - 'StatisticalVariable'].map(STATVARS) - data.to_csv(output_file_path, - float_format='%.6f', - index=False) + else: + census_tract = None + + for chunk in pd.read_csv(input_file_path, + chunksize=chunk_size): + chunk["date"] = pd.to_datetime( + chunk["date"], + format="%d%b%Y", + errors="raise").dt.strftime("%Y-%m-%d") + + if "Census" in input_file_name: + if "PM2.5" in input_file_name: + chunk = pd.melt( + chunk, + id_vars=[ + 'year', 'date', 'statefips', + 'countyfips', 'ctfips', + 'latitude', 'longitude' + ], + value_vars=[ + str(census_tract + '_pred'), + str(census_tract + '_stdd') + ], + var_name='StatisticalVariable', + value_name='Value') + elif "Ozone" in input_file_name: + chunk = pd.melt( + chunk, + id_vars=[ + 'year', 'date', 'statefips', + 'countyfips', 'ctfips', + 'latitude', 'longitude', + census_tract + '_stdd' + ], + value_vars=[ + str(census_tract + '_pred') + ], + var_name='StatisticalVariable', + value_name='Value') + chunk.rename(columns={ + census_tract + '_stdd': 'Error' + }, + inplace=True) + chunk['ctfips'] = chunk[ + 'ctfips'].astype(str).str.zfill(11) + chunk["dcid"] = "geoId/" + chunk[ + "ctfips"].astype(str) + chunk['StatisticalVariable'] = chunk[ + 'StatisticalVariable'].map(STATVARS) + if first_chunk: + chunk.to_csv(output_file_path, + float_format='%.6f', + index=False) + first_chunk = False + else: + chunk.to_csv(output_file_path, + mode='a', + header=False, + float_format='%.6f', + index=False) logging.info( f"Finished cleaning file {output_file_name}!" ) @@ -279,7 +293,6 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): def main(_): """Main function to generate the cleaned csv file.""" global _INPUT_FILE_PATH, _OUTPUT_FILE_PATH - _INPUT_FILE_PATH = _FLAGS.input_file_path _INPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.input_file_path) Path(_INPUT_FILE_PATH).mkdir(parents=True, exist_ok=True) _OUTPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.output_file_path) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json index fe9e41637e..e77e6bc4eb 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json @@ -1,6 +1,11 @@ { "schema_version": "1.0", "rules": [ + { + "rule_id": "check_max_date_consistent", + "description": "Checks if the MaxDate is the same for all StatVars.", + "validator": "MAX_DATE_CONSISTENT" + }, { "rule_id": "check_deleted_records_percent", "description": "Disable deleted records check because differ is decoupled for high-scale Census Tracts.", From 1e3870097d2c5e514b92153afc050bc51868feff Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 7 Sep 2026 06:09:35 +0000 Subject: [PATCH 20/21] fix(cdc): resolve P1/P2/P3 review items (exception raise, argv[1], zero-row guards, dead code, lint threshold) --- .../download_files.py | 4 +- .../parse_air_quality.py | 39 ++++++++----------- .../validation_config_county.json | 7 ++++ 3 files changed, 25 insertions(+), 25 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/download_files.py b/scripts/us_cdc/environmental_health_toxicology/download_files.py index 2c01f08967..cff3646a1c 100644 --- a/scripts/us_cdc/environmental_health_toxicology/download_files.py +++ b/scripts/us_cdc/environmental_health_toxicology/download_files.py @@ -82,12 +82,12 @@ def get_record_count_with_retry(count_url): raise -def main(_): +def main(argv): """Main function to download the csv files.""" global _INPUT_FILE_PATH _INPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.input_file_path) Path(_INPUT_FILE_PATH).mkdir(parents=True, exist_ok=True) - importname = sys.argv[1] + importname = argv[1] logging.info(f'Loading config: {_FLAGS.config_file}') with file_util.FileIO(_FLAGS.config_file, 'r') as f: config = json.load(f) diff --git a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py index 3932b832a9..b74836ffe9 100644 --- a/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py +++ b/scripts/us_cdc/environmental_health_toxicology/parse_air_quality.py @@ -49,27 +49,6 @@ "O3_pop_pred": "PopulationWeighted_Concentration_AirPollutant_Ozone" } -# Mapping of month abbreviations to month numbers. -MONTH_MAP = { - "JAN": 1, - "FEB": 2, - "MAR": 3, - "APR": 4, - "MAY": 5, - "JUN": 6, - "JUL": 7, - "AUG": 8, - "SEP": 9, - "OCT": 10, - "NOV": 11, - "DEC": 12 -} - - -# this method is applicable only for "census tract PM25" -def add_prefix_zero(value, length=11): - return str(value).zfill(length) - def clean_air_quality_data(configs, importname, inputpath, outputpath): """ @@ -85,8 +64,10 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): """ try: logging.info(f"import name from command line {importname}") + import_found = False for config in configs: if config["import_name"] == importname: + import_found = True files = config["files"] for file_info in files: output_file_name = file_info["output_file_name"] @@ -209,6 +190,10 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): header=False, float_format='%.6f', index=False) + if first_chunk: + raise ValueError( + f"Input file {input_file_path} contains no data rows." + ) logging.info( f"Finished cleaning file {output_file_name}!" ) @@ -279,6 +264,10 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): header=False, float_format='%.6f', index=False) + if first_chunk: + raise ValueError( + f"Input file {input_file_path} contains no data rows." + ) logging.info( f"Finished cleaning file {output_file_name}!" ) @@ -286,18 +275,22 @@ def clean_air_quality_data(configs, importname, inputpath, outputpath): logging.error( f"Error cleaning {input_file_name}: {e}") raise + if not import_found: + raise ValueError( + f"Import name '{importname}' not found in configuration") except Exception as e: logging.fatal(f"Error while processing the data: {e}") + raise -def main(_): +def main(argv): """Main function to generate the cleaned csv file.""" global _INPUT_FILE_PATH, _OUTPUT_FILE_PATH _INPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.input_file_path) Path(_INPUT_FILE_PATH).mkdir(parents=True, exist_ok=True) _OUTPUT_FILE_PATH = os.path.join(_MODULE_DIR, _FLAGS.output_file_path) Path(_OUTPUT_FILE_PATH).mkdir(parents=True, exist_ok=True) - importname = sys.argv[1] + importname = argv[1] logging.info(f'Loading config: {_FLAGS.config_file}') with file_util.FileIO(_FLAGS.config_file, 'r') as f: config = json.load(f) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json index bac5d68f65..594ff88eb6 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json @@ -13,6 +13,13 @@ "params": { "threshold": 0.05 } + }, + { + "rule_id": "check_lint_error_count", + "description": "Tolerate transient remote network RPC drops during multi-hour existence checks.", + "params": { + "threshold": 2000 + } } ] } From d48c6c29c1e28c50b540226c06e040d9a9e43f4c Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Tue, 8 Sep 2026 12:19:36 +0000 Subject: [PATCH 21/21] fix(us_cdc/environmental_health_toxicology): remove check_lint_error_count rule from validation configs --- .../validation_config_census_tract.json | 7 ------- .../validation_config_county.json | 7 ------- 2 files changed, 14 deletions(-) diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json index e77e6bc4eb..91f4418647 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_census_tract.json @@ -10,13 +10,6 @@ "rule_id": "check_deleted_records_percent", "description": "Disable deleted records check because differ is decoupled for high-scale Census Tracts.", "enabled": false - }, - { - "rule_id": "check_lint_error_count", - "description": "Tolerate transient remote network RPC drops during multi-hour existence checks.", - "params": { - "threshold": 500 - } } ] } diff --git a/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json index 594ff88eb6..bac5d68f65 100644 --- a/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json +++ b/scripts/us_cdc/environmental_health_toxicology/validation_config_county.json @@ -13,13 +13,6 @@ "params": { "threshold": 0.05 } - }, - { - "rule_id": "check_lint_error_count", - "description": "Tolerate transient remote network RPC drops during multi-hour existence checks.", - "params": { - "threshold": 2000 - } } ] }