Skip to content

2 Preprocesamiento

Javi edited this page May 6, 2025 · 3 revisions

Entrega I

Decodificación

Tras descargar los datos proporcionados en el Campus Virtual, el primer paso fue decodificar la información del dataframe inicial. De las dos columnas que había, tskafka indicaba el tiempo en milisegundos desde la época UNIX, que transformamos a formato estándar datetime para aumentar la legibilidad. La otra columna, message, contenía un mensaje codificado en base64, que pasamos a formato hexadecimal para después decodificarlo en función de los códigos proporcionados en el manual “The 1090 Megahertz Riddle”.

Apoyándonos en esta documentación, decodificamos los mensajes con funciones de pyModeS tras detectar el tipo de cada mensaje a partir del formato del enlace descendente. Había múltiples tipos de mensajes, cada uno con información distinta. Para recoger los distintos mensajes generamos un dataframe general con todas las posibles columnas, rellenando con nulos los valores que faltaban según el tipo de mensaje.

Para poder procesarlo después creamos una clase DataframeProcessor con distintos métodos estáticos que filtraban el conjunto de datos general según las necesidades de visualización, devolviendo dataframes más pequeños y específicos.

Anatomía de los mensajes

Los mensajes tienen pueden ser de cuatro tipos (identidad, altitud, ADS-B y Mode-S), codificado como ya hemos visto con el formato de enlace descendiente. Todos los mensajes tienen asociados el timestamp de emisión. Cada uno de estos tipos de mensaje proporcionan una información concreta, siendo común a todos ellos la ICAO, que es el identificador único de cada aeronave. Más concretamente:

  • Los de tipo identidad proporcionan el squack code, un código de cuatro dígitos utilizado en los transpondedores de las aeronaves para identificarlas y comunicarse con los controladores de tráfico aéreo.
  • Los de tipo altitud proporcionan la altitud barométrica de la aeronave.
  • Los de tipo ADS-B ofrecen información importante como la posición de la aeronave, la velocidad, el estado de vuelo o la categoría de turbulencia. Dentro de ellos hay más subtipos en función del campo typecode.
  • Los de tipo Mode-S proporcionan otra información relativa a la aeronave, incluyendo también datos meteorológicos. Dentro de ellos hay más subtipos en función del campo BDS.

Tanto los mensajes de tipo ADS-B como los de Mode-S pueden indicar el identificador único de vuelo, un campo llamado Callsign.

Para más detalles, ver el módulo decoder.py dentro de la carpeta preprocess.

Dificultades

Sin embargo, como cada tipo de mensaje proporcionaba una información incompleta, había ciertas filas que el código rellenaba de manera arbitraria y esto llevó a errores de coherencia en los datos. Para resolverlo, buscamos estas filas de outliers y las eliminamos del conjunto de datos, ya que la información de ese tipo de mensaje no era necesaria para esta entrega.

Uno de los principales problemas a los que nos hemos enfrentado en este proyecto ha sido el gran tamaño de los datos, tanto a la hora de procesar y decodificar el dataframe inicial como a la hora de representar las gráficas. Nuestros ordenadores no eran lo suficientemente potentes, así que utilizamos el ordenador de la facultad para decodificar los datos. Dividimos el conjunto en particiones para procesarlas, pero tuvimos problemas la primera vez que ejecutamos el código porque en cada partición saltaba un error distinto, y teníamos que modificar el código para acomodar cada caso. Había algunos datos que estaban almacenados como tuplas en vez de un solo dato, lo que hacía saltar el error al ejecutar el decodificador. Para solucionarlo dividimos el dataframe en más columnas, pero a medida que avanzaban las particiones iban saliendo nuevas tuplas de datos corruptos. La siguiente medida fue convertir esas columnas a objetos para evitar errores, y luego filtrar los datos corruptos y convertir la columna al tipo correcto después de decodificar todos los datos, manualmente. Además de esta conversión de tipos, los nulos se guardaron como tipo string, lo cual dificultaba el procesado y hacía que algunas columnas se guardaran como objeto, que también tuvimos que convertir al tipo correcto y marcar los nulos con su tipo correspondiente. Para facilitar todo este proceso creamos el módulo utilities.py con las funciones necesarias de conversión de tipos, procesado de nulos y creación de columnas auxiliares.

Estos problemas se agravaron especialmente al intentar procesar la semana completa debido a que la librería pandas no maneja bien grandes volúmenes de datos en términos de uso de memoria. Esto provoca que el entorno kernel se quede sin memoria disponible, lo que interrumpe la operación. Probamos usando la librería dask, pero hay métodos de pandas como merge_asof que no son compatibles con ella. Esto nos dificultó el proceso ya que eran necesarias esas funcionalidades específicas a la hora de representar las gráficas. Como solución, decidimos separar los datos en 10 particiones, asegurándonos de que los mensajes del mismo ICAO estuvieran en la misma partición. Aplicamos las funciones pertinentes a cada particición por separado y luego juntamos los resultados en un único dataframe, que era de menor tamaño y ya se podía cargar y representar sin sobrepasar la memoria del ordenador.

