Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 03.06 · 14 мин
Средний
Spark SQLCatalogcreateOrReplaceTempViewspark.sqlHive MetastoreDatabaseTable

Spark SQL и Catalog API

Два пути к одному результату

Spark предоставляет два равноценных API для работы с данными: DataFrame API и Spark SQL. Оба проходят через один и тот же Catalyst optimizer и генерируют идентичный физический план выполнения:

# DataFrame API
result_df = (
    employees
    .filter(col("age") > 30)
    .join(departments, employees.dept_id == departments.id)
    .groupBy("dept_name")
    .agg(avg("salary").alias("avg_salary"))
)

# Spark SQL -- абсолютно тот же результат и план
result_sql = spark.sql("""
    SELECT d.dept_name, AVG(e.salary) as avg_salary
    FROM employees e
    JOIN departments d ON e.dept_id = d.id
    WHERE e.age > 30
    GROUP BY d.dept_name
""")

Оба запроса пройдут одинаковые стадии Catalyst: парсинг, анализ, логическая оптимизация (predicate pushdown, constant folding), физическое планирование. Результат — один и тот же DAG из stages и tasks.

TIP

Когда использовать SQL, а когда DataFrame API?

  • SQL: удобнее для аналитиков с SQL-опытом, лучше для сложных JOIN + подзапросов, легче читать при множественных агрегациях
  • DataFrame API: compile-time type checking (ошибки в названиях колонок видны при сборке), лучше для цепочки трансформаций, проще комбинировать с бизнес-логикой на Python/Scala
  • Совет: используйте то, что читабельнее для конкретного запроса. Можно свободно смешивать оба подхода в одном приложении.

spark.sql() и регистрация представлений

Чтобы использовать SQL-синтаксис, DataFrame нужно зарегистрировать как временное представление (temporary view):

# Регистрируем DataFrame как SQL-таблицу
employees.createOrReplaceTempView("employees")
departments.createOrReplaceTempView("departments")

# Теперь можно писать SQL
top_earners = spark.sql("""
    SELECT name, salary, department
    FROM employees
    WHERE salary > 70000
    ORDER BY salary DESC
""")
top_earners.show()

Методы регистрации представлений:

МетодОбласть видимостиКогда использовать
createOrReplaceTempView("name")Текущая SparkSessionСтандартный выбор для интерактивной работы
createTempView("name")Текущая SparkSessionЕсли хотите ошибку при дублировании имени
createOrReplaceGlobalTempView("name")Все SparkSession в приложенииМежсессионный обмен данными
createGlobalTempView("name")Все SparkSessionС проверкой дублирования

Глобальные представления доступны через схему global_temp:

df.createOrReplaceGlobalTempView("shared_data")
spark.sql("SELECT * FROM global_temp.shared_data")
Spark3.5

Метод registerTempTable() устарел с Spark 2.0. Используйте createOrReplaceTempView() — он явно указывает, что создаётся представление (view), а не физическая таблица.

Catalog API: метаданные под контролем

spark.catalog предоставляет программный доступ к метаданным — базам данных, таблицам, колонкам:

# Список баз данных
spark.catalog.listDatabases()

# Список таблиц в текущей базе
spark.catalog.listTables()

# Список колонок таблицы
spark.catalog.listColumns("employees")

# Проверить существование таблицы
spark.catalog.tableExists("employees")

# Текущая база данных
spark.catalog.currentDatabase()

# Сменить базу данных
spark.catalog.setCurrentDatabase("analytics")

Catalog API особенно полезен для:

  • Автоматизации ETL: проверка существования таблиц перед записью
  • Data discovery: перебор всех таблиц и их схем
  • Тестирования: верификация, что DDL-операции создали нужные объекты
Проверка знанийKnowledge check
Чем отличается createOrReplaceTempView() от createGlobalTempView()?
ОтветAnswer
createOrReplaceTempView() создаёт представление, видимое только в текущей SparkSession. createGlobalTempView() создаёт представление, доступное из всех SparkSession в приложении, через схему global_temp (например, SELECT * FROM global_temp.my_view). Для интерактивной работы в notebook обычно достаточно createOrReplaceTempView().

Managed vs External Tables

Spark поддерживает два типа постоянных таблиц (в отличие от временных представлений, которые живут только в сессии):

# Managed table -- Spark управляет и данными, и метаданными
df.write.saveAsTable("analytics.sales_summary")

# External table -- Spark управляет только метаданными
df.write.option("path", "/data/warehouse/sales/").saveAsTable("analytics.sales_raw")

