mirror of
https://github.com/PlaneQuery/OpenAirframes.git
synced 2026-09-18 11:52:15 +02:00
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 <noreply@anthropic.com>
This commit is contained in:
@@ -4,6 +4,10 @@ import polars as pl
|
|||||||
|
|
||||||
COLUMNS = ['dbFlags', 'ownOp', 'year', 'desc', 'aircraft_category', 'r', 't']
|
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:
|
def compress_df_polars(df: pl.DataFrame, icao: str) -> pl.DataFrame:
|
||||||
"""Compress a single ICAO group to its most informative row using Polars."""
|
"""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}")
|
print(f"Loading from parquet: {parquet_file}")
|
||||||
df = pl.read_parquet(
|
df = pl.read_parquet(
|
||||||
parquet_file,
|
parquet_file,
|
||||||
columns=['time', 'icao', 'r', 't', 'dbFlags', 'ownOp', 'year', 'desc', 'aircraft_category']
|
columns=FINAL_COLUMN_ORDER
|
||||||
)
|
)
|
||||||
|
|
||||||
# Convert to timezone-naive datetime
|
# Convert to timezone-naive datetime
|
||||||
|
|||||||
@@ -2,8 +2,10 @@ from pathlib import Path
|
|||||||
import polars as pl
|
import polars as pl
|
||||||
import argparse
|
import argparse
|
||||||
import os
|
import os
|
||||||
|
|
||||||
|
from src.adsb.compress_adsb_to_aircraft_data import FINAL_COLUMN_ORDER
|
||||||
|
|
||||||
OUTPUT_DIR = Path("./data/output")
|
OUTPUT_DIR = Path("./data/output")
|
||||||
CORRECT_ORDER_OF_COLUMNS = ["time", "icao", "r", "t", "dbFlags", "ownOp", "year", "desc", "aircraft_category"]
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
parser = argparse.ArgumentParser(description="Concatenate compressed parquet files for a single day")
|
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 = pl.concat(frames, how="vertical", rechunk=True)
|
||||||
|
|
||||||
df = df.sort(["time", "icao"])
|
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"
|
output_path = OUTPUT_DIR / f"openairframes_adsb_{args.date}.parquet"
|
||||||
print(f"Writing combined parquet to {output_path} with {df.height} rows")
|
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"
|
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")
|
print(f"Writing combined csv.gz to {csv_output_path} with {df.height} rows")
|
||||||
df.write_csv(csv_output_path, compression="gzip")
|
df.write_csv(csv_output_path, compression="gzip")
|
||||||
|
else:
|
||||||
|
print(f"No parquet files found in {date_dir}")
|
||||||
|
|
||||||
if args.concat_with_latest_csv:
|
if args.concat_with_latest_csv:
|
||||||
print("Loading latest CSV from GitHub releases to concatenate with...")
|
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")
|
print("Writing latest CSV directly without concatenation to avoid duplicates")
|
||||||
os.makedirs(OUTPUT_DIR, exist_ok=True)
|
os.makedirs(OUTPUT_DIR, exist_ok=True)
|
||||||
final_csv_output_path = OUTPUT_DIR / f"openairframes_adsb_{csv_start_date}_{csv_end_date}.csv.gz"
|
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")
|
df_latest_csv.write_csv(final_csv_output_path, compression="gzip")
|
||||||
else:
|
else:
|
||||||
print(f"Concatenating latest CSV (through {csv_end_date}) with new data ({args.date})")
|
print(f"Concatenating latest CSV (through {csv_end_date}) with new data ({args.date})")
|
||||||
# Ensure column order matches before concatenating
|
# 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
|
from src.adsb.compress_adsb_to_aircraft_data import concat_compressed_dfs
|
||||||
df_final = concat_compressed_dfs(df_latest_csv, df)
|
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"
|
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")
|
df_final.write_csv(final_csv_output_path, compression="gzip")
|
||||||
print(f"Final CSV written to {final_csv_output_path}")
|
print(f"Final CSV written to {final_csv_output_path}")
|
||||||
|
|||||||
Reference in New Issue
Block a user