A la hora de representar los datos era necesario extraer la hora y el día de la semana a partir de la fecha, también con funciones de utilities.py.

Posiciones válidas

Creamos un módulo data_processor.py para determinar si la posición de un avión es válida o no. Según el manual, las posiciones dadas se pueden considerar no válidas si se cumple uno de los siguientes casos:

  • Si el avión se encuentra fuera del rango máximo del radar. Este rango no es fijo, pues depende, entre otras cosas, de la altitud del avión.
  • Si la distancia entre el avión y el radar es mayor de 180 millas náuticas (equivalen a unos 1800 km).

Cálculo de los tiempos de espera

Los mensajes asociados a los Callsigns (un identificador único de vuelo) no mostraban el estado de vuelo (flight status) como airborne, lo que complicó la correcta identificación de los vuelos activos. Para resolverlo, utilizamos la función merge_asof, combinando dos dataframes complementarios utilizando el ICAO como clave y relacionándolos por proximidad de tiempo.

Después, seleccionamos para cada vuelo el momento en el que despegaba y el momento en el que estaba en pista con velocidad cero y los restamos para obtener el tiempo de espera, que posteriormente convertimos a segundos.

Es preciso indicar que las funciones de ciertas librerías fallaron al intentar procesar el gran volumen de datos (> 3 GB). Para solucionarlo, implementamos un procesamiento por particiones, guardando los resultados de forma incremental en archivos Parquet separados, evitando así cargar todo el conjunto de datos en memoria al mismo tiempo. Un pequeño detalle a tener en cuenta es que todas las señales procedentes de una misma aeronave (misma ICAO) se destinaron a la misma partición para que, de este modo, los vuelos (identificados por el Callsign) también estuviesen en la misma. De este modo, no habría errores a la hora de calcular los tiempos de espera.

Entrega II

Transición a PySpark

En la primera entrega utilizamos la librería pandas para procesar y limpiar todos los datos. Sin embargo, no era eficiente ya que la cantidad de información era muy elevada y los portátiles no tenían suficiente memoria para ejecutarlo. Además, la propia librería también ponía límites al tamaño de los datos que se pueden pasar como argumento a sus funciones.

La solución a esto fue refactorizar nuestros módulos existentes de procesado de datos a la librería PySpark. Con esto sacamos ventaja del procesamiento distribuido y conseguimos evitar los problemas por el tamaño de los datos. También fue útil a la hora de trabajar en el clúster, ya que toda la ejecución del código era mucho más eficiente al explotar el paralelismo proporcionado por la librería usando los múltiples nodos accesibles en el sistema de Cloudera.

Lectura y decodificación

Como primera etapa, esta es la encargada de leer los datos crudos (csv), aplicar nuestro decodificador, y guardar los datos en formato parquet. El desafío estaba en que los datos se encuentran en el sistema de datos distribuidos HDFS de Cloudera. En un principio es una gran herramienta, ya que aprovechamos las ventajas de un sistema distribuido y la FDI nos otorgó tres nodos en los que podemos distribuir nuestra carga computacional de manera eficiente, sin embargo, nos costó un tiempo realizar la transición a PySpark (explicado anteriormente).

Tras realizar la transición, ya todo fue mucho más fluido y notamos la gran eficiencia de estas herramientas en comparación a la práctica anterior donde empleamos sistemas no distribuidos sobre entornos locales.

Entonces, quitando las limitaciones de disco y runtime que nos saltaban en determinadas ocasiones, logramos leer, decodificar y guardar los datos en estos simples pasos:

  • Lectura de datos: se leen archivos .csv codificados en base64 desde rutas por fecha usando Spark con sep=";" y header=True.
  • Decodificación de mensajes: se aplica una UDF decode_message() que transforma el mensaje base64 en un diccionario con atributos relevantes.
  • Transformación a filas: se convierte la columna decoded (tipo dict) en filas tipo Row usando dict_to_row(), seleccionando solo las columnas necesarias.
  • Creación del DataFrame estructurado: se usa un esquema con todos los campos como StringType() para cargar datos evitando errores de tipo.
  • Limpieza de valores nulos: se reemplazan null, "None" y valores ausentes por "NaN" con when(), .replace() y .na.fill().
  • Conversión de tipos: se hace cast explícito de las columnas necesarias a su tipo correcto (timestamp, int, float, etc.).
  • Escritura optimizada: se guarda en formato Parquet con .repartition(10), .option("dfs.replication", 1) y .mode("overwrite") en una ruta por día.
  • Ejecución diaria: se recorre un rango de fechas con un while, procesando los datos día por día de forma independiente.

El tiempo de lectura y decodificación por día es de aproximadamente 40 minutos en el clúster de Cloudera.