# Или через SQL
spark.sql("""
    CREATE EXTERNAL TABLE sales_raw (
        id BIGINT, amount DOUBLE, sale_date DATE
    )
    STORED AS PARQUET
    LOCATION '/data/warehouse/sales/'
""")
ТипДанныеМетаданныеDROP TABLE
ManagedУправляет Spark (warehouse dir)В Hive MetastoreУдаляет данные и метаданные
ExternalНа внешнем storage (вы управляете)В Hive MetastoreУдаляет только метаданные, данные остаются

Практическое правило: используйте managed tables для промежуточных результатов, external tables для данных, которыми пользуются другие системы.

Hive Metastore Integration

Spark использует Hive Metastore для хранения метаданных о постоянных таблицах, базах данных, партициях. По умолчанию Spark запускает встроенный Derby-based metastore (файл metastore_db в рабочей директории).

В production обычно настраивают внешний Hive Metastore (MySQL/PostgreSQL):

spark = (
    SparkSession.builder
    .appName("production-etl")
    .config("spark.sql.warehouse.dir", "/data/warehouse/")
    .config("hive.metastore.uris", "thrift://metastore-host:9083")
    .enableHiveSupport()
    .getOrCreate()
)

Это позволяет нескольким Spark-приложениям и другим инструментам (Hive, Presto, Trino) видеть одни и те же таблицы.

Смешивание SQL и DataFrame API

Одно из мощных свойств Spark — возможность свободно переключаться между SQL и DataFrame API в одном пайплайне:

# Шаг 1: DataFrame API для ETL-логики
clean_data = (
    raw_data
    .filter(col("amount").isNotNull())
    .withColumn("amount_usd", col("amount") * col("exchange_rate"))
)

# Шаг 2: Регистрируем промежуточный результат как SQL view
clean_data.createOrReplaceTempView("clean_sales")

# Шаг 3: SQL для аналитического запроса
report = spark.sql("""
    SELECT
        region,
        COUNT(*) as total_sales,
        SUM(amount_usd) as total_revenue,
        AVG(amount_usd) as avg_deal_size
    FROM clean_sales
    GROUP BY region
    HAVING total_revenue > 100000
    ORDER BY total_revenue DESC
""")

# Шаг 4: Обратно к DataFrame API для записи
report.write.mode("overwrite").parquet("/reports/regional_sales/")

Результат spark.sql() — обычный DataFrame, с которым можно продолжить работу через DataFrame API.

Проверка знанийKnowledge check
Почему DataFrame API и Spark SQL генерируют одинаковый физический план выполнения?
ОтветAnswer
Оба API проходят через один и тот же Catalyst optimizer. DataFrame API строит логический план из вызовов методов (filter, join, groupBy), SQL-запрос парсится в логический план из текстового синтаксиса. После этого оба логических плана проходят одинаковые стадии оптимизации: analysis, logical optimization (predicate pushdown, constant folding), physical planning (выбор join-стратегии). Результат -- идентичный физический план и DAG из stages/tasks.

Spark 4.0 SQL: новые возможности

Spark4.0

В Spark 4.0 SQL-движок получил три крупных языковых расширения, которые меняют сам способ написания запросов: pipe syntax для функциональной композиции, native SQL UDFs и first-class string collations.

Pipe syntax (|>): функциональная композиция в SQL

Spark 4.0 ввёл pipe-оператор |>, который позволяет читать SQL-запрос сверху вниз — так же, как трансформации в DataFrame API. Каждый pipe-этап получает на вход результат предыдущего и применяет к нему один SQL-оператор.

-- Классический SQL: читать снаружи внутрь, FROM-WHERE-GROUP-HAVING-SELECT-ORDER
SELECT region, COUNT(*) AS sales_count, SUM(amount_usd) AS revenue
FROM clean_sales
WHERE amount_usd > 0
GROUP BY region
HAVING revenue > 100000
ORDER BY revenue DESC;

-- Pipe syntax: тот же запрос как линейная цепочка
FROM clean_sales
|> WHERE amount_usd > 0
|> AGGREGATE COUNT(*) AS sales_count, SUM(amount_usd) AS revenue
   GROUP BY region
|> WHERE revenue > 100000
|> ORDER BY revenue DESC;

Поддерживаемые операторы в pipe-цепочке: SELECT, EXTEND, DROP, SET, WHERE, AGGREGATE, JOIN, ORDER BY, LIMIT, UNION, TABLESAMPLE, PIVOT, UNPIVOT. Каждый pipe-шаг — независимая лексическая единица, и предыдущий результат адресуется через позицию в цепочке.

Главное преимущество — инкрементальное построение колонок:

