diff --git a/scripts/us_epa/national_emissions_inventory/README.md b/scripts/us_epa/national_emissions_inventory/README.md index a5b96ab415..2a334e2b29 100644 --- a/scripts/us_epa/national_emissions_inventory/README.md +++ b/scripts/us_epa/national_emissions_inventory/README.md @@ -37,7 +37,7 @@ These are the attributes that we will use | pollutant type(s) | The type of Gas generated which pollutes the Air. | ### Cleaned Data -Cleaned data will be inside [output/national_emissions.csv] as a CSV file with the following columns. +Cleaned data will be inside [gcs_output/output_files/national_emissions.csv] as a CSV file with the following columns. - year - geo_Id @@ -48,8 +48,8 @@ Cleaned data will be inside [output/national_emissions.csv] as a CSV file with t ### MCFs and Template MCFs -- [output/national_emissions.mcf] -- [output/national_emissions.tmcf] +- [gcs_output/output_files/national_emissions.mcf] +- [gcs_output/output_files/national_emissions.tmcf] ### Running Tests diff --git a/scripts/us_epa/national_emissions_inventory/__init__.py b/scripts/us_epa/national_emissions_inventory/__init__.py new file mode 100644 index 0000000000..f8cb6d9d20 --- /dev/null +++ b/scripts/us_epa/national_emissions_inventory/__init__.py @@ -0,0 +1,13 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. diff --git a/scripts/us_epa/national_emissions_inventory/config.py b/scripts/us_epa/national_emissions_inventory/config.py index 8416e9267e..82cc50b46b 100644 --- a/scripts/us_epa/national_emissions_inventory/config.py +++ b/scripts/us_epa/national_emissions_inventory/config.py @@ -68,7 +68,6 @@ 'pollutant_code': 'pollutant code', 'total_emissions': 'total emissions', 'emissions_uom': 'emissions uom', - 'total emissions': 'observation', 'pollutant_type': 'pollutant type(s)' } @@ -90,14 +89,10 @@ 'fips state/county code': 'fips code', 'scc': 'scc', 'pollutant code': 'pollutant code', - 'total emissions': 'observation', - 'uom': 'unit' + 'total emissions': 'total emissions', + 'uom': 'emissions uom' } -drop_tribes = [ - 'state', 'fips state code', 'data category', 'reporting period', - 'emissions operating type', 'pollutant desc', 'data set' -] drop_df = [ 'scc', 'pollutant code', 'emissions type code', 'pollutant type(s)', 'fips code' @@ -1500,7 +1495,6 @@ "1": "External Combustion", "2": "Internal Combustion Engines", "3": "Industrial Processes", - "4": "Chemical Evaporation", "4": "Petroleum And Solvent Evaporation", "5": "Waste Disposal", "6": "MACT Source Categories", diff --git a/scripts/us_epa/national_emissions_inventory/manifest.json b/scripts/us_epa/national_emissions_inventory/manifest.json index aa997fecf3..e1608e2a68 100644 --- a/scripts/us_epa/national_emissions_inventory/manifest.json +++ b/scripts/us_epa/national_emissions_inventory/manifest.json @@ -12,7 +12,9 @@ "process.py" ], "source_files": [ - "gcs_output/input_files/*/*.csv" + "gcs_output/input_files/*/*.csv", + "manifest.json", + "validation_config.json" ], "import_inputs": [ { @@ -21,13 +23,12 @@ } ], "cron_schedule": "0 0 1 1-12/3 *", + "validation_config_file": "validation_config.json", "resource_limits": { "cpu": 8, "memory": 128, - "disk": 100 + "disk": 300 } } ] } - - diff --git a/scripts/us_epa/national_emissions_inventory/process.py b/scripts/us_epa/national_emissions_inventory/process.py index 3a14191517..72b7a54ed2 100644 --- a/scripts/us_epa/national_emissions_inventory/process.py +++ b/scripts/us_epa/national_emissions_inventory/process.py @@ -16,11 +16,13 @@ and generates cleaned CSV, MCF, TMCF file. """ +import gc import os +import shutil import sys import time -# import shutil -# import tempfile +import traceback +import uuid import concurrent.futures from absl import app, flags, logging import pandas as pd @@ -71,7 +73,7 @@ "value: C:national_emissions->observation\n") TRIBAL_GEOCODE_START_RANGE = 80000 -MAX_WORKERS = os.cpu_count() +MAX_WORKERS = min(8, os.cpu_count() or 1) class USAirEmissionTrends: @@ -125,41 +127,54 @@ def _regularize_columns(self, df: pd.DataFrame, df.rename(columns=replacement_08_11, inplace=True) df['pollutant type(s)'] = 'nan' if 'event' in file_path: - df.loc[:, 'emissions type code'] = '' + df['emissions type code'] = '' elif 'process' in file_path: df = df.dropna(subset=['fips code']) - df.loc[:, 'emissions type code'] = '' + df['emissions type code'] = '' if '2008' in file_path: - df.loc[:, 'year'] = '2008' + df['year'] = '2008' else: - df.loc[:, 'year'] = '2011' + df['year'] = '2011' elif '2017' in file_path: if 'Event' in file_path: df['pollutant type(s)'] = 'nan' - elif 'point' in file_path: - if 'unknown' in file_path or '678910' in file_path: - df.rename(columns=replacement_point_17, inplace=True) - df.loc[:, 'emissions type code'] = '' + elif 'point_' in os.path.basename( + file_path) or 'facility_process' in file_path: + df.rename(columns=replacement_point_17, inplace=True) + df['emissions type code'] = '' + elif 'nonpoint' in file_path: + df['emissions type code'] = '' df['year'] = '2017' elif '2020' in file_path: if 'Event' in file_path: df['pollutant type(s)'] = 'nan' - elif 'point' in file_path: - if 'unknown' in file_path: - df.rename(columns=replacement_20, inplace=True) - df.loc[:, 'emissions type code'] = '' + elif 'point_' in os.path.basename( + file_path) or 'facility_process' in file_path: + df.rename(columns=replacement_20, inplace=True) + df['emissions type code'] = '' + elif 'nonpoint' in file_path: + df['emissions type code'] = '' df['year'] = '2020' elif 'tribes' in file_path: - df.rename(columns=replacement_tribes, inplace=True) - df = self._data_standardize(df, 'fips code') + if 'fips code' not in df.columns and 'tribal name' in df.columns: + df.rename(columns=replacement_tribes, inplace=True) + df = self._data_standardize(df, 'fips code') + else: + df.rename(columns=replacement_14, inplace=True) + if 'event' in file_path or 'process' in file_path: + df['emissions type code'] = '' df['pollutant type(s)'] = 'nan' df['year'] = '2014' - else: + elif '2014' in file_path: df.rename(columns=replacement_14, inplace=True) if 'event' in file_path or 'process' in file_path: - df.loc[:, 'emissions type code'] = '' + df['emissions type code'] = '' df['pollutant type(s)'] = 'nan' df['year'] = '2014' + else: + raise ValueError( + f"Unhandled file path or unexpected survey year in: {file_path}" + ) # Ensure all expected columns exist before subsetting for col in df_columns: @@ -178,57 +193,48 @@ def _national_emissions(self, file_path: str) -> pd.DataFrame: Returns: df (pd.DataFrame): provides the cleaned df as output """ - try: - logging.info(f"Processing file: {file_path}") - df = pd.read_csv(file_path, header=0, low_memory=False) - - pd.set_option('display.max_columns', 14) - df = self._regularize_columns(df, file_path) - df['pollutant code'] = df['pollutant code'].astype(str) - df['geo_Id'] = ([f'{x:05}' for x in df['fips code']]) - - # Convert geo_Id to numeric and filter based on range - df['geo_Id'] = pd.to_numeric( - df['geo_Id'], errors='coerce' - ) # Convert to numeric, invalid parsing will be set as NaN - df = df[df['geo_Id'] <= TRIBAL_GEOCODE_START_RANGE] - - # Remove if Tribal Details are needed - df['geo_Id'] = df['geo_Id'].astype(float).astype(int) - df = df.drop(df[df.geo_Id > TRIBAL_GEOCODE_START_RANGE].index) - df['geo_Id'] = ([f'{x:05}' for x in df['geo_Id']]) - df['geo_Id'] = df['geo_Id'].astype(str) - - # Remove if Tribal Details are needed - df['scc'] = df['scc'].astype(str) - df['scc'] = np.where(df['scc'].str.len() == 10, df['scc'].str[0:2], - df['scc'].str[0]) - df['geo_Id'] = 'geoId/' + df['geo_Id'] - df.rename(columns=replacement_17, inplace=True) - df_pollutants = df[df['pollutant code'].isin(pollutants)] - df_pollutants = self._data_standardize(df_pollutants, - 'pollutant code') - df['pollutant code'] = '' - df = pd.concat([df, df_pollutants]) - df = self._data_standardize(df, 'unit') - df['scc_name'] = df['scc'].astype(str) - df = df.replace({'scc_name': replace_source_metadata}) - df['scc_name'] = df['scc_name'].str.replace(' ', '') - df['SV'] = ('Annual_Amount_Emissions_' + - df['pollutant code'].astype(str) + '_SCC_' + - df['scc'].astype(str)) + '_' + df['scc_name'] - - df['Measurement_Method'] = 'dcAggregate/EPA_NationalEmissionInventory' - df['SV'] = df['SV'].str.replace('_nan', '').str.replace('__', '_') - df = df.drop(columns=drop_df) - df = df.drop(df[df['observation'] == '.'].index) - # safely turn any non-numeric values into NaN - df['observation'] = pd.to_numeric(df['observation'], - errors='coerce') - return df - except Exception as e: - logging.error(f"Error processing file {file_path}: {e}") - return pd.DataFrame() + logging.info(f"Processing file: {file_path}") + df = pd.read_csv(file_path, header=0, low_memory=False) + + pd.set_option('display.max_columns', 14) + df = self._regularize_columns(df, file_path) + df['pollutant code'] = df['pollutant code'].astype(str) + + # Convert fips code to numeric, filter out invalid/tribal codes, + # and format as 5-digit string + df['fips_num'] = pd.to_numeric(df['fips code'], errors='coerce') + df = df.dropna(subset=['fips_num']) + df = df[(df['fips_num'] > 0) & + (df['fips_num'] <= TRIBAL_GEOCODE_START_RANGE)] + df['geo_Id'] = 'geoId/' + df['fips_num'].astype(int).astype( + str).str.zfill(5) + df = df.drop(columns=['fips_num']) + + # Strip trailing .0 from float-parsed SCC codes before extracting level 1 + df['scc'] = df['scc'].astype(str).str.replace(r'\.0$', '', + regex=True).str.strip() + df['scc'] = np.where(df['scc'].str.len() == 10, df['scc'].str[0:2], + df['scc'].str[0]) + df.rename(columns=replacement_17, inplace=True) + df_pollutants = df[df['pollutant code'].isin(pollutants)] + df_pollutants = self._data_standardize(df_pollutants, 'pollutant code') + df['pollutant code'] = '' + df = pd.concat([df, df_pollutants], ignore_index=True) + df = self._data_standardize(df, 'unit') + df['scc_name'] = df['scc'].astype(str) + df = df.replace({'scc_name': replace_source_metadata}) + df['scc_name'] = df['scc_name'].str.replace(' ', '') + df['SV'] = ('Annual_Amount_Emissions_' + + df['pollutant code'].astype(str) + '_SCC_' + + df['scc'].astype(str)) + '_' + df['scc_name'] + + df['Measurement_Method'] = 'dcAggregate/EPA_NationalEmissionInventory' + df['SV'] = df['SV'].str.replace('_nan', '').str.replace('__', '_') + df = df.drop(columns=drop_df) + # safely turn any non-numeric values into NaN and drop them + df['observation'] = pd.to_numeric(df['observation'], errors='coerce') + df = df.dropna(subset=['observation']) + return df def _process_file(self, file_path: str) -> None: """ @@ -236,16 +242,29 @@ def _process_file(self, file_path: str) -> None: """ try: df = self._national_emissions(file_path) - if not df.empty: + if df is not None and not df.empty: + df = df.sort_values(by=[ + 'geo_Id', 'year', 'SV', 'Measurement_Method', 'observation' + ]) + df.dropna(subset=['observation'], inplace=True) + df['observation'] = np.where(df['unit'] == 'Pound', + df['observation'] / 2000, + df['observation']) + df['unit'] = "Ton" + if 'scc_name' in df.columns: + df = df.drop(columns=['scc_name']) + df = df.groupby(['geo_Id', 'year', 'Measurement_Method', 'SV'], + as_index=False)['observation'].sum() + df['unit'] = "Ton" intermediate_file_path = os.path.join( self.temp_dir, - f"{str(datetime.now().timestamp()).replace('.', '_')}_{os.path.basename(file_path)}" - ) - df.to_csv(intermediate_file_path, index=False) + f"{uuid.uuid4().hex}_{os.path.basename(file_path)}.pkl") + df.to_pickle(intermediate_file_path) logging.info( f"Saved intermediate file at : {intermediate_file_path}") except Exception as e: - logging.error(f"Error processing file {file_path}: {e}") + logging.exception(f"Error processing file {file_path}: {e}") + raise def _mcf_property_generator(self) -> None: """ @@ -291,7 +310,7 @@ def _mcf_property_generator(self) -> None: emission_type=code) + "\n" logging.info("MCF properties generation complete.") - def _process(self): + def _process(self) -> None: """ This Method processes the input files to generate the final df. @@ -301,40 +320,82 @@ def _process(self): None """ logging.info("Starting data processing across all input files.") + if self.temp_dir and not os.path.exists(self.temp_dir): + os.makedirs(self.temp_dir, exist_ok=True) with concurrent.futures.ThreadPoolExecutor( max_workers=MAX_WORKERS) as executor: - executor.map(self._process_file, self._input_files) + list(executor.map(self._process_file, self._input_files)) logging.info("Consolidating intermediate files.") intermediate_files = [ - os.path.join(self.temp_dir, f) for f in os.listdir(self.temp_dir) + os.path.join(self.temp_dir, f) + for f in os.listdir(self.temp_dir) + if not f.startswith('.') ] - dfs = [] - for f in intermediate_files: - try: - dfs.append(pd.read_csv(f, low_memory=False)) - logging.info(f"Appending {f}") - except Exception as e: - logging.error(f"Error reading intermediate file {f}: {e}") - - if not dfs: - logging.error("No dataframes to concatenate. Exiting.") - return - - self.final_df = pd.concat(dfs, ignore_index=True) + if not intermediate_files: + raise FileNotFoundError( + "No intermediate files found to concatenate.") + + chunk_size = 10 + chunk_dfs = [] + for i in range(0, len(intermediate_files), chunk_size): + batch = intermediate_files[i:i + chunk_size] + batch_dfs = [] + for f in batch: + try: + if f.endswith('.pkl'): + df = pd.read_pickle(f) + else: + df = pd.read_csv(f, low_memory=False) + batch_dfs.append(df) + logging.info(f"Appending {f}") + except Exception as e: + logging.error( + f"Error reading intermediate file {f}: {e}\n{traceback.format_exc()}" + ) + raise + if batch_dfs: + batch_concat = pd.concat(batch_dfs, ignore_index=True) + del batch_dfs + batch_concat = batch_concat.sort_values(by=[ + 'geo_Id', 'year', 'SV', 'Measurement_Method', 'observation' + ]) + batch_concat.dropna(subset=['observation'], inplace=True) + batch_concat['observation'] = np.where( + batch_concat['unit'] == 'Pound', + batch_concat['observation'] / 2000, + batch_concat['observation']) + if 'scc_name' in batch_concat.columns: + batch_concat = batch_concat.drop(columns=['scc_name']) + batch_concat = batch_concat.groupby( + ['geo_Id', 'year', 'Measurement_Method', 'SV'], + as_index=False)['observation'].sum() + batch_concat['unit'] = "Ton" + chunk_dfs.append(batch_concat) + del batch_concat + gc.collect() + + if not chunk_dfs: + raise RuntimeError( + "No dataframes to concatenate after processing intermediate files." + ) + + self.final_df = pd.concat(chunk_dfs, ignore_index=True) + del chunk_dfs + gc.collect() self.final_df = self.final_df.sort_values( by=['geo_Id', 'year', 'SV', 'Measurement_Method', 'observation']) - self.final_df['observation'].replace('', np.nan, inplace=True) self.final_df.dropna(subset=['observation'], inplace=True) self.final_df['observation'] = np.where( self.final_df['unit'] == 'Pound', self.final_df['observation'] / 2000, self.final_df['observation']) - self.final_df = self.final_df.groupby( - ['geo_Id', 'year', 'Measurement_Method', 'SV']).sum().reset_index() - self.final_df['unit'] = "Ton" if 'scc_name' in self.final_df.columns: self.final_df = self.final_df.drop(columns=['scc_name']) + self.final_df = self.final_df.groupby( + ['geo_Id', 'year', 'Measurement_Method', 'SV'], + as_index=False)['observation'].sum() + self.final_df['unit'] = "Ton" logging.info("Data processing complete.") def generate_tmcf(self) -> None: @@ -407,9 +468,10 @@ def process_files(input_path: str, output_file_path: str, if file.lower().endswith('.csv') ] except Exception as e: - logging.fatal( - f"Error finding input files: {e}. Run the download script first.\n") - sys.exit(1) + logging.error( + f"Error finding input files: {e}. Run the download script first.\n" + f"{traceback.format_exc()}") + raise # Defining Output Files logging.info( @@ -418,8 +480,9 @@ def process_files(input_path: str, output_file_path: str, mcf_name = "national_emissions.mcf" tmcf_name = "national_emissions.tmcf" cleaned_csv_path = os.path.join(output_file_path, csv_name) - if not os.path.exists(intermediate_path): - os.makedirs(intermediate_path, exist_ok=True) + if os.path.exists(intermediate_path): + shutil.rmtree(intermediate_path) + os.makedirs(intermediate_path, exist_ok=True) mcf_path = os.path.join(output_file_path, mcf_name) tmcf_path = os.path.join(output_file_path, tmcf_name) @@ -431,17 +494,25 @@ def process_files(input_path: str, output_file_path: str, loader.generate_mcf() loader.generate_tmcf() except Exception as e: - logging.error(f"An unexpected error occurred: {e}") + logging.error( + f"An unexpected error occurred: {e}\n{traceback.format_exc()}") + raise -def main(_): +def main(argv): """ The main function for the script. """ + if len(argv) > 1: + raise app.UsageError("Too many command-line arguments.") logging.set_verbosity(1) logging.info("Started process script") start_time = time.time() - process_files(FLAGS.input_path, FLAGS.output_path, FLAGS.intermediate_path) + try: + process_files(FLAGS.input_path, FLAGS.output_path, + FLAGS.intermediate_path) + except Exception as e: + logging.fatal(f"Process script failed: {e}", exc_info=True) elapsed_time = time.time() - start_time logging.info(f"Total execution time: {elapsed_time:.2f} seconds") diff --git a/scripts/us_epa/national_emissions_inventory/process_test.py b/scripts/us_epa/national_emissions_inventory/process_test.py index 7b70884dd4..2305051fd6 100644 --- a/scripts/us_epa/national_emissions_inventory/process_test.py +++ b/scripts/us_epa/national_emissions_inventory/process_test.py @@ -12,12 +12,22 @@ # See the License for the specific language governing permissions and # limitations under the License. -import unittest -import os -import tempfile import filecmp +import os import shutil -from .process import * +import sys +import tempfile +import unittest +from unittest import mock + +import numpy as np +import pandas as pd + +_MODULE_DIR = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, _MODULE_DIR) + +from config import df_columns +from process import USAirEmissionTrends, process_files class ProcessEnhancedTest(unittest.TestCase): @@ -30,6 +40,8 @@ def setUp(self): def tearDown(self): shutil.rmtree(self.temp_dir) + if os.path.exists(self.intermediate_path): + shutil.rmtree(self.intermediate_path) def test_script(self): input_path = os.path.join(self.test_data_dir, 'input') @@ -55,5 +67,435 @@ def test_script(self): f"File content mismatch: {expected_file} and {generated_file}") +class RegularizeColumnsTest(unittest.TestCase): + + def setUp(self): + self.loader = USAirEmissionTrends([], '', '', '', '') + self.test_data_dir = os.path.join(_MODULE_DIR, 'test_data') + + def test_regularize_columns_2017_nonpoint_float64_nan(self): + # 2017 nonpoint file with float64 NaN emissions type code + df = pd.DataFrame({ + 'fips code': [1001, 1003], + 'scc': [10100101, 10100201], + 'pollutant code': ['CO', 'NOX'], + 'total emissions': [12.5, 34.2], + 'emissions uom': ['TON', 'TON'], + 'emissions type code': [np.nan, np.nan], + }) + self.assertEqual(df['emissions type code'].dtype, np.float64) + + result = self.loader._regularize_columns(df, + '/path/to/2017_nonpoint.csv') + + self.assertEqual(list(result.columns), df_columns) + self.assertTrue((result['year'] == '2017').all()) + self.assertTrue((result['emissions type code'] == '').all()) + self.assertEqual(result['fips code'].tolist(), [1001, 1003]) + + def test_regularize_columns_2020_nonpoint_float64_nan(self): + # 2020 nonpoint file with float64 NaN emissions type code + df = pd.DataFrame({ + 'fips code': [2013, 2016], + 'scc': [20100101, 20100201], + 'pollutant code': ['SO2', 'VOC'], + 'total emissions': [5.1, 8.7], + 'emissions uom': ['TON', 'TON'], + 'emissions type code': [np.nan, np.nan], + }) + self.assertEqual(df['emissions type code'].dtype, np.float64) + + result = self.loader._regularize_columns( + df, '/path/to/2020_nonpoint_data.csv') + + self.assertEqual(list(result.columns), df_columns) + self.assertTrue((result['year'] == '2020').all()) + self.assertTrue((result['emissions type code'] == '').all()) + self.assertEqual(result['fips code'].tolist(), [2013, 2016]) + + def test_regularize_columns_point_files(self): + # 2017 point_ file with unknown + df_pt17_unk = pd.DataFrame({ + 'fips': [1001], + 'pollutant_code': ['CO'], + 'total_emissions': [15.0], + 'emissions_uom': ['TON'], + 'scc': [10100101], + }) + res_pt17_unk = self.loader._regularize_columns( + df_pt17_unk, '/path/to/2017_point_unknown.csv') + self.assertEqual(list(res_pt17_unk.columns), df_columns) + self.assertEqual(res_pt17_unk['year'].iloc[0], '2017') + self.assertEqual(res_pt17_unk['emissions type code'].iloc[0], '') + self.assertEqual(res_pt17_unk['fips code'].iloc[0], 1001) + + # 2017 point_ file with 678910 + df_pt17_num = pd.DataFrame({ + 'fips': [1002], + 'pollutant_code': ['NOX'], + 'total_emissions': [22.0], + 'emissions_uom': ['TON'], + 'scc': [10100102], + }) + res_pt17_num = self.loader._regularize_columns( + df_pt17_num, '/path/to/2017_point_678910.csv') + self.assertEqual(list(res_pt17_num.columns), df_columns) + self.assertEqual(res_pt17_num['year'].iloc[0], '2017') + self.assertEqual(res_pt17_num['emissions type code'].iloc[0], '') + self.assertEqual(res_pt17_num['fips code'].iloc[0], 1002) + + # 2020 point_ file with unknown + df_pt20_unk = pd.DataFrame({ + 'fips state/county code': [1003], + 'pollutant code': ['SO2'], + 'total emissions': [30.0], + 'uom': ['TON'], + 'scc': [10100103], + }) + res_pt20_unk = self.loader._regularize_columns( + df_pt20_unk, '/path/to/2020_point_unknown.csv') + self.assertEqual(list(res_pt20_unk.columns), df_columns) + self.assertEqual(res_pt20_unk['year'].iloc[0], '2020') + self.assertEqual(res_pt20_unk['emissions type code'].iloc[0], '') + self.assertEqual(res_pt20_unk['fips code'].iloc[0], 1003) + self.assertEqual(res_pt20_unk['total emissions'].iloc[0], 30.0) + self.assertEqual(res_pt20_unk['emissions uom'].iloc[0], 'TON') + + # 2017 point_ file for regions 1-5 (point_12345.csv) with underscore columns + df_pt17_reg15 = pd.DataFrame({ + 'fips': [1005], + 'pollutant_code': ['CO'], + 'total_emissions': [18.0], + 'emissions_uom': ['TON'], + 'scc': [10100105], + }) + res_pt17_reg15 = self.loader._regularize_columns( + df_pt17_reg15, + '/path/to/2017neiJan_facility_process_byregions/point_12345.csv') + self.assertEqual(list(res_pt17_reg15.columns), df_columns) + self.assertEqual(res_pt17_reg15['year'].iloc[0], '2017') + self.assertEqual(res_pt17_reg15['emissions type code'].iloc[0], '') + self.assertEqual(res_pt17_reg15['fips code'].iloc[0], 1005) + self.assertEqual(res_pt17_reg15['total emissions'].iloc[0], 18.0) + self.assertEqual(res_pt17_reg15['emissions uom'].iloc[0], 'TON') + + # 2017 point_ file for regions 1-5 (point_12345.csv) with standard spaced columns + df_pt17_reg15_spaces = pd.DataFrame({ + 'fips code': [1001], + 'pollutant code': ['75070'], + 'pollutant type(s)': ['HAP'], + 'total emissions': [0.265867], + 'emissions uom': ['TON'], + 'scc': [10100101], + }) + res_pt17_reg15_spaces = self.loader._regularize_columns( + df_pt17_reg15_spaces, + '/path/to/2017neiJan_facility_process_byregions/point_12345.csv') + self.assertEqual(list(res_pt17_reg15_spaces.columns), df_columns) + self.assertEqual(res_pt17_reg15_spaces['year'].iloc[0], '2017') + self.assertEqual(res_pt17_reg15_spaces['fips code'].iloc[0], 1001) + self.assertEqual(res_pt17_reg15_spaces['total emissions'].iloc[0], + 0.265867) + self.assertEqual(res_pt17_reg15_spaces['emissions uom'].iloc[0], 'TON') + + # 2020 point_ file for regions 1-10 (point_1.csv ... point_10.csv) + df_pt20_reg1 = pd.DataFrame({ + 'fips state/county code': [1006], + 'pollutant code': ['NOX'], + 'total emissions': [25.0], + 'uom': ['TON'], + 'scc': [10100106], + }) + res_pt20_reg1 = self.loader._regularize_columns( + df_pt20_reg1, + '/path/to/2020nei_facility_process_byregions/point_1.csv') + self.assertEqual(list(res_pt20_reg1.columns), df_columns) + self.assertEqual(res_pt20_reg1['year'].iloc[0], '2020') + self.assertEqual(res_pt20_reg1['emissions type code'].iloc[0], '') + self.assertEqual(res_pt20_reg1['fips code'].iloc[0], 1006) + self.assertEqual(res_pt20_reg1['total emissions'].iloc[0], 25.0) + self.assertEqual(res_pt20_reg1['emissions uom'].iloc[0], 'TON') + + # facility_process file + df_fac = pd.DataFrame({ + 'fips code': [1004], + 'pollutant code': ['PM10-PRI'], + 'total emissions': [4.5], + 'emissions uom': ['TON'], + 'scc': [10100104], + }) + res_fac = self.loader._regularize_columns( + df_fac, '/path/to/2017_facility_process.csv') + self.assertEqual(list(res_fac.columns), df_columns) + self.assertEqual(res_fac['year'].iloc[0], '2017') + self.assertEqual(res_fac['emissions type code'].iloc[0], '') + + def test_regularize_columns_unhandled_year_raises_value_error(self): + df = pd.DataFrame({ + 'fips code': [1001], + 'pollutant code': ['CO'], + 'total emissions': [10.0], + 'emissions uom': ['TON'], + 'scc': [10100101], + }) + with self.assertRaises(ValueError): + self.loader._regularize_columns(df, '/path/to/2023_data.csv') + + def test_regularize_columns_2014_tribes(self): + # Tribal file with 'tribal name' and without 'fips code' + df_tribes = pd.DataFrame({ + 'tribal name': ['Navajo Nation'], + 'scc': [10100101], + 'pollutant code': ['VOC'], + 'total emissions': [10.0], + 'emissions uom': ['TON'], + }) + res_tribes = self.loader._regularize_columns( + df_tribes, '/path/to/2014_tribes_data.csv') + self.assertEqual(list(res_tribes.columns), df_columns) + self.assertEqual(res_tribes['year'].iloc[0], '2014') + self.assertEqual(res_tribes['fips code'].iloc[0], 'Navajo Nation') + self.assertEqual(res_tribes['pollutant type(s)'].iloc[0], 'nan') + + # Tribal process file process_tribes.csv with existing 'fips' column + df_proc_tribes = pd.DataFrame({ + 'fips': [1005], + 'scc': [10100105], + 'pollutant_cd': ['CO'], + 'total_emissions': [18.0], + 'uom': ['TON'], + }) + res_proc_tribes = self.loader._regularize_columns( + df_proc_tribes, '/path/to/process_tribes.csv') + self.assertEqual(list(res_proc_tribes.columns), df_columns) + self.assertEqual(res_proc_tribes['year'].iloc[0], '2014') + self.assertEqual(res_proc_tribes['emissions type code'].iloc[0], '') + self.assertEqual(res_proc_tribes['fips code'].iloc[0], 1005) + self.assertEqual(res_proc_tribes['pollutant type(s)'].iloc[0], 'nan') + + # Tribal event file event_tribes.csv + df_event_tribes = pd.DataFrame({ + 'fips': [1006], + 'scc': [10100106], + 'pollutant_cd': ['PM25-PRI'], + 'total_emissions': [2.5], + 'uom': ['TON'], + }) + res_event_tribes = self.loader._regularize_columns( + df_event_tribes, '/path/to/event_tribes.csv') + self.assertEqual(list(res_event_tribes.columns), df_columns) + self.assertEqual(res_event_tribes['year'].iloc[0], '2014') + self.assertEqual(res_event_tribes['emissions type code'].iloc[0], '') + + def test_national_emissions_fips_padding_and_scc_cleaning(self): + # Verify 4-digit string FIPS ('1001') -> 'geoId/01001', + # float SCC ('10100101.0') -> stripped of '.0' and Level 1 mapped to ExternalCombustion, + # non-numeric observation '.' -> dropped + temp_dir = tempfile.mkdtemp() + try: + csv_path = os.path.join(temp_dir, '2020nei_point_1.csv') + df = pd.DataFrame({ + 'fips state/county code': ['1001', '01003', '1001'], + 'pollutant code': ['CO', 'CO', 'CO'], + 'total emissions': ['12.5', '.', '15.0'], + 'uom': ['TON', 'TON', 'TON'], + 'scc': ['10100101.0', '10100101', '10100101.0'], + }) + df.to_csv(csv_path, index=False) + result = self.loader._national_emissions(csv_path) + + # Check that 4-digit FIPS was padded to 'geoId/01001', NOT 'geoId/10010' + self.assertTrue((result['geo_Id'] == 'geoId/01001').all()) + # Observation '.' was dropped. The 2 remaining rows each generate + # aggregate and pollutant-specific StatVars (total 4 rows) + self.assertEqual(len(result), 4) + self.assertCountEqual(result['observation'].tolist(), + [12.5, 15.0, 12.5, 15.0]) + # SCC was correctly extracted as '1' (External Combustion), not '10' + self.assertTrue( + all('SCC_1_ExternalCombustion' in sv for sv in result['SV'])) + finally: + shutil.rmtree(temp_dir) + + def test_mcf_property_generator(self): + loader = USAirEmissionTrends([], '', '', '', '') + loader.final_df = pd.DataFrame({ + 'SV': [ + 'Annual_Amount_Emissions_CarbonMonoxide_SCC_1_ExternalCombustion', + 'Annual_Amount_Emissions_SCC_21_StationarySourceFuelCombustion' + ], + 'observation': [10.0, 20.0], + 'geo_Id': ['geoId/01001', 'geoId/01001'], + 'year': ['2020', '2020'], + 'Measurement_Method': [ + 'dcAggregate/EPA_NationalEmissionInventory', + 'dcAggregate/EPA_NationalEmissionInventory' + ], + 'unit': ['Ton', 'Ton'] + }) + loader._mcf_property_generator() + mcf = loader.final_mcf_template + self.assertIn( + 'Node: dcid:Annual_Amount_Emissions_CarbonMonoxide_SCC_1_ExternalCombustion', + mcf) + self.assertIn( + 'Node: dcid:Annual_Amount_Emissions_SCC_21_StationarySourceFuelCombustion', + mcf) + self.assertIn('epaSccCode: dcs:EPA_SCC/1', mcf) + self.assertIn('epaSccCode: dcs:EPA_SCC/21', mcf) + self.assertIn('emittedThing: dcs:CarbonMonoxide', mcf) + + def test_worker_exception_bubbling(self): + temp_dir = tempfile.mkdtemp() + inter_dir = tempfile.mkdtemp() + # Unhandled year raises ValueError in _regularize_columns which must bubble up through _process() + bad_file = os.path.join(temp_dir, '2023_unhandled_year.csv') + df = pd.DataFrame({'fips code': [1001]}) + df.to_csv(bad_file, index=False) + try: + loader = USAirEmissionTrends([bad_file], '', '', '', inter_dir) + with self.assertRaises(ValueError): + loader._process() + finally: + shutil.rmtree(temp_dir) + if os.path.exists(inter_dir): + shutil.rmtree(inter_dir) + + def test_intermediate_directory_purged_on_rerun(self): + inter_dir = tempfile.mkdtemp() + temp_out = tempfile.mkdtemp() + stale_file = os.path.join(inter_dir, 'stale_intermediate.csv') + with open(stale_file, 'w') as f: + f.write("stale data") + self.assertTrue(os.path.exists(stale_file)) + + input_path = os.path.join(self.test_data_dir, 'input') + process_files(input_path, temp_out, inter_dir) + + self.assertFalse(os.path.exists(stale_file)) + shutil.rmtree(temp_out) + if os.path.exists(inter_dir): + shutil.rmtree(inter_dir) + + def test_multi_year_processing(self): + temp_in = tempfile.mkdtemp() + temp_out = tempfile.mkdtemp() + inter_dir = tempfile.mkdtemp() + try: + dir_14 = os.path.join(temp_in, '2014') + dir_17 = os.path.join(temp_in, '2017neiJan_facility_process') + dir_20 = os.path.join(temp_in, '2020nei_facility_process') + os.makedirs(dir_14) + os.makedirs(dir_17) + os.makedirs(dir_20) + + df_14 = pd.DataFrame({ + 'state_and_county_fips_code': [1001], + 'pollutant_cd': ['CO'], + 'total_emissions': [5.0], + 'uom': ['TON'], + 'SCC': [10100101], + }) + df_14.to_csv(os.path.join(dir_14, 'process_14.csv'), index=False) + + df_17 = pd.DataFrame({ + 'fips': [1001], + 'pollutant_code': ['CO'], + 'total_emissions': [10.0], + 'emissions_uom': ['TON'], + 'scc': [10100101], + }) + df_17.to_csv(os.path.join(dir_17, 'point_12345.csv'), index=False) + + df_20 = pd.DataFrame({ + 'fips state/county code': [1001], + 'pollutant code': ['CO'], + 'total emissions': [20.0], + 'uom': ['TON'], + 'scc': [10100101], + }) + df_20.to_csv(os.path.join(dir_20, 'point_1.csv'), index=False) + + csv_out = os.path.join(temp_out, 'national_emissions.csv') + mcf_out = os.path.join(temp_out, 'national_emissions.mcf') + tmcf_out = os.path.join(temp_out, 'national_emissions.tmcf') + + ip_files = [ + os.path.join(dir_14, 'process_14.csv'), + os.path.join(dir_17, 'point_12345.csv'), + os.path.join(dir_20, 'point_1.csv') + ] + loader = USAirEmissionTrends(ip_files, csv_out, mcf_out, tmcf_out, + inter_dir) + loader.generate_csv() + loader.generate_mcf() + loader.generate_tmcf() + + self.assertTrue(os.path.exists(csv_out)) + res_df = pd.read_csv(csv_out) + # Each year generates aggregate and CO-specific StatVars (3 * 2 = 6 rows) + self.assertEqual(len(res_df), 6) + self.assertCountEqual( + res_df['year'].astype(str).tolist(), + ['2014', '2014', '2017', '2017', '2020', '2020']) + self.assertCountEqual(res_df['observation'].tolist(), + [5.0, 5.0, 10.0, 10.0, 20.0, 20.0]) + self.assertTrue((res_df['geo_Id'] == 'geoId/01001').all()) + finally: + shutil.rmtree(temp_in) + shutil.rmtree(temp_out) + if os.path.exists(inter_dir): + shutil.rmtree(inter_dir) + + def test_process_raises_on_corrupted_intermediate_file(self): + temp_in = tempfile.mkdtemp() + temp_out = tempfile.mkdtemp() + inter_dir = tempfile.mkdtemp() + try: + dir_17 = os.path.join(temp_in, '2017') + os.makedirs(dir_17) + df_17 = pd.DataFrame({ + 'fips': [1001], + 'pollutant_code': ['CO'], + 'total_emissions': [10.0], + 'emissions_uom': ['TON'], + 'scc': [10100101], + }) + dummy_file = os.path.join(dir_17, 'point_1.csv') + df_17.to_csv(dummy_file, index=False) + + loader = USAirEmissionTrends([dummy_file], + os.path.join(temp_out, 'out.csv'), + os.path.join(temp_out, 'out.mcf'), + os.path.join(temp_out, 'out.tmcf'), + inter_dir) + with mock.patch('pandas.read_pickle', + side_effect=IOError('Corrupt pickle')): + with self.assertRaises(IOError): + loader.generate_csv() + finally: + shutil.rmtree(temp_in) + shutil.rmtree(temp_out) + if os.path.exists(inter_dir): + shutil.rmtree(inter_dir) + + def test_process_files_raises_on_unhandled_year(self): + temp_in = tempfile.mkdtemp() + temp_out = tempfile.mkdtemp() + inter_dir = tempfile.mkdtemp() + try: + # Input file with unhandled year + bad_file = os.path.join(temp_in, '1999_emissions.csv') + with open(bad_file, 'w') as f: + f.write("col1,col2\n1,2\n") + with self.assertRaises(Exception): + process_files(temp_in, temp_out, inter_dir) + finally: + shutil.rmtree(temp_in) + shutil.rmtree(temp_out) + if os.path.exists(inter_dir): + shutil.rmtree(inter_dir) + + if __name__ == '__main__': unittest.main() diff --git a/scripts/us_epa/national_emissions_inventory/test_data/input/esg_cty_scc_12961.csv b/scripts/us_epa/national_emissions_inventory/test_data/input/2014/esg_cty_scc_12961.csv similarity index 100% rename from scripts/us_epa/national_emissions_inventory/test_data/input/esg_cty_scc_12961.csv rename to scripts/us_epa/national_emissions_inventory/test_data/input/2014/esg_cty_scc_12961.csv diff --git a/scripts/us_epa/national_emissions_inventory/test_data/input/nonroad_123.csv b/scripts/us_epa/national_emissions_inventory/test_data/input/2014/nonroad_123.csv similarity index 100% rename from scripts/us_epa/national_emissions_inventory/test_data/input/nonroad_123.csv rename to scripts/us_epa/national_emissions_inventory/test_data/input/2014/nonroad_123.csv diff --git a/scripts/us_epa/national_emissions_inventory/test_data/input/onroad_123.csv b/scripts/us_epa/national_emissions_inventory/test_data/input/2014/onroad_123.csv similarity index 100% rename from scripts/us_epa/national_emissions_inventory/test_data/input/onroad_123.csv rename to scripts/us_epa/national_emissions_inventory/test_data/input/2014/onroad_123.csv diff --git a/scripts/us_epa/national_emissions_inventory/validation_config.json b/scripts/us_epa/national_emissions_inventory/validation_config.json new file mode 100644 index 0000000000..b9b98eb9a2 --- /dev/null +++ b/scripts/us_epa/national_emissions_inventory/validation_config.json @@ -0,0 +1,22 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_deleted_records_percent", + "description": "Verifies that the percentage of deleted records does not exceed the 0.1% threshold, ensuring unintended observation drops (such as missing nonpoint sources which drop >20% of records) are caught while accommodating minor upstream EPA revisions.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 0.1 + } + }, + { + "rule_id": "check_date_freshness", + "description": "Verifies that observation dates span from 2008 to at least 2020 across active StatVars", + "validator": "SQL_VALIDATOR", + "params": { + "query": "SELECT COUNT(DISTINCT StatVar) AS active_svs, MIN(MinDate) AS min_date FROM stats WHERE MaxDate >= '2020'", + "condition": "active_svs >= 385 AND min_date <= '2008'" + } + } + ] +}