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.
Когда использовать 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-операции создали нужные объекты
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.
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 и т.д.).
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.