- Зачем: - продемонстрировать возможность доступа нескольких движков к одной Iceberg-таблице через общий каталог. - Что: - создан notebooks/06_spark_and_trino_on_same_table.ipynb с примерами write (Spark) -> read (Trino). - в jupyter/Dockerfile добавлен клиент trino. - в START_HERE.md и docs/stack_reference.md добавлена инструкция по подключению DBeaver к Trino. - статус модуля 6 в планах обновлен до Ready for validation. - Проверка: - ручная проверка выполнения ячеек в ноутбуке (Python trino client). - проверка наличия всех инструкций в документации.
26 KiB
26 KiB
In [ ]:
import pandas as pd
from pyspark.sql import SparkSession
from trino.dbapi import connect
import trino
CATALOG_NAME = "lakehouse"
BRONZE_TABLE = f"{CATALOG_NAME}.bronze.nyc_taxi_yellow"
SILVER_TABLE = f"{CATALOG_NAME}.silver.nyc_taxi_yellow"
spark = SparkSession.builder \
.appName("module-06-spark-and-trino") \
.getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
assert spark.catalog.tableExists(BRONZE_TABLE), f"Таблица {BRONZE_TABLE} не найдена (пройди Модуль 4)."
assert spark.catalog.tableExists(SILVER_TABLE), f"Таблица {SILVER_TABLE} не найдена (пройди Модуль 5)."
def trino_query(sql: str) -> pd.DataFrame:
"""Выполняет SQL-запрос в Trino и возвращает результат как Pandas DataFrame."""
conn = connect(
host="trino",
port=8080,
user="jupyter",
catalog="lakehouse"
)
cur = conn.cursor()
try:
cur.execute(sql)
rows = cur.fetchall()
columns = [desc[0] for desc in cur.description]
return pd.DataFrame(rows, columns=columns)
except trino.exceptions.ProgrammingError as e:
if "No nodes available to run query" in str(e):
print("Ошибка: Trino ещё не готов. Подождите пару минут после старта контейнеров.")
raise
# Запросы типа CREATE/DROP не возвращают строк
return pd.DataFrame([{"status": "Success"}])
finally:
conn.close()
print("Spark готов.")
print(f"Версия Trino клиента: {trino.__version__}")
print("Trino-клиент готов. Проверка...")
display(trino_query("SELECT 1 AS test_col"))In [ ]:
spark_metrics = spark.sql(f"""
SELECT
count(*) AS row_count,
avg(fare_amount) AS avg_fare,
avg(trip_distance) AS avg_distance,
count(DISTINCT pickup_borough) AS boroughs
FROM {SILVER_TABLE}
""").toPandas()
display(spark_metrics)
spark.table(SILVER_TABLE).select("tpep_pickup_datetime", "fare_amount", "pickup_borough").show(5)In [ ]:
trino_query("SHOW SCHEMAS FROM lakehouse")In [ ]:
trino_query("SHOW TABLES FROM lakehouse.silver")In [ ]:
trino_query("DESCRIBE lakehouse.silver.nyc_taxi_yellow")In [ ]:
trino_query("SELECT tpep_pickup_datetime, fare_amount, pickup_borough FROM lakehouse.silver.nyc_taxi_yellow LIMIT 5")In [ ]:
trino_metrics = trino_query("""
SELECT
count(*) AS row_count,
avg(fare_amount) AS avg_fare,
avg(trip_distance) AS avg_distance,
count(DISTINCT pickup_borough) AS boroughs
FROM lakehouse.silver.nyc_taxi_yellow
""")
print("Метрики Trino:")
display(trino_metrics)
print("Метрики Spark (для сравнения):")
display(spark_metrics)In [ ]:
spark.sql("""
CREATE OR REPLACE TABLE lakehouse.default.borough_summary
USING iceberg AS
SELECT pickup_borough,
count(*) AS trips,
avg(fare_amount) AS avg_fare,
avg(trip_distance) AS avg_distance
FROM lakehouse.silver.nyc_taxi_yellow
WHERE pickup_borough IS NOT NULL
GROUP BY pickup_borough
""")
spark.table("lakehouse.default.borough_summary").show()In [ ]:
trino_query("SELECT * FROM lakehouse.default.borough_summary ORDER BY trips DESC")In [ ]:
trino_query("""
SELECT
pickup_borough,
count(*) AS trips,
avg(fare_amount) AS avg_fare,
avg(trip_distance) AS avg_distance,
avg(tip_amount) AS avg_tip
FROM lakehouse.silver.nyc_taxi_yellow
WHERE pickup_borough IS NOT NULL
GROUP BY pickup_borough
ORDER BY trips DESC
""")In [ ]:
trino_top_zones = trino_query("""
SELECT
pickup_zone,
count(*) AS trips,
avg(total_amount) AS avg_total
FROM lakehouse.silver.nyc_taxi_yellow
WHERE pickup_zone IS NOT NULL
GROUP BY pickup_zone
ORDER BY trips DESC
LIMIT 10
""")
display(trino_top_zones)In [ ]:
spark_top_zones = spark.sql("""
SELECT
pickup_zone,
count(*) AS trips,
avg(total_amount) AS avg_total
FROM lakehouse.silver.nyc_taxi_yellow
WHERE pickup_zone IS NOT NULL
GROUP BY pickup_zone
ORDER BY trips DESC
LIMIT 10
""").toPandas()
display(spark_top_zones)In [ ]:
trino_query("SHOW TABLES FROM lakehouse.bronze")In [ ]:
trino_query("SELECT count(*) AS row_count FROM lakehouse.bronze.nyc_taxi_yellow")In [ ]:
trino_query("SELECT * FROM lakehouse.bronze.nyc_taxi_yellow LIMIT 5")In [ ]:
# Ваш код (Spark): CREATE TABLE zone_summary
In [ ]:
# Ваш код (Trino): чтение zone_summary
In [ ]:
# Ваш код (Trino): тот же агрегат напрямую по silver
In [ ]:
spark.sql("DROP TABLE IF EXISTS lakehouse.default.zone_summary")In [ ]:
# Ваш код (Trino): аналитический запрос по bronze
In [ ]:
# Ваш код (Spark): тот же запрос по bronze
In [ ]:
spark.sql("DROP TABLE IF EXISTS lakehouse.default.borough_summary")
spark.stop()