From bcbd9f2985b381799e64298610f5dd6e570054df Mon Sep 17 00:00:00 2001 From: Ashley Childress Date: Fri, 28 Aug 2026 22:33:16 -0400 Subject: [PATCH] refactor: define the released ADS-B column order in one place - replace two independent copies with FINAL_COLUMN_ORDER - polars concatenates by position, so a forked copy corrupted releases silently Generated-by: Claude Opus 5 --- src/adsb/compress_adsb_to_aircraft_data.py | 6 +++++- src/adsb/concat_parquet_to_final.py | 14 +++++++++----- 2 files changed, 14 insertions(+), 6 deletions(-) diff --git a/src/adsb/compress_adsb_to_aircraft_data.py b/src/adsb/compress_adsb_to_aircraft_data.py index cf995b7..d32cb54 100644 --- a/src/adsb/compress_adsb_to_aircraft_data.py +++ b/src/adsb/compress_adsb_to_aircraft_data.py @@ -4,6 +4,10 @@ import polars as pl COLUMNS = ['dbFlags', 'ownOp', 'year', 'desc', 'aircraft_category', 'r', 't'] +# Positional contract for every released ADS-B artifact. polars concatenates by +# position after .select(), so a divergent copy corrupts output without erroring. +FINAL_COLUMN_ORDER = ['time', 'icao', 'r', 't', 'dbFlags', 'ownOp', 'year', 'desc', 'aircraft_category'] + def compress_df_polars(df: pl.DataFrame, icao: str) -> pl.DataFrame: """Compress a single ICAO group to its most informative row using Polars.""" @@ -164,7 +168,7 @@ def load_parquet_part(part_id: int, date: str) -> pl.DataFrame: print(f"Loading from parquet: {parquet_file}") df = pl.read_parquet( parquet_file, - columns=['time', 'icao', 'r', 't', 'dbFlags', 'ownOp', 'year', 'desc', 'aircraft_category'] + columns=FINAL_COLUMN_ORDER ) # Convert to timezone-naive datetime diff --git a/src/adsb/concat_parquet_to_final.py b/src/adsb/concat_parquet_to_final.py index 21b86e4..1284cb2 100644 --- a/src/adsb/concat_parquet_to_final.py +++ b/src/adsb/concat_parquet_to_final.py @@ -2,8 +2,10 @@ from pathlib import Path import polars as pl import argparse import os + +from src.adsb.compress_adsb_to_aircraft_data import FINAL_COLUMN_ORDER + OUTPUT_DIR = Path("./data/output") -CORRECT_ORDER_OF_COLUMNS = ["time", "icao", "r", "t", "dbFlags", "ownOp", "year", "desc", "aircraft_category"] def main(): parser = argparse.ArgumentParser(description="Concatenate compressed parquet files for a single day") @@ -23,7 +25,7 @@ def main(): df = pl.concat(frames, how="vertical", rechunk=True) df = df.sort(["time", "icao"]) - df = df.select(CORRECT_ORDER_OF_COLUMNS) + df = df.select(FINAL_COLUMN_ORDER) output_path = OUTPUT_DIR / f"openairframes_adsb_{args.date}.parquet" print(f"Writing combined parquet to {output_path} with {df.height} rows") @@ -32,6 +34,8 @@ def main(): csv_output_path = OUTPUT_DIR / f"openairframes_adsb_{args.date}.csv.gz" print(f"Writing combined csv.gz to {csv_output_path} with {df.height} rows") df.write_csv(csv_output_path, compression="gzip") + else: + print(f"No parquet files found in {date_dir}") if args.concat_with_latest_csv: print("Loading latest CSV from GitHub releases to concatenate with...") @@ -50,15 +54,15 @@ def main(): print("Writing latest CSV directly without concatenation to avoid duplicates") os.makedirs(OUTPUT_DIR, exist_ok=True) final_csv_output_path = OUTPUT_DIR / f"openairframes_adsb_{csv_start_date}_{csv_end_date}.csv.gz" - df_latest_csv = df_latest_csv.select(CORRECT_ORDER_OF_COLUMNS) + df_latest_csv = df_latest_csv.select(FINAL_COLUMN_ORDER) df_latest_csv.write_csv(final_csv_output_path, compression="gzip") else: print(f"Concatenating latest CSV (through {csv_end_date}) with new data ({args.date})") # Ensure column order matches before concatenating - df_latest_csv = df_latest_csv.select(CORRECT_ORDER_OF_COLUMNS) + df_latest_csv = df_latest_csv.select(FINAL_COLUMN_ORDER) from src.adsb.compress_adsb_to_aircraft_data import concat_compressed_dfs df_final = concat_compressed_dfs(df_latest_csv, df) - df_final = df_final.select(CORRECT_ORDER_OF_COLUMNS) + df_final = df_final.select(FINAL_COLUMN_ORDER) final_csv_output_path = OUTPUT_DIR / f"openairframes_adsb_{csv_start_date}_{args.date}.csv.gz" df_final.write_csv(final_csv_output_path, compression="gzip") print(f"Final CSV written to {final_csv_output_path}")