-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtransform.py
More file actions
161 lines (134 loc) · 6.69 KB
/
Copy pathtransform.py
File metadata and controls
161 lines (134 loc) · 6.69 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
"""
Pure PySpark transform for the flight ETL workshop.
This module contains the transformation logic without any awsglue dependencies.
Both the AWS Glue entry point (main.py) and the local runner (main_local.py)
import from this module so the same transforms run identically in both contexts.
"""
from __future__ import annotations
from pyspark.sql import DataFrame
from pyspark.sql.functions import col, count, concat_ws, mean
from pyspark.sql.window import Window
# Column type spec for the raw flights CSV — used to coerce types after read.
FLIGHT_COLUMN_TYPES: list[tuple[str, str]] = [
("id", "int"),
("year", "int"),
("month", "int"),
("day", "int"),
("dep_time", "double"),
("sched_dep_time", "int"),
("dep_delay", "double"),
("arr_time", "double"),
("sched_arr_time", "int"),
("arr_delay", "double"),
("carrier", "string"),
("flight", "int"),
("tailnum", "string"),
("origin", "string"),
("dest", "string"),
("air_time", "double"),
("distance", "int"),
("hour", "int"),
("minute", "int"),
("time_hour", "timestamp"),
("name", "string"),
]
def coerce_types(df: DataFrame) -> DataFrame:
"""Cast every column to the FLIGHT_COLUMN_TYPES spec, preserving column order."""
for column, dtype in FLIGHT_COLUMN_TYPES:
if column in df.columns:
df = df.withColumn(column, col(column).cast(dtype))
return df
def fill_nulls(df: DataFrame) -> DataFrame:
"""Drop string-null rows and fill numeric nulls with column means."""
null_columns: list[str] = [c for c in df.columns if df.filter(col(c).isNull()).count() > 0]
string_cols = [
f.name
for f in df.schema.fields
if f.dataType.simpleString() == "string" and f.name in null_columns
]
numeric_cols = [
f.name
for f in df.schema.fields
if f.dataType.simpleString() in ("int", "double") and f.name in null_columns
]
if string_cols:
df = df.dropna(subset=string_cols)
if numeric_cols:
means_row = df.select(*[mean(col(c)).alias(c) for c in numeric_cols]).collect()[0]
df = df.fillna({c: means_row[c] for c in numeric_cols if means_row[c] is not None})
return df
def add_derived_columns(df: DataFrame) -> DataFrame:
"""Add total_delay and route columns used by the downstream aggregations."""
df = df.withColumn("total_delay", col("arr_delay") + col("dep_delay"))
df = df.withColumn("route", concat_ws("-", col("origin"), col("dest")))
return df
def add_window_features(df: DataFrame) -> DataFrame:
"""Add carrier_* and route_* window-function feature columns."""
carrier_year_month = Window.partitionBy("year", "month", "carrier")
carrier_year_hour = Window.partitionBy("year", "hour", "carrier")
carrier_year = Window.partitionBy("year", "carrier")
route_year_month = Window.partitionBy("year", "month", "route")
route_year_hour = Window.partitionBy("year", "hour", "route")
route_year = Window.partitionBy("year", "route")
carrier_windows = {
"year_hour": carrier_year_hour,
"year_month": carrier_year_month,
"year_avg": carrier_year,
}
route_windows = {
"year_hour": route_year_hour,
"year_month": route_year_month,
"year_avg": route_year,
}
delay_columns = ["arr_delay", "dep_delay", "total_delay"]
count_columns = ["flight"]
for column in delay_columns:
for partition, window in carrier_windows.items():
df = df.withColumn(f"carrier_{partition}_avg_{column}", mean(column).over(window))
for partition, window in route_windows.items():
df = df.withColumn(f"route_{partition}_avg_{column}", mean(column).over(window))
for column in count_columns:
for partition, window in carrier_windows.items():
df = df.withColumn(f"carrier_{partition}_total_{column}", count(column).over(window))
for partition, window in route_windows.items():
df = df.withColumn(f"route_{partition}_total_{column}", count(column).over(window))
return df
def carrier_analysis(df: DataFrame) -> DataFrame:
"""Aggregate by year/month/hour/carrier with all delay + flight metrics."""
return df.groupBy("year", "month", "hour", "carrier").agg(
mean("arr_delay").alias("carrier_year_month_hour_avg_arr_delay"),
mean("carrier_year_month_avg_arr_delay").alias("carrier_year_month_avg_arr_delay"),
mean("carrier_year_hour_avg_arr_delay").alias("carrier_year_hour_avg_arr_delay"),
mean("dep_delay").alias("carrier_year_month_hour_avg_dep_delay"),
mean("carrier_year_month_avg_dep_delay").alias("carrier_year_month_avg_dep_delay"),
mean("carrier_year_hour_avg_dep_delay").alias("carrier_year_hour_avg_dep_delay"),
mean("total_delay").alias("carrier_year_month_hour_avg_total_delay"),
mean("carrier_year_month_avg_total_delay").alias("carrier_year_month_avg_total_delay"),
mean("carrier_year_hour_avg_total_delay").alias("carrier_year_hour_avg_total_delay"),
count("flight").alias("carrier_year_month_hour_total_flights"),
mean("carrier_year_month_total_flight").alias("carrier_year_month_total_flights"),
mean("carrier_year_hour_total_flight").alias("carrier_year_hour_total_flights"),
)
def route_analysis(df: DataFrame) -> DataFrame:
"""Aggregate by year/month/hour/route with all delay + flight metrics."""
return df.groupBy("year", "month", "hour", "route").agg(
mean("arr_delay").alias("route_year_month_hour_avg_arr_delay"),
mean("route_year_month_avg_arr_delay").alias("route_year_month_avg_arr_delay"),
mean("route_year_hour_avg_arr_delay").alias("route_year_hour_avg_arr_delay"),
mean("dep_delay").alias("route_year_month_hour_avg_dep_delay"),
mean("route_year_month_avg_dep_delay").alias("route_year_month_avg_dep_delay"),
mean("route_year_hour_avg_dep_delay").alias("route_year_hour_avg_dep_delay"),
mean("total_delay").alias("route_year_month_hour_avg_total_delay"),
mean("route_year_month_avg_total_delay").alias("route_year_month_avg_total_delay"),
mean("route_year_hour_avg_total_delay").alias("route_year_hour_avg_total_delay"),
count("flight").alias("route_year_month_hour_total_flights"),
mean("route_year_month_total_flight").alias("route_year_month_total_flights"),
mean("route_year_hour_total_flight").alias("route_year_hour_total_flights"),
)
def transform(df: DataFrame) -> tuple[DataFrame, DataFrame]:
"""Run the full transformation chain. Returns (carrier_df, route_df)."""
df = coerce_types(df)
df = fill_nulls(df)
df = add_derived_columns(df)
df = add_window_features(df)
return carrier_analysis(df), route_analysis(df)