From b37ef9a2ea0f45602df71d080f23cc8036d96f3f Mon Sep 17 00:00:00 2001 From: Abhishek Jaiswal Date: Mon, 31 Aug 2026 09:58:33 +0000 Subject: [PATCH 1/7] 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 2/7] 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 3/7] 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 4/7] 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 5/7] 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 6/7] 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 7/7] 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": [ {