diff --git a/quantaq_cli/schema.py b/quantaq_cli/schema.py index 9bf0e51..06c6ce6 100644 --- a/quantaq_cli/schema.py +++ b/quantaq_cli/schema.py @@ -1,41 +1,43 @@ import contextlib +import json from loguru import logger import numpy as np +import pandas as pd import pandera.pandas as pa # Default dtypes COLUMN_DEFINITIONS = [ # --- OPC colunms --- - ('opc_bin0', np.float64), - ('opc_bin1', np.float64), - ('opc_bin2', np.float64), - ('opc_bin3', np.float64), - ('opc_bin4', np.float64), - ('opc_bin5', np.float64), - ('opc_bin6', np.float64), - ('opc_bin7', np.float64), - ('opc_bin8', np.float64), - ('opc_bin9', np.float64), - ('opc_bin10', np.float64), - ('opc_bin11', np.float64), - ('opc_bin12', np.float64), - ('opc_bin13', np.float64), - ('opc_bin14', np.float64), - ('opc_bin15', np.float64), - ('opc_bin16', np.float64), - ('opc_bin17', np.float64), - ('opc_bin18', np.float64), - ('opc_bin19', np.float64), - ('opc_bin20', np.float64), - ('opc_bin21', np.float64), - ('opc_bin22', np.float64), - ('opc_bin23', np.float64), - ('opc_bin1MToF', np.float64), - ('opc_bin3MToF', np.float64), - ('opc_bin5MToF', np.float64), - ('opc_bin7MToF', np.float64), + ('bin0', np.float64), + ('bin1', np.float64), + ('bin2', np.float64), + ('bin3', np.float64), + ('bin4', np.float64), + ('bin5', np.float64), + ('bin6', np.float64), + ('bin7', np.float64), + ('bin8', np.float64), + ('bin9', np.float64), + ('bin10', np.float64), + ('bin11', np.float64), + ('bin12', np.float64), + ('bin13', np.float64), + ('bin14', np.float64), + ('bin15', np.float64), + ('bin16', np.float64), + ('bin17', np.float64), + ('bin18', np.float64), + ('bin19', np.float64), + ('bin20', np.float64), + ('bin21', np.float64), + ('bin22', np.float64), + ('bin23', np.float64), + ('bin1MToF', np.float64), + ('bin3MToF', np.float64), + ('bin5MToF', np.float64), + ('bin7MToF', np.float64), ('opc_temp', np.float64), ('opc_rh', np.float64), ('opc_pm1', np.float64), @@ -116,6 +118,14 @@ ('iteration', np.int16), ('dd_measurement_state', np.float64), # needs to be nullable for older data ('dd_operating_state', np.float64), # needs to be nullable for older data + + # -- Wind columns -- + ('wx_u', np.float64), + ('wx_u', np.float64), + ('wx_wd', np.float64), + ('wx_ws', np.float64), + ('wx_ws_scalar', np.float64), + ] STATIC_COLUMN_RENAMES = { @@ -155,11 +165,20 @@ "rh": "sample_rh", # --- Device / metadata columns --- - "operating_state": "dd_operating_state" + "operating_state": "dd_operating_state", + + # -- Wind columns --- + "w_u": "wx_u", + "w_v": "wx_v", + "u": "wx_u", + "v": "wx_v", + "wd": "wx_wd", + "ws_vector": "wx_ws", + "ws" : "wx_ws_scalar" } # Prefixes for unstandardized column names -COLUMN_RENAME_PREFIXES = ("bin", "opc.bin", "met.", "gases.", "geo.") +COLUMN_RENAME_PREFIXES = ("opc.bin", "met.", "gases.", "geo.") def validate_schema(df, nullable=True, required=False, coerce_dtypes=True, coerce_rename=True): @@ -179,16 +198,23 @@ def validate_schema(df, nullable=True, required=False, coerce_dtypes=True, coerc Returns: df (pd.DataFrame): the DataFrame with validated dtypes and column names """ - df = df.copy() columns = { - name: pa.Column(dtype, nullable=nullable, required=required) + name: pa.Column(dtype, nullable=nullable, required=required, coerce=coerce_dtypes) for name, dtype in COLUMN_DEFINITIONS } + expected_dtypes = dict(COLUMN_DEFINITIONS) legacy_names = set(STATIC_COLUMN_RENAMES.keys()) prefix_patterns = COLUMN_RENAME_PREFIXES + def _wrong_dtype_columns(df): + """Return columns whose dtype doesn't match their expected dtype.""" + return sorted( + col for col in df.columns + if col in expected_dtypes and df[col].dtype != np.dtype(expected_dtypes[col]) + ) + def _has_legacy_column_names(df): """Return True if any legacy names/prefixes remain.""" has_legacy_name = bool(legacy_names.intersection(df.columns)) @@ -201,6 +227,40 @@ def _no_legacy_column_names(df): """Pandera check: True if no legacy names/prefixes remain.""" return not _has_legacy_column_names(df) + def _expected_rename(col): + """Predict what standardize_columns() would rename this column to.""" + if col in STATIC_COLUMN_RENAMES: + return STATIC_COLUMN_RENAMES[col] + if col.startswith("opc.bin"): + return col.replace("opc.", "") + if col.startswith("met."): + return col.replace("met.", "") + if col.startswith("gases."): + return "ox_diff" if col == "gases.o3.diff" else col.removeprefix("gases.").replace(".", "_") + if col.startswith("geo."): + return col.removeprefix("geo.") + return "unknown rename rule" + + def _log_failures(df): + """Log one line per failure: dtype mismatches and legacy column names.""" + wrong_dtype_cols = _wrong_dtype_columns(df) + for col in wrong_dtype_cols: + logger.info( + "Schema validation failed - wrong dtype: '{}' should be {}, got {}", + col, np.dtype(expected_dtypes[col]), df[col].dtype, + ) + + if _has_legacy_column_names(df): + legacy_cols = ( + sorted(legacy_names.intersection(df.columns)) + + sorted(col for col in df.columns for p in prefix_patterns if col.startswith(p)) + ) + for col in legacy_cols: + logger.info( + "Schema validation failed - unstandardized column name: should be '{}', got '{}'", + _expected_rename(col), col, + ) + # index = None means no index is specified # strict = False allows missing and extra columns in the DataFrame schema = pa.DataFrameSchema( @@ -214,34 +274,40 @@ def _no_legacy_column_names(df): ) try: - # print all schema errors instead of raising on the first error schema.validate(df, lazy=True) except pa.errors.SchemaErrors as err: - logger.error("Schema validation failed.") - logger.error("{}", err.failure_cases.to_string()) + logger.bind(schema_errors=err.message).error("Schema validation failed") + _log_failures(df) - # Check the actual condition directly rather than parsing - # failure_cases["check"], since pandera's naming of anonymous - # check functions in that column isn't a stable contract. if coerce_rename and _has_legacy_column_names(df): logger.warning("Standardizing unstandardized column names.") df = standardize_columns(df) - if coerce_dtypes: + if coerce_dtypes and _wrong_dtype_columns(df): logger.warning("Coercing dtypes to expected types.") - dtype_map = { - col: dtype for col, dtype in COLUMN_DEFINITIONS - if col in df.columns - } - df = df.astype(dtype_map) + for col, dtype in COLUMN_DEFINITIONS: + if col not in df.columns: + continue + if np.issubdtype(np.dtype(dtype), np.number): + before_na = df[col].isna().sum() + df[col] = pd.to_numeric(df[col], errors="coerce").astype(dtype) + new_na = df[col].isna().sum() - before_na + if new_na > 0: + logger.warning( + "Coerced {} unparseable value(s) in '{}' to NaN.", + new_na, col, + ) + else: + df[col] = df[col].astype(dtype) if coerce_rename or coerce_dtypes: try: schema.validate(df, lazy=True) logger.info("Schema validation passed after coercion.") except pa.errors.SchemaErrors as second_err: - logger.error("Schema validation still failing after coercion.") - logger.error("{}", second_err.failure_cases.to_string()) + logger.bind(schema_errors=second_err.message).error("Schema validation still failing after coercion.") + _log_failures(df) + raise ValueError(f"Schema validation failed: {json.dumps(second_err.message)}") from second_err return df @@ -273,11 +339,6 @@ def standardize_columns(df): # STATIC_COLUMN_RENAMES dict used by validate_schema column_renames = dict(STATIC_COLUMN_RENAMES) - # bin0 --> opc_bin0, etc.. - for column in df.columns: - if column.startswith("bin"): # opc_bin0, opc_bin23 - column_renames[column] = column.replace("bin", "opc_bin") - # Add diff columns (we minus ae) if they don't already exist if not any('diff' in col for col in df.columns): for pollutant in ("co", "no", "no2"): @@ -291,8 +352,8 @@ def standardize_columns(df): # CloudAPI schema to database schema for column in df.columns: - if column.startswith("opc.bin"): # opc_bin0, opc_bin23 - column_renames[column] = column.replace(".", "_") + if column.startswith("opc.bin"): + column_renames[column] = column.replace("opc.", "") elif column.startswith("met."): column_renames[column] = column.replace("met.", "") elif column.startswith("gases."): diff --git a/quantaq_cli/toolkit/munge.py b/quantaq_cli/toolkit/munge.py index 4a3c4be..c712bf5 100644 --- a/quantaq_cli/toolkit/munge.py +++ b/quantaq_cli/toolkit/munge.py @@ -24,6 +24,6 @@ def clean_dataframe(df, coerce_dtypes=True, coerce_rename=True): # Validate the schema # by default this also coerces dtypes and column names, but that can be overrided - df = validate_schema(df, coerce_dtypes=True, coerce_rename=True) + df = validate_schema(df, coerce_dtypes=coerce_dtypes, coerce_rename=coerce_rename) return df diff --git a/quantaq_cli/toolkit/resample.py b/quantaq_cli/toolkit/resample.py index 75fddf3..bf54847 100644 --- a/quantaq_cli/toolkit/resample.py +++ b/quantaq_cli/toolkit/resample.py @@ -205,7 +205,7 @@ def resample_dataframe( if have_uv and (df[u_col].isna().all() or df[v_col].isna().all()): logger.debug( "All wind components contain NaNs ({}: {}, {}: {}); " - "Deriving them from the averaged u/v components", + "cannot vector-average from u/v - will derive from speed/direction instead", u_col, int(df[u_col].isna().sum()), v_col, int(df[v_col].isna().sum()), ) diff --git a/tests/test_schema.py b/tests/test_schema.py index 121be84..b84db7a 100644 --- a/tests/test_schema.py +++ b/tests/test_schema.py @@ -85,34 +85,34 @@ def test_schema_migration_modpm_rawsd(self): 'sample_rh', 'sample_temp', 'sample_pres', - 'opc_bin0', - 'opc_bin1', - 'opc_bin2', - 'opc_bin3', - 'opc_bin4', - 'opc_bin5', - 'opc_bin6', - 'opc_bin7', - 'opc_bin8', - 'opc_bin9', - 'opc_bin10', - 'opc_bin11', - 'opc_bin12', - 'opc_bin13', - 'opc_bin14', - 'opc_bin15', - 'opc_bin16', - 'opc_bin17', - 'opc_bin18', - 'opc_bin19', - 'opc_bin20', - 'opc_bin21', - 'opc_bin22', - 'opc_bin23', - 'opc_bin1MToF', - 'opc_bin3MToF', - 'opc_bin5MToF', - 'opc_bin7MToF', + 'bin0', + 'bin1', + 'bin2', + 'bin3', + 'bin4', + 'bin5', + 'bin6', + 'bin7', + 'bin8', + 'bin9', + 'bin10', + 'bin11', + 'bin12', + 'bin13', + 'bin14', + 'bin15', + 'bin16', + 'bin17', + 'bin18', + 'bin19', + 'bin20', + 'bin21', + 'bin22', + 'bin23', + 'bin1MToF', + 'bin3MToF', + 'bin5MToF', + 'bin7MToF', 'opc_sample_period', 'opc_sample_flow', 'opc_temp', @@ -236,30 +236,30 @@ def test_schema_migration_modx_api(self): 'neph_pm1_env', 'neph_pm10_env', 'neph_pm25_env', - 'opc_bin0', - 'opc_bin1', - 'opc_bin10', - 'opc_bin11', - 'opc_bin12', - 'opc_bin13', - 'opc_bin14', - 'opc_bin15', - 'opc_bin16', - 'opc_bin17', - 'opc_bin18', - 'opc_bin19', - 'opc_bin2', - 'opc_bin20', - 'opc_bin21', - 'opc_bin22', - 'opc_bin23', - 'opc_bin3', - 'opc_bin4', - 'opc_bin5', - 'opc_bin6', - 'opc_bin7', - 'opc_bin8', - 'opc_bin9', + 'bin0', + 'bin1', + 'bin10', + 'bin11', + 'bin12', + 'bin13', + 'bin14', + 'bin15', + 'bin16', + 'bin17', + 'bin18', + 'bin19', + 'bin2', + 'bin20', + 'bin21', + 'bin22', + 'bin23', + 'bin3', + 'bin4', + 'bin5', + 'bin6', + 'bin7', + 'bin8', + 'bin9', 'opc_pm1', 'opc_pm10', 'opc_pm25'