Estructuración y enriquecimiento de datos

Esta etapa es quizá la más costosa a nivel computacional. Por ello, se han pensado meticulosamente las fases del pipeline de tal forma que se trate de hacer de la forma más eficiente posible. En concreto, hemos tratado de retrasar al máximo las operaciones de join para de esta forma evitar la explosión combinatoria de casos.

  • Separación del resultado de la etapa anterior en varios dataframes de acuerdo a la tipología de los mensajes: posiciones, identificadores de vuelo, velocidades y categorías de turbulencia.
  • Filtrado para quedarnos con las posiciones cercanas al aeropuerto en un radio de 5 kilómetros (se eliminan aproximadamente un 60% de los mensajes).
  • Combinación de los dataframes de posiciones e identificadores de vuelo.
  • Detección y asignación de aeronaves situadas en puntos de espera, así como la pista correspondiente.
  • Combinación del resultado con las categorías de turbulencia.
    dataframe A
  • Eliminación de aquellos vuelos que no se les ha detectado en un punto de espera.
  • Selección de las filas con posiciones pertenecientes al intervalo desde que se detecta a la aeronave en un punto de espera hasta la primera vez que se le detecta en el aire y, por lo tanto, ha despegado.
  • Combinación del resultado con los mensajes de velocidades (por proximidad de timestamps de 1 segundo).
  • Filtrado para quedarnos solo con las aeronaves que en algún momento se paran en un punto de espera (velocidad cero).
  • Calculamos los tiempos de espera para todos los mensajes que están en algún punto de espera:
    • Tiempo desde la primera vez que se identifica el vuelo.
    • Tiempo que lleva situado en el punto de espera.
    • Tiempo que hasta el despegue (variable objetivo).
      dataframe B
  • Con A, creamos otro dataframe que indique cada 10 segundos qué pistas y puntos de espera están ocupados.
    dataframe C
  • Con A, creamos otro dataframe que indique todos los despegues y aterrizajes por pista, incluyendo la categoría de turbulencia del avión.
    dataframe D
  • Agrupamos D por minuto para conseguir tasa de despegues y de aterrizajes por minuto.
    dataframe E
  • Combinamos B con C, D y E para obtener:
    • Pistas y puntos de espera ocupados.
    • Tiempo desde el último evento (aterrizaje / despegue) en esa pista.
    • Categoría de turbulencia de esa aeronave.
    • Tasas de aterrizaje y despegue en el último minuto.

Es preciso indicar que los timestamps de los dataframes C y E se han redondeado hacia arriba, mientras que los del dataframe B se han redondeado hacia abajo, de esta forma evitamos data leakage.

dataframe F

  • En F, extraemos la hora, día de la semana y si es festivo o no.
  • Extraemos la aerolínea (primeros tres dígitos del identificador de vuelo).
  • Renombramos variables y reorganizamos las columnas.
  • Combinamos con los datos meteorológicos.

Depuración y selección de variables

Esta última fase del preprocesamiento tiene lugar tras la exploración de variables (ver siguiente sección), aplicando las conclusiones extraídas de la misma:

  • Limpiamos las variables turbulence_category y last_event_turb_cat eliminando lo que hay entre paréntesis (de cara a las gráficas).
  • Eliminamos aquellos outliers de la variable respuesta que consideramos que eran errores de medición / físicamente imposibles:
    • Tiempos de espera inferiores a los 20 segundos: el runway time occupancy (ROT) es en promedio de 30 segundos, es imposible que aviones despeguen tan rápido.
    • Tiempos de espera superiores a los 1.250 segundos: procedentes de aeronaves con ICAO nula ‘______’ o que consideramos que pasaron por el punto pero a continuación se cambiaron de posición.
    • Se eliminaron 559 filas. 16 vuelos distintos. Un 0.36% del total.

img

- Eliminamos variables por las siguientes razones: - **Presencia de nulos**: `snow_depth (m)` - **Ausencia de varianza**: `last_event` - **Evitar colinealidad**: - `lat`, `lon`, `apparent_temperature (°C)`, `rain (mm)`, `pressure_msl (hPa)`, `wind_speed_100m (km/h)`, `wind_gusts_10m (km/h)`, `soil_temperature_0_to_7cm (°C)`, `soil_temperature_7_to_28cm (°C)`, `soil_temperature_28_to_100cm (°C)`, `soil_moisture_7_to_28cm (m³/m³)`, `soil_moisture_28_to_100cm (m³/m³)`, `vapour_pressure_deficit (kPa)` - **Innecesarias**: `event_timestamp`, `first_holding_time`, `first_airborne_time`, `first_on_ground_time`, `time`

Las variables timestamp, icao y callsign las dejamos para poder mantener un registro de la procedencia de los datos.

Hemos decidido no realizar ningún preprocesado adicional en esta fase. Cualquier transformación adicional a los datos será propia de cada uno de los modelos.

Clone this wiki locally