- Зачем: - демо-ячейки последующих модулей не должны зависеть от домашних заданий предыдущих — модуль 5 использует lookup в демо-коде для JOIN. - Что: - модуль 4: секция 8 из домашки (5 пустых ячеек) превращена в демо с заполненным кодом (чтение CSV, CTAS). - модуль 4: новая домашка (секция 9) — осмотр metadata lookup-таблицы, сравнение с основной. - модуль 4: секции перенумерованы (9→10, 10→11), завершение упоминает обе bronze-таблицы. - модуль 5: assert-сообщение для lookup обновлено (убрана ссылка на «самостоятельное задание»). - планы модулей 4 и 5: обновлены дизайн-решения, структура секций и раздел рисков. - Проверка: - прогнать готовые ячейки модуля 4 → lookup-таблица создаётся автоматически → модуль 5 проходит без ошибок. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
28 KiB
28 KiB
In [ ]:
import pandas as pd
from pyspark.sql import SparkSession, functions as F
CATALOG_NAME = "lakehouse"
BRONZE_NAMESPACE = f"{CATALOG_NAME}.bronze"
BRONZE_TABLE = f"{BRONZE_NAMESPACE}.nyc_taxi_yellow"
BRONZE_LOOKUP_TABLE = f"{BRONZE_NAMESPACE}.taxi_zone_lookup"
SILVER_NAMESPACE = f"{CATALOG_NAME}.silver"
SILVER_TABLE = f"{SILVER_NAMESPACE}.nyc_taxi_yellow"
spark = SparkSession.builder \
.appName("module-05-silver-layer") \
.getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
print("Spark готов")
print(f"Bronze table: {BRONZE_TABLE}")
print(f"Silver table: {SILVER_TABLE}")In [ ]:
assert spark.catalog.tableExists(BRONZE_TABLE), (
f"Таблица {BRONZE_TABLE} не найдена. Вернись в Модуль 4 и выполни загрузку в bronze."
)
assert spark.catalog.tableExists(BRONZE_LOOKUP_TABLE), (
f"Таблица {BRONZE_LOOKUP_TABLE} не найдена. "
"Вернись в Модуль 4 и выполни Секцию 8."
)
bronze_count = spark.table(BRONZE_TABLE).count()
assert bronze_count > 0, f"Таблица {BRONZE_TABLE} пуста."
print(f"Обе bronze-таблицы на месте. В основной таблице {bronze_count:,} строк. Теперь строим silver.")In [ ]:
bronze_df = spark.table(BRONZE_TABLE)
profile_df = bronze_df.select(
F.count("*").alias("total_rows"),
(F.sum(F.when(F.col("fare_amount") < 0, 1).otherwise(0)) / F.count("*") * 100).alias("negative_fare_pct"),
(F.sum(F.when(F.col("total_amount") > 1000, 1).otherwise(0)) / F.count("*") * 100).alias("extreme_total_pct"),
(F.sum(F.when((F.col("tpep_pickup_datetime") < "2024-01-01") | (F.col("tpep_pickup_datetime") >= "2025-01-01"), 1).otherwise(0)) / F.count("*") * 100).alias("out_of_range_date_pct"),
(F.sum(F.when(F.col("passenger_count").isNull(), 1).otherwise(0)) / F.count("*") * 100).alias("null_passengers_pct"),
(F.sum(F.when(F.col("RatecodeID").isNull(), 1).otherwise(0)) / F.count("*") * 100).alias("null_ratecode_pct"),
(F.sum(F.when(F.col("congestion_surcharge").isNull(), 1).otherwise(0)) / F.count("*") * 100).alias("null_congestion_pct"),
(F.sum(F.when(F.col("Airport_fee").isNull(), 1).otherwise(0)) / F.count("*") * 100).alias("null_airport_fee_pct")
)
profile_pd = profile_df.toPandas().T
profile_pd.columns = ["value"]
profile_pdIn [ ]:
bronze_df = spark.table(BRONZE_TABLE)
filtered_df = (
bronze_df
.filter(F.col("fare_amount") >= 0) # F1
.filter(F.col("total_amount") <= 1000) # F2
.filter((F.col("tpep_pickup_datetime") >= "2024-01-01") & (F.col("tpep_pickup_datetime") < "2025-01-01")) # F3
)
filtered_count = filtered_df.count()
removed_count = bronze_count - filtered_count
removed_pct = (removed_count / bronze_count) * 100
print(f"Строк до фильтрации: {bronze_count:,}")
print(f"Строк после фильтрации: {filtered_count:,}")
print(f"Удалено аномалий: {removed_count:,} ({removed_pct:.2f}%)")In [ ]:
normalized_df = (
filtered_df
.withColumn("passenger_count", F.coalesce(F.col("passenger_count").cast("int"), F.lit(0))) # N1
.withColumn("RatecodeID", F.coalesce(F.col("RatecodeID").cast("int"), F.lit(99))) # N2
.withColumn("congestion_surcharge", F.coalesce(F.col("congestion_surcharge"), F.lit(0.0))) # N3
.withColumn("Airport_fee", F.coalesce(F.col("Airport_fee"), F.lit(0.0))) # N4
)
normalized_df.select("passenger_count", "RatecodeID", "congestion_surcharge", "Airport_fee").show(5)In [ ]:
zone_lookup_df = spark.table(BRONZE_LOOKUP_TABLE)
enriched_df = (
normalized_df
.withColumn("trip_duration_minutes", (F.unix_timestamp("tpep_dropoff_datetime") - F.unix_timestamp("tpep_pickup_datetime")) / 60) # E1
.join(zone_lookup_df.alias("pu_zone"), F.col("PULocationID") == F.col("pu_zone.LocationID"), "left") # E2
.withColumnRenamed("Borough", "pickup_borough")
.withColumnRenamed("Zone", "pickup_zone")
.drop("LocationID", "service_zone")
.join(zone_lookup_df.alias("do_zone"), F.col("DOLocationID") == F.col("do_zone.LocationID"), "left") # E3
.withColumnRenamed("Borough", "dropoff_borough")
.withColumnRenamed("Zone", "dropoff_zone")
.drop("LocationID", "service_zone")
)
print(f"Количество колонок после обогащения: {len(enriched_df.columns)}")
enriched_df.select("tpep_pickup_datetime", "trip_duration_minutes", "pickup_zone", "dropoff_zone").show(5, truncate=False)In [ ]:
spark.sql(f"CREATE NAMESPACE IF NOT EXISTS {SILVER_NAMESPACE}")
enriched_df.createOrReplaceTempView("silver_ready")
ctas_sql = f"""
CREATE OR REPLACE TABLE {SILVER_TABLE}
USING iceberg
AS
SELECT
VendorID,
tpep_pickup_datetime,
tpep_dropoff_datetime,
passenger_count,
trip_distance,
RatecodeID,
store_and_fwd_flag,
PULocationID,
DOLocationID,
payment_type,
fare_amount,
extra,
mta_tax,
tip_amount,
tolls_amount,
improvement_surcharge,
total_amount,
congestion_surcharge,
Airport_fee,
trip_duration_minutes,
pickup_borough,
pickup_zone,
dropoff_borough,
dropoff_zone
FROM silver_ready
"""
print(f"Выполняем CTAS для {SILVER_TABLE}...")
spark.sql(ctas_sql)
silver_df = spark.table(SILVER_TABLE)
print(f"Silver-таблица готова. Строк: {silver_df.count():,}")
silver_df.printSchema()In [ ]:
assert silver_df.filter(F.col("fare_amount") < 0).count() == 0, "Найдены отрицательные тарифы!"
assert silver_df.filter(F.col("total_amount") > 1000).count() == 0, "Найдены экстремальные суммы!"
assert silver_df.filter(F.col("passenger_count").isNull()).count() == 0, "Найдены NULL в passenger_count!"
assert silver_df.filter(F.col("RatecodeID").isNull()).count() == 0, "Найдены NULL в RatecodeID!"
print("Row-level checks: OK")In [ ]:
expected_new_columns = ["trip_duration_minutes", "pickup_borough", "pickup_zone", "dropoff_borough", "dropoff_zone"]
for col in expected_new_columns:
assert col in silver_df.columns, f"Колонка {col} отсутствует в silver!"
types = dict(silver_df.dtypes)
assert types["passenger_count"] == "int", f"Тип passenger_count должен быть int, а не {types['passenger_count']}"
assert types["RatecodeID"] == "int", f"Тип RatecodeID должен быть int, а не {types['RatecodeID']}"
print("Schema-level checks: OK")In [ ]:
comparison = spark.sql(f"""
WITH stats AS (
SELECT
'bronze' as layer,
count(*) as rows,
avg(fare_amount) as avg_fare,
avg(trip_distance) as avg_dist
FROM {BRONZE_TABLE}
UNION ALL
SELECT
'silver' as layer,
count(*) as rows,
avg(fare_amount) as avg_fare,
avg(trip_distance) as avg_dist
FROM {SILVER_TABLE}
)
SELECT
*,
CASE
WHEN layer = 'silver' THEN
(1 - rows * 1.0 / (SELECT rows FROM stats WHERE layer = 'bronze')) * 100
ELSE 0
END as filtered_pct
FROM stats
""").toPandas()
comparisonIn [ ]:
# Ваш код: фильтрация F4
In [ ]:
# Ваш код: нормализация N5
In [ ]:
# Ваш код: обогащение E4
In [ ]:
# Ваш код: пересоздание silver-таблицы (CTAS)
In [ ]:
# Ваш код: новые проверки качества
In [ ]:
# Ваш код: сравнение количества строк до и после добавления F4
In [ ]:
spark.stop()