-- EXTEND добавляет колонку, опираясь на предыдущие шаги
FROM orders
|> EXTEND amount * 0.1 AS commission
|> EXTEND commission + amount AS total_payable
|> SELECT order_id, total_payable
|> ORDER BY total_payable DESC;

Pipe syntax не заменяет классический SQL — они полностью совместимы и могут смешиваться в одном запросе. На уровне Catalyst оба синтаксиса парсятся в одинаковый logical plan, поэтому никакого performance-overhead нет.

Native SQL UDFs

До Spark 4.0 пользовательские функции на чистом SQL были невозможны: для повторного использования логики нужно было писать Scala/Python UDF, что несло overhead и ломало оптимизатор. Spark 4.0 добавил CREATE FUNCTION ... LANGUAGE SQL — native SQL-функции, которые встраиваются прямо в Catalyst.

-- Temporary SQL UDF
CREATE OR REPLACE TEMPORARY FUNCTION discount_amount(price DECIMAL(10,2), discount_pct INT)
RETURNS DECIMAL(10,2)
LANGUAGE SQL
RETURN price - (price * discount_pct / 100);

-- Permanent SQL UDF (сохраняется в catalog)
CREATE OR REPLACE FUNCTION analytics.tier_label(amount DECIMAL(10,2))
RETURNS STRING
LANGUAGE SQL
RETURN CASE
    WHEN amount < 100 THEN 'bronze'
    WHEN amount < 500 THEN 'silver'
    WHEN amount < 2000 THEN 'gold'
    ELSE 'platinum'
END;

-- Использование
SELECT order_id,
       discount_amount(price, 15) AS price_after,
       analytics.tier_label(price) AS tier
FROM orders;

Ключевое отличие от Python/Scala UDF: SQL UDF inline в logical plan при analysis-фазе. Catalyst видит тело функции как обычное выражение, применяет predicate pushdown, constant folding и whole-stage codegen. Никакой serialization-overhead Python-worker’а или JVM-callback не возникает.

String collations: language-aware сравнение

В Spark 4.0 STRING-тип получил атрибут collation — набор правил сравнения и сортировки, специфичный для locale или режима (case/accent insensitive).

-- Case-insensitive сравнение через UTF8_BINARY_LCASE
SELECT * FROM users
WHERE email = '[email protected]' COLLATE UTF8_BINARY_LCASE;

-- Unicode case-insensitive + accent-insensitive (для multi-lingual)
SELECT name FROM customers
WHERE name = 'andre' COLLATE UNICODE_CI_AI;
-- Найдёт 'André', 'ANDRE', 'andré' одной операцией

-- Collation на уровне колонки (вся таблица)
CREATE TABLE products (
    sku    STRING,
    name   STRING COLLATE UNICODE_CI,  -- case-insensitive
    region STRING COLLATE UTF8_BINARY  -- byte-exact (default)
);

-- Group by с правильным locale-aware ordering
SELECT region, COUNT(*) FROM customers
GROUP BY region COLLATE UNICODE_CI_AI;

Поддерживаемые collations: UTF8_BINARY (default, byte-exact), UTF8_BINARY_LCASE (case-insensitive ASCII), UNICODE (full Unicode-aware), UNICODE_CI (case-insensitive), UNICODE_AI (accent-insensitive), UNICODE_CI_AI (оба) — плюс locale-specific варианты (fr-FR, de-DE, tr-TR и т.д.).

INFO

Collation встроены в catalyst-операторы, а не реализуются через runtime-преобразование. При фильтрации с COLLATE UTF8_BINARY_LCASE Spark не делает LOWER(col) = LOWER(literal) — он использует специализированный операторный код, что критично для производительности на больших таблицах. Это первоклассная фича на уровне типа, а не runtime-помощник.

-- Сравнение производительности: collation vs UPPER()
-- Без collation (классический подход): два UPPER() per row
SELECT * FROM users WHERE UPPER(name) = UPPER('john');
-- С collation: специализированный native сравниватель
SELECT * FROM users WHERE name = 'john' COLLATE UTF8_BINARY_LCASE;

В типичных бенчмарках collation-based сравнение на 30-60% быстрее UPPER()/LOWER() обхода — разница накапливается при join’ах с string-ключами.

Как Catalyst превращает запрос в оптимизированный план — на уровне исходников — в senior-курсе:

Spark Internals: архитектура Catalyst

Что дальше?

В следующем уроке мы разберём UDF (User-Defined Functions) — как расширить Spark собственными функциями, когда встроенных недостаточно, и почему Python UDF значительно медленнее pandas UDF.

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 8