Files
mini-lakehouse-lab/notebooks/06_spark_and_trino_on_same_table.ipynb
ddadmin 65d57ec6d6 feat(module-06): реализован модуль по совместной работе Spark и Trino
- Зачем:
  - продемонстрировать возможность доступа нескольких движков к одной 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).
  - проверка наличия всех инструкций в документации.
2026-03-07 23:58:36 +03:00

26 KiB
Raw Permalink Blame History

Модуль 6. Одна таблица, два движка: Spark и Trino

В этом модуле мы увидим главное практическое следствие архитектуры Lakehouse: таблица, записанная одним движком, может быть прочитана другим без копирования данных.

Цели модуля:

  • записать таблицу через Spark и тут же прочитать её через Trino;
  • выполнить SQL-запросы к Iceberg-таблицам из Trino;
  • понять роль общего каталога и общего хранилища;
  • осознать отличие Lakehouse-модели (decoupled compute) от классического DWH.

Prerequisite: пройден Модуль 5 (silver-таблица существует).

Связка с привычным DWH (Greenplum)

Классический DWH (Greenplum) Lakehouse
Где лежат данные Внутри СУБД, на управляемых дисках В объектном хранилище (MinIO), отдельно от движков
Где лежат метаданные pg_catalog внутри той же СУБД Внешний JDBC-каталог (PostgreSQL) + metadata в MinIO
Кто может читать таблицу Только сама СУБД Любой движок с доступом к каталогу и хранилищу
Чтобы дать доступ другому инструменту Подключиться к СУБД или скопировать данные Подключить тот же каталог — данные уже доступны

Ключевая идея: в Lakehouse данные не заперты в одном движке. Это называется decoupled compute.

Как устроен этот ноутбук: Spark-код и Trino-запросы выполняются здесь. Для Trino используется Python-клиент trino. Если у вас установлен DBeaver (и подключен по инструкции из START_HERE.md) — каждый Trino-запрос можно выполнить и там (в markdown будут подсказки).

1. Spark-сессия, Trino-клиент и проверка таблиц

Инициализируем Spark и настраиваем helper для отправки запросов в Trino.

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"))

Если ячейка выше упала с ImportError: trino — нужно пересобрать образ Jupyter (в терминале: docker compose build jupyter && docker compose up -d).

Обе таблицы на месте. Spark может их читать. Trino-клиент готов. Вопрос: может ли Trino прочитать те же таблицы?

2. Читаем silver через Spark — фиксируем метрики

Сначала получим числа из Spark. Потом сравним с Trino.

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)

Запомни эти числа. Сейчас выполним тот же запрос через Trino.

3. Навигация по каталогу через Trino

DBeaver: тот же запрос можно выполнить в DBeaver, если он подключён к Trino по инструкции из START_HERE.md.

In [ ]:
trino_query("SHOW SCHEMAS FROM lakehouse")
In [ ]:
trino_query("SHOW TABLES FROM lakehouse.silver")
In [ ]:
trino_query("DESCRIBE lakehouse.silver.nyc_taxi_yellow")

Trino видит те же namespace (bronze, silver, default) и ту же таблицу с не менее чем 24 колонками. Различия в нотации типов (varchar vs STRING, double vs DOUBLE) нормальны — это разница синтаксиса движков, а не данных.

4. «Aha-момент» — одни и те же данные

Центральный момент модуля. Trino читает таблицу, записанную Spark в Модуле 5.

DBeaver: SELECT * FROM lakehouse.silver.nyc_taxi_yellow LIMIT 10;

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)

Строки совпадают. Оба движка читают одни и те же data files из MinIO (с точностью до floating-point округления в функции avg).

5. Почему это работает — роль общего каталога

Общий каталог: оба движка подключены к одному инстансу PostgreSQL (postgres-iceberg:5432/iceberg). Когда Spark создаёт таблицу, он записывает ссылку на неё в PostgreSQL. Trino читает тот же PostgreSQL и находит эту ссылку.

Общее хранилище: data files (Parquet) и metadata-файлы лежат в MinIO (s3://lakehouse/warehouse). Оба движка имеют сетевой доступ к бакету и могут прочитать файлы.

Параллель с Greenplum: В классическом DWH данные лежат внутри СУБД на её собственных дисках. Доступ возможен только через саму СУБД. В Lakehouse данные лежат снаружи, а движки взаимозаменяемы.

Decoupled compute: Вычислительные ресурсы отделены от хранения. Если нам потребуется добавить третий движок (например, Apache Flink для streaming), нужно будет просто подключить его к тому же каталогу и хранилищу.

  ┌─────────────────┐     ┌──────────────────┐
  │  Spark (Jupyter)│     │  Trino (DBeaver) │
  └────────┬────────┘     └────────┬─────────┘
           │                       │
           ▼                       ▼
  ┌──────────────────────────────────────────┐
  │  PostgreSQL (JDBC catalog)               │
  │  metadata_location -> s3://...           │
  └────────────────────┬─────────────────────┘
                       ▼
  ┌──────────────────────────────────────────┐
  │  MinIO (S3-compatible storage)           │
  │  metadata/ + data/ (Parquet files)       │
  └──────────────────────────────────────────┘

6. Spark записывает — Trino читает (live write → read)

До этого мы читали таблицы, созданные в предыдущих модулях. Теперь — живой цикл: Spark создаёт таблицу прямо сейчас, и Trino её тут же видит.

Создаём агрегатную таблицу lakehouse.default.borough_summary через Spark:

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()

Spark записал таблицу. Мы не делали никакой синхронизации. Видит ли её Trino?

DBeaver: SELECT * FROM lakehouse.default.borough_summary ORDER BY trips DESC;

In [ ]:
trino_query("SELECT * FROM lakehouse.default.borough_summary ORDER BY trips DESC")

Trino видит таблицу мгновенно. Spark записал metadata в PostgreSQL и data files в MinIO. Trino прочитал тот же каталог — увидел таблицу. Никакого копирования, никакой синхронизации.

Это и есть decoupled compute: один движок создал, другой прочитал. В Greenplum для этого пришлось бы либо подключиться к тому же движку, либо скопировать данные.

7. Аналитические запросы в Trino

Trino как аналитический SQL-движок. Выполним несколько аналитических запросов.

DBeaver: скопируй запрос ниже без обёртки trino_query().

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)

