"""Orchestration entrypoint — runs load → clean → features end-to-end. Writes all parquet snapshots + meta_summary.json (the contract consumed by the site's 'approach' page). """ import json from pathlib import Path import pandas as pd from .clean import clean_inpatient, clean_outpatient from .features import ( add_features_inpatient, add_features_outpatient, build_provider_combined, validate_observed_mappings, ) from .config import ( COMPATIBILITY_RATIO_ALIAS, DATASET_METADATA, DATA_PROC, SITE_DATA, ) def _provider_overlap(inp: pd.DataFrame, out: pd.DataFrame, basis: str) -> dict: inp_providers = set(inp["provider_ccn"].dropna()) out_providers = set(out["provider_ccn"].dropna()) in_both = inp_providers & out_providers return { "basis": basis, "inpatient_providers": int(len(inp_providers)), "outpatient_providers": int(len(out_providers)), "in_both": int(len(in_both)), "providers_in_both": int(len(in_both)), "inpatient_only": int(len(inp_providers - out_providers)), "outpatient_only": int(len(out_providers - inp_providers)), } def _denominator( df: pd.DataFrame, name: str, definition: str, field: str | None = None, condition: str = "none", ) -> dict: return { "name": name, "unit": "deduplicated provider-service row", "definition": definition, "field": field, "condition": condition, "rows": int(len(df)), "providers": int(df["provider_ccn"].nunique()), } def _geography_conflict_summary( inp_feat: pd.DataFrame, out_full_feat: pd.DataFrame ) -> dict: """Count shared providers whose inpatient geography conflicts with full outpatient geography. `build_provider_combined` resolves conflicts by inpatient precedence, then full outpatient presence, then cost-observed outpatient; this helper records how often each field actually conflicts. """ fields = ("state", "state_fips", "census_region", "urban_rural") inp_geo = ( inp_feat.groupby("provider_ccn")[list(fields)].first().reset_index() ) out_geo = ( out_full_feat.groupby("provider_ccn")[list(fields)].first().reset_index() ) merged = inp_geo.merge(out_geo, on="provider_ccn", suffixes=("_inp", "_out")) conflicts = { field: int((merged[f"{field}_inp"].astype(str).fillna("") != merged[f"{field}_out"].astype(str).fillna("")).sum()) for field in fields } inp_arr = merged[[f"{field}_inp" for field in fields]].astype(str).fillna("").to_numpy() out_arr = merged[[f"{field}_out" for field in fields]].astype(str).fillna("").to_numpy() return { "shared_providers": int(len(merged)), "conflicts_by_field": conflicts, "any_conflict": int((inp_arr != out_arr).any(axis=1).sum()), "precedence": "inpatient geography wins, then full outpatient presence, then cost-observed outpatient summary", "sensitivity_status": "inpatient precedence applied; a reversal sensitivity is not computed", } def run_pipeline() -> dict: print("[1/4] Loading + cleaning inpatient ...") inp = clean_inpatient() print(f" → {inp.shape[0]:,} rows, {inp['provider_ccn'].nunique():,} providers") print("[2/4] Loading + cleaning outpatient (both cost & full) ...") out_cost = clean_outpatient(drop_suppressed=True) out_full = clean_outpatient(drop_suppressed=False) print(f" → cost subset: {out_cost.shape[0]:,} rows | full: {out_full.shape[0]:,} rows") print("[3/4] Feature engineering ...") inp_feat = add_features_inpatient(inp) out_feat = add_features_outpatient(out_cost) out_full_feat = add_features_outpatient(out_full) print(f" → inpatient features: {inp_feat.shape[1]} cols") print(f" → outpatient features: {out_feat.shape[1]} cols") print("[4/4] Building provider_combined ...") provider_combined = build_provider_combined(inp_feat, out_feat, out_full_feat) print(f" → {len(provider_combined):,} providers") # Persist parquet snapshots inp.to_parquet(DATA_PROC / "inpatient_clean.parquet", index=False) out_cost.to_parquet(DATA_PROC / "outpatient_clean.parquet", index=False) out_full.to_parquet(DATA_PROC / "outpatient_full.parquet", index=False) inp_feat.to_parquet(DATA_PROC / "inpatient_features.parquet", index=False) out_feat.to_parquet(DATA_PROC / "outpatient_features.parquet", index=False) out_full_feat.to_parquet(DATA_PROC / "outpatient_full_features.parquet", index=False) provider_combined.to_parquet(DATA_PROC / "provider_combined.parquet", index=False) print(" parquet snapshots saved.") full_denominator = _denominator( out_full, "outpatient_full", "All deduplicated outpatient provider-APC rows, including rows with suppressed avg_allowed_amount.", field="avg_allowed_amount", condition="all rows retained; suppression is represented by null", ) cost_denominator = _denominator( out_cost, "outpatient_cost_observed", "Deduplicated outpatient provider-APC rows with an observed avg_allowed_amount.", field="avg_allowed_amount", condition="avg_allowed_amount is not null", ) full_overlap = _provider_overlap(inp, out_full, "outpatient_full") cost_overlap = _provider_overlap(inp, out_cost, "outpatient_cost_observed") mapping_validation = validate_observed_mappings(inp, out_full) print( " mapping validation: " f"DRG {mapping_validation['drg']['observed_code_count']:,} observed, " f"{mapping_validation['drg']['unmapped_code_count']:,} unmapped, " f"{mapping_validation['drg']['ambiguous_code_count']:,} ambiguous | " f"APC {mapping_validation['apc']['observed_code_count']:,} observed, " f"{mapping_validation['apc']['unmapped_code_count']:,} unmapped, " f"{mapping_validation['apc']['ambiguous_code_count']:,} ambiguous" ) def _maryland_observation(df: pd.DataFrame) -> dict: is_maryland = df["state"].eq("MD") return { "observed": bool(is_maryland.any()), "rows": int(is_maryland.sum()), "providers": int(df.loc[is_maryland, "provider_ccn"].nunique()), } # Meta summary for the website's approach page summary = { "periods": { "inpatient": DATASET_METADATA["inpatient"]["period"], "outpatient": DATASET_METADATA["outpatient"]["period"], }, "inpatient": { **DATASET_METADATA["inpatient"], "rows": int(len(inp)), "providers": int(inp["provider_ccn"].nunique()), "drgs": int(inp["drg_code"].nunique()), "states": int(inp["state"].nunique()), "missing_values": 0, "suppression": "none", "cost_measure": "avg_total_payment", "payment_measure_field": "avg_total_payment", }, "outpatient_full": { **DATASET_METADATA["outpatient"], "rows": int(len(out_full)), "providers": int(out_full["provider_ccn"].nunique()), "apcs": int(out_full["apc_code"].nunique()), "states": int(out_full["state"].nunique()), "denominator": full_denominator, "cost_suppressed_rows": int(out_full["avg_allowed_amount"].isna().sum()), "cost_suppressed_pct": round( out_full["avg_allowed_amount"].isna().mean() * 100, 2), "outlier_suppressed_rows": int(out_full["outlier_services"].isna().sum()), "cost_measure": "avg_allowed_amount", "payment_measure_field": "avg_allowed_amount", }, "outpatient_cost_subset": { **DATASET_METADATA["outpatient"], "rows": int(len(out_cost)), "providers": int(out_cost["provider_ccn"].nunique()), "apcs": int(out_cost["apc_code"].nunique()), "denominator": cost_denominator, "dataset_role": "cost_observed_subset", "note": "rows with avg_allowed_amount observed; CMS cost suppression is not imputed", "cost_measure": "avg_allowed_amount", "payment_measure_field": "avg_allowed_amount", }, "outpatient_denominators": { "full": full_denominator, "cost_observed": cost_denominator, }, "provider_overlap": { # Direct fields retain the previous full-outpatient contract. **full_overlap, "full_outpatient": full_overlap, "cost_observed_outpatient": cost_overlap, "cost_observed": cost_overlap, }, "provider_overlap_full": full_overlap, "provider_overlap_cost_observed": cost_overlap, "provider_combined": { "rows": int(len(provider_combined)), "providers": int(provider_combined["provider_ccn"].nunique()), "has_inpatient": int(provider_combined["has_inpatient"].sum()), "has_outpatient_full": int(provider_combined["has_outpatient_full"].sum()), "has_outpatient_cost_observed": int( provider_combined["has_outpatient_cost_observed"].sum() ), "definition": ( "Full inpatient plus full outpatient provider union; outpatient " "summary metrics are populated only for cost-observed rows." ), "geography_conflicts": _geography_conflict_summary(inp_feat, out_full_feat), }, "maryland_exception": { "flag_field": "is_maryland_exception", "state": "MD", "reason": "Maryland's all-payer rate-setting system is an exception to the standard Medicare outpatient payment context.", "inpatient": _maryland_observation(inp), "outpatient_full": _maryland_observation(out_full), "outpatient_cost_observed": _maryland_observation(out_cost), "inpatient_included": _maryland_observation(inp)["observed"], "outpatient_included": _maryland_observation(out_full)["observed"], "outpatient_excluded": not _maryland_observation(out_full)["observed"], "outpatient_exclusion_note": "Maryland is absent from the CMS outpatient file used here; do not interpret that absence as a zero or an imputed outpatient observation.", }, "payment_measures": { "inpatient": DATASET_METADATA["inpatient"]["payment_measure"], "outpatient": DATASET_METADATA["outpatient"]["payment_measure"], }, "feature_formulas": { "inpatient": DATASET_METADATA["inpatient"]["ratio"], "outpatient": DATASET_METADATA["outpatient"]["ratio"], }, "engineered_features": { "inpatient": ["census_region", "urban_rural", "drg_mdc", "is_maryland_exception", "charge_to_total_payment_ratio"], "outpatient": ["census_region", "urban_rural", "apc_family", "is_maryland_exception", "charge_to_allowed_amount_ratio"], "provider_combined": ["inp_drgs_seen", "out_apcs_seen", "inp_avg_charge_to_total_payment_ratio", "out_avg_charge_to_allowed_amount_ratio", "is_maryland_exception", "has_inpatient", "has_outpatient", "has_outpatient_full", "has_outpatient_cost_observed", "outpatient_presence_basis"], }, "compatibility_aliases": { COMPATIBILITY_RATIO_ALIAS: { "inpatient": "charge_to_total_payment_ratio", "outpatient": "charge_to_allowed_amount_ratio", "status": "legacy alias retained", "reason": "Unchanged analysis modules still consume the historical field name; use the denominator-specific fields for new work.", }, }, "mapping_validation": mapping_validation, } SITE_DATA.mkdir(parents=True, exist_ok=True) with open(SITE_DATA / "meta_summary.json", "w") as f: json.dump(summary, f, indent=2) with open(SITE_DATA / "mapping_validation.json", "w") as f: json.dump(mapping_validation, f, indent=2) with open(SITE_DATA / "denominator_audit.json", "w") as f: json.dump({ "periods": summary["periods"], "outpatient_denominators": summary["outpatient_denominators"], "provider_overlap_full": summary["provider_overlap_full"], "provider_overlap_cost_observed": summary["provider_overlap_cost_observed"], "provider_combined": summary["provider_combined"], "maryland_exception": summary["maryland_exception"], "estimands": { "outpatient_full": "All outpatient provider-APC rows, including suppressed allowed amounts.", "outpatient_cost_observed": "Outpatient provider-APC rows with observed allowed amounts.", "cross_dataset_cost_observed": "Provider overlap using inpatient rows and outpatient cost-observed providers.", "provider_combined_presence": "Full-frame outpatient presence is separate from cost-observed outpatient summary metrics.", }, }, f, indent=2) print(f" meta_summary.json saved.") return { "inpatient_features": inp_feat, "outpatient_features": out_feat, "provider_combined": provider_combined, } if __name__ == "__main__": # pragma: no cover run_pipeline()