Результаты совпадают (с точностью до floating-point). Два движка, одна таблица.

8. Навигация по bronze из Trino

До этого bronze читали только из Spark. Проверим через Trino.

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")

Trino видит все таблицы из всех namespace каталога. Каталог — общий, хранилище — общее.

DBeaver: те же запросы работают аналогично.

9. DBeaver — SQL-клиент для Trino (рекомендация)

В этом ноутбуке мы работали с Trino через Python-клиент — для воспроизводимости. В реальной аналитической работе Trino чаще используют через SQL-клиенты: DBeaver, DataGrip, DbVisualizer. DBeaver — де-факто стандарт для SQL, знакомый всем, кто работал с PostgreSQL/Greenplum.

Если DBeaver установлен и подключён (инструкция в START_HERE.md), попробуйте выполнить в нём любой запрос из этого модуля. Результат будет тот же — тот же протокол, тот же движок, та же таблица. Разница: DBeaver — для интерактивной SQL-работы, Python-клиент — для автоматизации и ноутбуков.

10. Самостоятельное задание

Четыре задачи. Центральная — самостоятельный цикл write → read, как в Секции 6.

Задача 1 (Spark → write). Создай через Spark новую агрегатную таблицу lakehouse.default.zone_summary с CTAS: для каждой pickup_zone посчитай count(*), avg(total_amount), avg(tip_amount) по silver-таблице (WHERE pickup_zone IS NOT NULL).

Задача 2 (Trino → read). Прочитай lakehouse.default.zone_summary через Trino (Python-клиент или DBeaver). Убедись, что Trino видит таблицу без синхронизации.

Задача 3 (сравнение). Выполни тот же агрегатный запрос напрямую по silver через Trino (без промежуточной таблицы). Сравни результат с zone_summary. Числа должны совпасть.

Задача 4 (ответ). Ответь в markdown-ячейке ниже: почему Trino увидел zone_summary сразу после создания через Spark? Какие три компонента стенда это обеспечивают? Что произошло бы, если бы у Trino был другой каталог?

In [ ]:
# Ваш код (Spark): CREATE TABLE zone_summary
In [ ]:
# Ваш код (Trino): чтение zone_summary
In [ ]:
# Ваш код (Trino): тот же агрегат напрямую по silver

Ваш ответ на Задачу 4: ...

In [ ]:
spark.sql("DROP TABLE IF EXISTS lakehouse.default.zone_summary")

Дополнительная задача (по желанию): аналитический запрос по bronze через Trino и через Spark, сравнение.

In [ ]:
# Ваш код (Trino): аналитический запрос по bronze
In [ ]:
# Ваш код (Spark): тот же запрос по bronze

11. Что мы НЕ сравниваем

Ограничение модуля: мы не сравниваем Spark и Trino как движки. Не обсуждаем: какой быстрее, какой лучше, каковы их внутренние различия оптимизаторов.

Цель — показать, что оба работают с одной таблицей через общий каталог. Это фундаментальное свойство архитектуры Lakehouse, а не конкретного движка.

12. Checkpoint

Проверь себя:

  1. Почему Trino видит таблицы, созданные Spark, без синхронизации?
  2. Какие два компонента стенда общие для Spark и Trino?
  3. Чем доступ к данным в Lakehouse отличается от Greenplum?
  4. Нужно ли копировать данные, чтобы Trino прочитал таблицу Spark?
  5. Что такое «decoupled compute» в одном предложении?
  6. Если добавить третий движок (Flink), что нужно сделать, чтобы он увидел те же таблицы?
  7. Что произошло, когда Spark создал borough_summary, а Trino её тут же прочитал?

13. Завершение

Удаляем демо-таблицу, она больше не нужна. Таблицы bronze и silver остаются для следующих модулей.

В Модуле 7 — schema evolution и time travel. В Модуле 8 — финальная практика.

In [ ]:
spark.sql("DROP TABLE IF EXISTS lakehouse.default.borough_summary")
spark.stop()