Catalog API: CREATE TABLE, ALTER TABLE, DESCRIBE DETAIL
Управление метаданными в Spark: Hive Metastore, managed vs external, полный жизненный цикл таблиц через DDL и programmatic Catalog API для production lakehouse.
Зачем нужен каталог¶
В самой простой архитектуре Data Lake аналитик работает с путями: spark.read.parquet("s3a://bucket/data/orders/2024/01/"). Это работает, но плохо масштабируется. При росте команды и числа датасетов возникают вопросы:
- Какие датасеты вообще существуют? Где их искать?
- Какая у них схема? Нужно ли читать файл чтобы узнать?
- Кто создал этот датасет? Когда? Для чего?
- Какого формата данные? Как они партиционированы?
- Можно ли на них делать SQL-запросы без знания физического пути?
Catalog (каталог) - ответ на все эти вопросы. Это слой метаданных поверх хранилища. Вместо spark.read.parquet("s3a://...") аналитик пишет SELECT * FROM gold.orders - и Spark сам знает где лежат файлы, в каком формате, как их прочитать.
Когда аналитик пишет SQL-запрос к gold.orders, Spark обращается к каталогу: находит физический путь (s3a://datalake/gold/orders/), формат (delta), схему, информацию о партициях. Затем читает нужные файлы. Аналитику не нужно знать ни пути, ни формат.
Что хранится в каталоге¶
Каталог хранит метаданные - не сами данные:
- Databases / Namespaces: логические группы таблиц (аналог схем в PostgreSQL)
- Tables: описание каждой таблицы: имя, схема, физическое расположение, формат, свойства
- Views: сохранённые SQL-запросы, которые выглядят как таблицы
- Partitions: информация о разделах партиционированных таблиц
- Statistics: статистика по таблицам и колонкам для оптимизатора
Архитектура каталога: Hive Metastore и его альтернативы¶
Hive Metastore как стандарт де-факто¶
Исторически Spark использует Hive Metastore (HMS) как хранилище метаданных. HMS - это сервис с собственной реляционной БД (обычно MySQL или PostgreSQL), где хранятся все метаданные таблиц.
Когда вы запускаете SparkSession.builder.enableHiveSupport().getOrCreate(), Spark подключается к HMS и получает доступ к уже зарегистрированным таблицам. Метаданные персистентны - перезапуск Spark не стирает информацию о таблицах.
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("catalog-demo") \
.config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
.enableHiveSupport() \
.getOrCreate()
# Проверяем текущий каталог
print(spark.catalog.currentDatabase()) # → default
print(spark.catalog.listDatabases()) # → все базы данных
Без enableHiveSupport() Spark использует in-memory session catalog - метаданные живут только в рамках SparkSession и исчезают при её завершении.
Современные каталоги¶
Hive Metastore - зрелое, но устаревающее решение. Современные lakehouse-платформы предлагают:
- AWS Glue Data Catalog - управляемый HMS-совместимый сервис от AWS. Spark может использовать его как каталог без самостоятельного развёртывания HMS
- Databricks Unity Catalog - централизованный каталог с fine-grained access control, data lineage и cross-workspace sharing
- Apache Polaris / Iceberg REST Catalog - открытый каталог для Iceberg-таблиц, реализующий Iceberg REST Catalog Specification
- Project Nessie - versioned catalog с Git-семантикой для Iceberg и Delta
В контексте этого урока мы работаем со стандартным Spark Catalog (совместимым с HMS), но все DDL-команды применимы и к современным каталогам.
Три уровня именования¶
Полное имя таблицы в Spark - трёхуровневое:
catalog_name.database_name.table_name
│ │ │
│ │ └── Имя таблицы: orders
│ └── База данных/схема: gold
└── Каталог: spark_catalog (default) или hive_metastore
В большинстве конфигураций catalog_name опускается и используется default (spark_catalog). Можно писать просто gold.orders.
Managed vs External таблицы¶
Это фундаментальное различие, которое нужно понимать до написания первой DDL-команды.
Managed (управляемые) таблицы¶
При создании managed-таблицы Spark берёт контроль над и метаданными, и данными. Данные хранятся в spark.sql.warehouse.dir - специальной директории, управляемой Spark:
-- Managed таблица
CREATE TABLE gold.orders (
order_id LONG,
customer_id LONG,
total_amount DOUBLE,
order_date DATE
) USING DELTA;
-- Данные будут в: <warehouse.dir>/gold.db/orders/
При DROP TABLE gold.orders:
- Метаданные удаляются из каталога ✓
- Физические файлы удаляются с диска ← данные утрачены безвозвратно
External (внешние) таблицы¶
При создании external-таблицы Spark регистрирует метаданные, но не управляет данными. Данные хранятся по пути, который вы указываете:
-- External таблица
CREATE TABLE gold.orders (
order_id LONG,
customer_id LONG,
total_amount DOUBLE,
order_date DATE
) USING DELTA
LOCATION 's3a://datalake/gold/orders/';
При DROP TABLE gold.orders:
- Метаданные удаляются из каталога ✓
- Физические файлы остаются ← данные сохранены
Когда что использовать¶
| Критерий | Managed | External |
|---|---|---|
| Данные создаются Spark-пайплайном | ✓ удобно | ✓ тоже можно |
| Данные уже существуют на S3/HDFS | ✗ нельзя | ✓ обязательно |
| Несколько инструментов читают данные | ✗ рискованно | ✓ безопасно |
| DROP TABLE удаляет данные | ✓ да | ✗ нет |
| Production lakehouse | ✗ не рекомендуется | ✓ предпочтительно |
В production почти всегда лучше external-таблицы. Причина: данные в lakehouse создаются одним инструментом, читаются другим. Если Spark-пайплайн перерегистрирует managed-таблицу, данные уничтожатся. External-таблица - просто "регистрация пути", данные живут независимо.
CREATE DATABASE¶
Прежде чем создавать таблицы, нужен namespace - база данных:
-- Простое создание
CREATE DATABASE IF NOT EXISTS silver;
-- С расположением (для external databases)
CREATE DATABASE IF NOT EXISTS gold
LOCATION 's3a://datalake/gold/'
COMMENT 'Gold layer: curated analytical datasets for BI and ML'
WITH DBPROPERTIES (
'owner' = 'data-platform-team',
'created_date' = '2024-01-01',
'sla' = '99.9%'
);
LOCATION для database определяет директорию по умолчанию для всех таблиц этой базы без явного LOCATION. При создании таблицы без LOCATION в managed-базе - данные идут в <db_location>/<table_name>/.
Просмотр всех баз данных:
SHOW DATABASES;
SHOW DATABASES LIKE 'gold*'; -- фильтрация по паттерну
DESCRIBE DATABASE EXTENDED gold;
-- Name: gold
-- Location: s3a://datalake/gold/
-- Comment: Gold layer: curated...
-- Properties: ((owner,data-platform-team), ...)
В PySpark:
# Переключиться на базу данных
spark.catalog.setCurrentDatabase("gold")
# Список баз данных
for db in spark.catalog.listDatabases():
print(f"{db.name}: {db.description} @ {db.locationUri}")
CREATE TABLE: полный синтаксис¶
Базовый CREATE TABLE с явной схемой¶
CREATE TABLE IF NOT EXISTS silver.order_items (
item_id BIGINT COMMENT 'Surrogate key',
order_id BIGINT NOT NULL COMMENT 'FK to orders',
sku STRING NOT NULL,
product_name STRING,
category STRING,
quantity INT NOT NULL,
unit_price DECIMAL(12,2),
discount_pct DECIMAL(5,4),
line_total DECIMAL(14,2) COMMENT 'qty * unit_price * (1 - discount_pct)',
order_date DATE NOT NULL COMMENT 'Partition column',
_ingested_at TIMESTAMP COMMENT 'ETL ingestion timestamp',
_source STRING COMMENT 'Source system identifier'
)
USING DELTA
PARTITIONED BY (order_date)
LOCATION 's3a://datalake/silver/order_items/'
COMMENT 'Silver: normalized order line items from Bronze orders JSON'
TBLPROPERTIES (
'delta.autoOptimize.optimizeWrite' = 'true',
'delta.autoOptimize.autoCompact' = 'true',
'pipeline.owner' = 'data-engineering',
'pipeline.version' = '2.1.0',
'pipeline.schedule' = '0 */6 * * *',
'data.retention_days' = '365',
'data.pii_columns' = 'none',
'docs.confluence' = 'https://confluence.example.com/display/DE/silver-order-items'
);
Разберём каждую часть:
IF NOT EXISTS - идемпотентность. Если таблица уже есть - DDL не падает с ошибкой. Обязательно в production-скриптах, где DDL выполняется при каждом деплое.
COMMENT на колонках - документация прямо в схеме. Это не просто "хорошая практика" - это метаданные, которые видны в Spark UI, Databricks Data Explorer, любом HMS-совместимом инструменте. Аналитик может сделать DESCRIBE silver.order_items и понять назначение каждой колонки.
USING DELTA - явное указание формата. Никогда не оставляйте формат по умолчанию: в разных конфигурациях Spark default может быть разным (Parquet, ORC, Hive SerDe). Явность - залог предсказуемости.
PARTITIONED BY (order_date) - физическое партиционирование. Spark создаст директории вида s3a://datalake/silver/order_items/order_date=2024-01-15/. Запрос WHERE order_date = '2024-01-15' прочитает только одну директорию.
LOCATION - путь к данным для external-таблицы. Данные могут уже существовать по этому пути - тогда Spark просто регистрирует их.
TBLPROPERTIES - произвольные key-value метаданные. В примере: настройки Delta Lake (delta.autoOptimize.*), информация о пайплайне (pipeline.*), политики данных (data.*), ссылки на документацию (docs.*).
Форматы таблиц: почему USING важен¶
-- Delta Lake: ACID, time travel, schema evolution
CREATE TABLE ... USING DELTA ...
-- Apache Iceberg: аналог Delta, vendor-neutral
CREATE TABLE ... USING ICEBERG ...
-- Чистый Parquet: нет ACID, нет time travel
CREATE TABLE ... USING PARQUET ...
-- ORC: исторически популярен в Hive-мире
CREATE TABLE ... USING ORC ...
-- Antиpattern: не указывать USING (Hive SerDe)
CREATE TABLE ... -- ← зависит от spark.sql.sources.default
В современном lakehouse выбор - Delta Lake или Iceberg. Оба предоставляют ACID-транзакции, schema evolution и time travel. Чистый Parquet - только если работаете со сторонними системами, не поддерживающими форматы lakehouse.
CREATE TABLE AS SELECT (CTAS)¶
CTAS создаёт таблицу и заполняет её данными за одну операцию. Очень удобно для создания производных таблиц:
CREATE TABLE gold.revenue_by_category
USING DELTA
LOCATION 's3a://datalake/gold/revenue_by_category/'
PARTITIONED BY (order_date)
TBLPROPERTIES ('pipeline.type' = 'aggregated')
AS
SELECT
order_date,
category,
SUM(line_total) AS total_revenue,
COUNT(DISTINCT order_id) AS order_count,
SUM(quantity) AS units_sold,
AVG(unit_price) AS avg_unit_price
FROM silver.order_items
WHERE order_date >= '2024-01-01'
GROUP BY order_date, category;
CTAS атомарен при использовании Delta: либо таблица создана полностью, либо не создана вообще. При ошибке в середине записи - транзакция откатывается.
Важный нюанс: схема CTAS-таблицы определяется автоматически из результата SELECT. Типы и nullable из запроса попадают в схему таблицы. Если хотите явные типы и комментарии - лучше сначала CREATE TABLE, потом INSERT INTO.
# То же самое в PySpark
revenue_df = spark.sql("""
SELECT order_date, category, SUM(line_total) AS total_revenue
FROM silver.order_items
GROUP BY order_date, category
""")
revenue_df.write \
.format("delta") \
.mode("overwrite") \
.option("overwriteSchema", "true") \
.partitionBy("order_date") \
.saveAsTable("gold.revenue_by_category")
Bucketing: детерминированная организация файлов¶
Bucketing - альтернатива партиционированию для ключей join. Вместо директорий по значению (как партиции) данные разбиваются на N файлов-корзин по хэшу ключа:
CREATE TABLE silver.users_bucketed
USING PARQUET
CLUSTERED BY (user_id) INTO 64 BUCKETS
LOCATION 's3a://datalake/silver/users_bucketed/'
AS SELECT * FROM silver.users;
Если обе стороны join забакетированы по одному ключу с одинаковым числом корзин - Spark может выполнить join без shuffle. Но bucketing требует тщательного планирования числа корзин: изменить его без полной перезаписи нельзя. В современных lakehouse-форматах (Delta, Iceberg) есть более гибкие альтернативы - ZORDER BY и CLUSTER BY.
Temporary Views: сессионный каталог¶
Temporary View - именованный SQL-запрос, доступный только в рамках текущей SparkSession. Не создаёт физических данных, не регистрируется в постоянном каталоге.
# DataFrame → Temporary View
df = spark.read.format("delta").load("s3a://datalake/silver/orders/")
df.createOrReplaceTempView("orders_silver")
# Теперь можно использовать в SQL
result = spark.sql("""
SELECT customer_id, COUNT(*) AS order_count
FROM orders_silver
GROUP BY customer_id
""")
# View живёт только пока жива SparkSession
# spark.stop() → view исчезает
Global Temporary View - доступна из любой SparkSession в рамках одного SparkContext. Находится в базе данных global_temp:
df.createOrReplaceGlobalTempView("shared_reference_data")
# Из другой сессии
spark2.sql("SELECT * FROM global_temp.shared_reference_data")
Global temp views полезны в многопоточных приложениях, где несколько сессий должны видеть общий вспомогательный датасет (справочники, конфиги).
ALTER TABLE: эволюция схемы без боли¶
Требования к данным меняются. Появляются новые колонки, переименовываются старые, обновляется документация. ALTER TABLE - инструмент для этих изменений.
ALTER TABLE ADD COLUMNS¶
Добавление новых колонок - безопасная аддитивная операция. В Delta Lake и Iceberg это операция над метаданными: физические файлы не перезаписываются. Новая колонка будет null во всех существующих строках.
-- Добавить одну колонку
ALTER TABLE silver.order_items
ADD COLUMN discount_code STRING
COMMENT 'Promo code applied at checkout, nullable';
-- Добавить несколько колонок сразу
ALTER TABLE silver.order_items
ADD COLUMNS (
campaign_id STRING COMMENT 'Marketing campaign identifier',
channel STRING COMMENT 'Acquisition channel: web/mobile/api',
is_gift BOOLEAN COMMENT 'Is this a gift order'
);
Мгновенная операция для Delta/Iceberg: обновляется только транзакционный лог (_delta_log/ или metadata/), файлы с данными не трогаются. При чтении старых файлов новые колонки будут возвращаться как null.
# Эквивалент в PySpark через SQL
spark.sql("""
ALTER TABLE silver.order_items
ADD COLUMN discount_code STRING COMMENT 'Promo code applied at checkout'
""")
# Или через DataFrame API с mergeSchema
df_with_new_col = df.withColumn("discount_code", lit(None).cast("string"))
df_with_new_col.write \
.format("delta") \
.mode("append") \
.option("mergeSchema", "true") \
.saveAsTable("silver.order_items")
ALTER TABLE RENAME COLUMN¶
-- Переименовать колонку (Delta Lake 2.0+ / Iceberg)
ALTER TABLE silver.order_items
RENAME COLUMN discount_pct TO discount_rate;
Переименование - операция метаданных. Физические файлы не изменяются. Delta Lake хранит mapping: при чтении старых файлов поле discount_pct маппится на discount_rate.
Важно: после переименования SQL-запросы должны использовать новое имя. Если сторонние системы обращаются к старому имени - это breaking change. Переименование требует синхронизации с командой и обновления downstream-кода.
ALTER TABLE CHANGE COLUMN TYPE¶
Изменение типа - потенциально опасная операция. Delta Lake разрешает только безопасные приведения типов (type widening):
-- Безопасно: расширение типа
ALTER TABLE silver.order_items
ALTER COLUMN quantity TYPE BIGINT;
-- INT → BIGINT: все существующие значения корректно конвертируются
-- Безопасно: от числа с меньшей точностью к большей
ALTER TABLE silver.order_items
ALTER COLUMN unit_price TYPE DECIMAL(14,2);
-- DECIMAL(12,2) → DECIMAL(14,2): расширение диапазона
-- НЕБЕЗОПАСНО (Delta заблокирует):
-- ALTER COLUMN unit_price TYPE DOUBLE; -- DECIMAL → DOUBLE потеря точности
-- ALTER COLUMN order_id TYPE INT; -- BIGINT → INT: потенциальное переполнение
-- ALTER COLUMN order_date TYPE STRING; -- DATE → STRING: разные семантики
При небезопасном изменении Delta выбросит ошибку:
AnalysisException: Cannot change column 'unit_price' from DECIMAL(12,2) to DOUBLE
Если всё же нужно изменить тип несовместимым образом - единственный путь: создать новую таблицу с нужной схемой и скопировать данные через CTAS.
ALTER TABLE SET TBLPROPERTIES¶
Обновление метаданных таблицы без изменения данных или схемы:
-- Обновить версию пайплайна и добавить новые свойства
ALTER TABLE silver.order_items
SET TBLPROPERTIES (
'pipeline.version' = '2.2.0',
'pipeline.last_updated' = '2024-03-01',
'data.quality_score' = '0.998',
'data.freshness_sla' = '1h'
);
-- Удалить свойство
ALTER TABLE silver.order_items
UNSET TBLPROPERTIES ('docs.confluence');
ALTER TABLE SET LOCATION¶
Перемещение external-таблицы на новый path без перемещения данных:
-- Переключить таблицу на новый S3-бакет (данные уже скопированы)
ALTER TABLE gold.revenue_by_category
SET LOCATION 's3a://new-datalake/gold/revenue_by_category/';
Данные физически не перемещаются - это чисто metadata-операция. Spark начнёт читать из нового пути. Полезно при миграции между bucket'ами.
DROP TABLE¶
-- External таблица: удалится только metadata, данные останутся
DROP TABLE IF EXISTS silver.order_items;
-- Файлы в s3a://datalake/silver/order_items/ остаются нетронутыми
-- Managed таблица: удалятся и metadata, и данные на диске!
DROP TABLE gold.temp_aggregation;
-- ← физические файлы УДАЛЕНЫ
-- TRUNCATE: удаляет данные, но оставляет metadata и схему
TRUNCATE TABLE silver.order_items;
-- Все строки удалены, таблица пуста, но существует в каталоге
DESCRIBE: инспекция метаданных¶
DESCRIBE TABLE - базовый обзор схемы¶
DESCRIBE TABLE silver.order_items;
+-------------------+-----------+--------------------------------------------------+
| col_name | data_type| comment|
+-------------------+-----------+--------------------------------------------------+
| item_id| bigint| Surrogate key|
| order_id| bigint| FK to orders |
| sku| string| |
| product_name| string| |
| category| string| |
| quantity| int| |
| unit_price|decimal(12,2)| |
| discount_pct|decimal(5,4)| |
| line_total|decimal(14,2)| qty * unit_price * (1 - discount_pct) |
| order_date| date| Partition column |
| _ingested_at| timestamp| ETL ingestion timestamp |
| _source| string| Source system identifier |
+-------------------+-----------+--------------------------------------------------+
Базовый DESCRIBE даёт схему: имена, типы, комментарии колонок. Этого достаточно для понимания структуры.
DESCRIBE EXTENDED / FORMATTED¶
DESCRIBE EXTENDED silver.order_items;
-- или
DESCRIBE FORMATTED silver.order_items;
Добавляет полную информацию о таблице помимо схемы:
# Detailed Table Information
Database: silver
Table: order_items
Owner: spark
Created Time: Mon Jan 15 10:00:00 UTC 2024
Last Access: UNKNOWN
Created By: Spark 3.5.0
Type: EXTERNAL
Provider: delta
Location: s3a://datalake/silver/order_items
Serde Library: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe
InputFormat: org.apache.hadoop.mapred.SequenceFileInputFormat
OutputFormat: org.apache.hadoop.hive.ql.io.HiveSequenceFileOutputFormat
# Partition Information
# col_name data_type comment
order_date date Partition column
# Table Properties
delta.autoOptimize.autoCompact true
delta.autoOptimize.optimizeWrite true
delta.minReaderVersion 1
delta.minWriterVersion 2
pipeline.owner data-engineering
pipeline.version 2.1.0
DESCRIBE EXTENDED или FORMATTED - это мощный инструмент аудита. Вы видите:
- Type: EXTERNAL - подтверждение что таблица внешняя (данные не удалятся при DROP)
- Provider: delta - формат хранения
- Location - точный физический путь на S3/MinIO
- Partition Information - какие колонки являются ключами партиционирования
- Table Properties - все TBLPROPERTIES, включая Delta-параметры и кастомные метаданные
DESCRIBE DETAIL (Delta Lake / Iceberg)¶
DESCRIBE DETAIL - специфичная для lakehouse-форматов команда. Возвращает детальную статистику о состоянии таблицы:
DESCRIBE DETAIL silver.order_items;
Результат (одна строка с множеством колонок):
+------+-----------------------------+------+------+----------+---------+----------+----------+
|format| location |nFiles|nParts|sizeInBytes|minReader|minWriter| properties|
+------+-----------------------------+------+------+----------+---------+----------+----------+
|delta |s3a://datalake/silver/... | 847 | 365 |4523456789| 1| 2| {delta.au..}|
+------+-----------------------------+------+------+----------+---------+----------+----------+
Ключевые поля:
- format:
delta- формат подтверждён - location: физический путь
- numFiles: 847 - количество физических Parquet-файлов (без файлов Delta-лога)
- numPartitions (или
sizeInBytesпартиций): 365 - количество уникальных партиций (по одной на каждый день за год) - sizeInBytes: 4 523 456 789 байт = ~4.2 GB
В PySpark эти данные доступны как DataFrame:
detail = spark.sql("DESCRIBE DETAIL silver.order_items")
detail.printSchema()
# Извлекаем ключевые метрики
row = detail.collect()[0]
print(f"Format: {row['format']}")
print(f"Location: {row['location']}")
print(f"Files: {row['numFiles']:,}")
print(f"Size: {row['sizeInBytes'] / 1024**3:.2f} GB")
print(f"Avg file: {row['sizeInBytes'] / row['numFiles'] / 1024**2:.1f} MB")
# Диагностика мелких файлов
avg_file_mb = row['sizeInBytes'] / row['numFiles'] / 1024**2
if avg_file_mb < 64:
print(f"⚠ Small file problem: avg file {avg_file_mb:.1f} MB < 64 MB")
print(f" Consider running OPTIMIZE on this table")
Диагностика Small File Problem через DESCRIBE DETAIL¶
Одна из главных причин медленных запросов - слишком много мелких файлов. S3 создаёт latency на каждый HTTP-запрос: много мелких файлов = много запросов = медленное чтение.
def audit_table_health(table_name: str) -> None:
"""Аудит здоровья Delta/Iceberg таблицы через DESCRIBE DETAIL."""
detail = spark.sql(f"DESCRIBE DETAIL {table_name}").collect()[0]
size_gb = detail['sizeInBytes'] / 1024**3
num_files = detail['numFiles']
avg_file_mb = detail['sizeInBytes'] / max(num_files, 1) / 1024**2
print(f"\n=== Table Health: {table_name} ===")
print(f" Format: {detail['format']}")
print(f" Location: {detail['location']}")
print(f" Total size: {size_gb:.2f} GB")
print(f" Num files: {num_files:,}")
print(f" Avg file size: {avg_file_mb:.1f} MB")
# Оценка здоровья файлов
if avg_file_mb >= 128:
file_health = "✓ Excellent (>128 MB avg)"
elif avg_file_mb >= 64:
file_health = "✓ Good (64-128 MB avg)"
elif avg_file_mb >= 16:
file_health = "⚠ Fair (16-64 MB avg) - consider OPTIMIZE"
else:
file_health = "✗ Poor (<16 MB avg) - run OPTIMIZE"
print(f" File health: {file_health}")
audit_table_health("silver.order_items")
=== Table Health: silver.order_items ===
Format: delta
Location: s3a://datalake/silver/order_items
Total size: 4.21 GB
Num files: 847
Avg file size: 5.1 MB
File health: ✗ Poor (<16 MB avg) - run OPTIMIZE
Средний размер файла 5.1 MB - типичный признак small file problem. После OPTIMIZE silver.order_items количество файлов сократится, средний размер вырастет до 64-128+ MB.
SHOW команды¶
-- Все базы данных
SHOW DATABASES;
SHOW DATABASES LIKE 'gold*';
-- Все таблицы в текущей или указанной БД
SHOW TABLES;
SHOW TABLES IN silver;
SHOW TABLES IN silver LIKE 'order*';
-- Партиции таблицы
SHOW PARTITIONS silver.order_items;
-- order_date=2024-01-01
-- order_date=2024-01-02
-- ...
-- Свойства таблицы
SHOW TBLPROPERTIES silver.order_items;
-- pipeline.owner data-engineering
-- pipeline.version 2.1.0
-- ...
-- Конкретное свойство
SHOW TBLPROPERTIES silver.order_items ('pipeline.version');
-- 2.1.0
-- Колонки (аналог DESCRIBE)
SHOW COLUMNS IN silver.order_items;
SparkSession.catalog: программный API¶
Помимо SQL, Spark предоставляет Python API для работы с каталогом. Это удобно для динамической генерации DDL, валидации в ETL-пайплайнах и автоматизации.
Навигация по каталогу¶
# Список баз данных
for db in spark.catalog.listDatabases():
print(f"Database: {db.name}")
print(f" Description: {db.description}")
print(f" Location: {db.locationUri}")
print()
# Список таблиц в базе данных
for table in spark.catalog.listTables("silver"):
print(f" {table.name} ({table.tableType}) - {table.description or 'no description'}")
# tableType: MANAGED, EXTERNAL, VIEW
# Список колонок
for col in spark.catalog.listColumns("silver", "order_items"):
nullable = "nullable" if col.nullable else "NOT NULL"
print(f" {col.name}: {col.dataType} [{nullable}]")
if col.description:
print(f" → {col.description}")
Проверка существования объектов¶
def ensure_database_exists(db_name: str, location: str, comment: str = "") -> None:
"""Создаёт базу данных если не существует."""
if not any(db.name == db_name for db in spark.catalog.listDatabases()):
spark.sql(f"""
CREATE DATABASE {db_name}
LOCATION '{location}'
COMMENT '{comment}'
""")
print(f"Created database: {db_name}")
else:
print(f"Database already exists: {db_name}")
def table_exists(db_name: str, table_name: str) -> bool:
"""Безопасная проверка существования таблицы."""
return spark.catalog.tableExists(f"{db_name}.{table_name}")
# Использование в ETL-пайплайне
ensure_database_exists("silver", "s3a://datalake/silver/", "Silver layer")
if not table_exists("silver", "order_items"):
spark.sql("""
CREATE TABLE silver.order_items (...)
USING DELTA
LOCATION 's3a://datalake/silver/order_items/'
""")
print("Created table: silver.order_items")
Refresh и repair: синхронизация метаданных¶
Иногда метаданные в каталоге расходятся с физическим состоянием хранилища:
-- REFRESH TABLE: сбрасывает кэш метаданных Spark
-- Нужно если файлы изменились вне Spark (например, другим инструментом)
REFRESH TABLE silver.order_items;
-- MSCK REPAIR TABLE: обнаруживает новые партиции в external Hive-таблицах
-- Нужно если данные добавлены напрямую в S3 без Spark
MSCK REPAIR TABLE silver.order_items;
Разница:
REFRESH TABLE- инвалидирует кэш Spark для конкретной таблицы. После этого следующий запрос прочитает актуальные метаданные из HMSMSCK REPAIR TABLE- сканирует физические пути на S3 и регистрирует в HMS новые партиции, которые были добавлены без Spark
Для Delta Lake MSCK REPAIR не нужен: Delta сам ведёт транзакционный лог и всегда знает актуальное состояние. Но REFRESH TABLE может понадобиться после операций из другого SparkSession.
# Программное обновление
spark.catalog.refreshTable("silver.order_items")
# Или через SQL
spark.sql("REFRESH TABLE silver.order_items")
Получение статистики через catalog¶
# Получить метаданные таблицы программно
table_meta = spark.catalog.getTable("silver", "order_items")
print(f"Name: {table_meta.name}")
print(f"Type: {table_meta.tableType}") # MANAGED или EXTERNAL
print(f"DB: {table_meta.database}")
# DESCRIBE DETAIL как DataFrame для программной обработки
detail_df = spark.sql("DESCRIBE DETAIL silver.order_items")
detail = detail_df.collect()[0].asDict()
print(f"Files: {detail.get('numFiles', 'N/A')}")
print(f"Size: {detail.get('sizeInBytes', 0) / 1024**3:.2f} GB")
print(f"Partitions:{detail.get('numPartitions', 'N/A')}")
Catalog и Catalyst: статистика для оптимизатора¶
Каталог участвует в оптимизации запросов: Catalyst-оптимизатор использует статистику о таблицах для выбора лучшего плана выполнения.
Почему статистика важна¶
Главное решение оптимизатора - выбор стратегии join:
- Broadcast Hash Join (BHJ): если одна из таблиц мала - отправить её целиком на все executors, избежать shuffle. Быстро
- Sort Merge Join (SMJ): обе стороны сортируются и мержируются. Требует shuffle, но работает для любых размеров
Если Catalyst не знает размер таблицы - он не может принять правильное решение. ANALYZE TABLE вычисляет и сохраняет статистику:
-- Базовая статистика: размер, количество строк
ANALYZE TABLE silver.order_items COMPUTE STATISTICS;
-- Расширенная: дополнительно гистограммы по каждой колонке
ANALYZE TABLE silver.order_items COMPUTE STATISTICS FOR ALL COLUMNS;
-- Только для конкретных колонок
ANALYZE TABLE silver.order_items COMPUTE STATISTICS FOR COLUMNS order_id, category;
После ANALYZE TABLE в DESCRIBE EXTENDED появится секция:
# Table Statistics
# Rows: 45,234,891
# Total size: 4,523,456,789 bytes
# Column Statistics:
# order_id: [min=1, max=45234891, distinct=45234891, null=0]
# category: [min=accessories, max=toys, distinct=23, null=42]
Catalyst видит что category имеет 23 уникальных значения - это помогает при cardinality estimation. Catalyst видит что orders (45 млн строк) гораздо больше order_items (1500 строк) - выбирает broadcast для малой таблицы.
Практика: жизненный цикл аналитической витрины¶
Разберём полный сценарий: создание новой витрины, наполнение данными, изменение требований, аудит состояния.
Шаг 1: Подготовка Medallion-структуры¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date, current_timestamp, lit, round as spark_round
spark = SparkSession.builder \
.appName("catalog-practice") \
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
.getOrCreate()
# Создаём базы данных для Bronze, Silver, Gold
for layer in ["bronze", "silver", "gold"]:
spark.sql(f"""
CREATE DATABASE IF NOT EXISTS {layer}
COMMENT '{layer.capitalize()} layer of Medallion architecture'
""")
spark.sql("SHOW DATABASES").show()
Шаг 2: Создание Silver-таблицы¶
CREATE TABLE IF NOT EXISTS silver.marketing_events (
event_id STRING NOT NULL COMMENT 'UUID from source system',
user_id BIGINT COMMENT 'User identifier, nullable for anonymous',
campaign_id STRING NOT NULL COMMENT 'Marketing campaign code',
channel STRING NOT NULL COMMENT 'Channel: email/push/sms/banner',
event_type STRING NOT NULL COMMENT 'click/open/conversion/unsubscribe',
revenue_usd DECIMAL(10,2) COMMENT 'Revenue attributed to this event',
event_date DATE NOT NULL COMMENT 'Partition key',
_ingested_at TIMESTAMP COMMENT 'ETL timestamp'
)
USING DELTA
PARTITIONED BY (event_date)
LOCATION 's3a://datalake/silver/marketing_events/'
COMMENT 'Silver: normalized marketing touchpoint events'
TBLPROPERTIES (
'pipeline.owner' = 'marketing-analytics',
'pipeline.version' = '1.0.0',
'pipeline.schedule' = '0 */4 * * *',
'data.retention_days' = '730',
'delta.autoOptimize.optimizeWrite' = 'true'
);
Шаг 3: Наполнение данными¶
from pyspark.sql.types import *
import pyspark.sql.functions as F
# Синтетические данные маркетинговых событий
import random
from datetime import date, timedelta
events_data = []
channels = ["email", "push", "sms", "banner"]
event_types = ["click", "open", "conversion", "unsubscribe"]
campaigns = ["SPRING24", "PROMO_Q1", "LOYALTY_GOLD", "REACTIVATION"]
for i in range(10000):
event_date = date(2024, 1, 1) + timedelta(days=random.randint(0, 89))
events_data.append((
f"evt_{i:08d}",
random.randint(1, 50000) if random.random() > 0.15 else None,
random.choice(campaigns),
random.choice(channels),
random.choice(event_types),
round(random.uniform(0, 500), 2) if random.random() > 0.7 else None,
event_date,
None, # _ingested_at
))
events_schema = StructType([
StructField("event_id", StringType()),
StructField("user_id", LongType()),
StructField("campaign_id", StringType()),
StructField("channel", StringType()),
StructField("event_type", StringType()),
StructField("revenue_usd", DecimalType(10, 2)),
StructField("event_date", DateType()),
StructField("_ingested_at", TimestampType()),
])
events_df = spark.createDataFrame(events_data, events_schema) \
.withColumn("_ingested_at", current_timestamp())
events_df.write \
.format("delta") \
.mode("append") \
.saveAsTable("silver.marketing_events")
print(f"Written {events_df.count()} rows")
Шаг 4: Аудит через DESCRIBE DETAIL¶
# Проверяем состояние таблицы после записи
detail = spark.sql("DESCRIBE DETAIL silver.marketing_events").collect()[0]
print("=== Table Audit ===")
print(f"Format: {detail['format']}")
print(f"Location: {detail['location']}")
print(f"Num files: {detail['numFiles']}")
print(f"Size: {detail['sizeInBytes'] / 1024:.1f} KB")
# Статистика партиций
spark.sql("SHOW PARTITIONS silver.marketing_events").show(10)
# Количество строк по партициям
spark.sql("""
SELECT event_date, COUNT(*) AS rows
FROM silver.marketing_events
GROUP BY event_date
ORDER BY event_date
LIMIT 5
""").show()
Шаг 5: Изменение требований - ALTER TABLE¶
Маркетинговый отдел просит добавить новые поля:
-- Добавляем колонку для промокода
ALTER TABLE silver.marketing_events
ADD COLUMN discount_code STRING
COMMENT 'Discount code applied, if any. Format: DISC_{TYPE}_{AMOUNT}';
-- Добавляем флаг A/B теста
ALTER TABLE silver.marketing_events
ADD COLUMN ab_variant STRING
COMMENT 'A/B test variant: A, B, or null if not in test';
-- Обновляем версию пайплайна
ALTER TABLE silver.marketing_events
SET TBLPROPERTIES (
'pipeline.version' = '1.1.0',
'pipeline.last_schema_change' = '2024-03-01',
'schema.changelog' = 'v1.1.0: added discount_code, ab_variant'
);
Шаг 6: Проверка эволюции¶
# Новые колонки видны в схеме
spark.sql("DESCRIBE TABLE silver.marketing_events").show(truncate=False)
# Старые строки возвращают null для новых колонок
spark.sql("""
SELECT event_id, campaign_id, discount_code, ab_variant
FROM silver.marketing_events
LIMIT 5
""").show()
# discount_code и ab_variant = null для старых строк - ОК
# Проверяем свойства
spark.sql("SHOW TBLPROPERTIES silver.marketing_events").show(truncate=False)
Шаг 7: Gold-агрегация с CTAS¶
CREATE TABLE IF NOT EXISTS gold.campaign_performance
USING DELTA
LOCATION 's3a://datalake/gold/campaign_performance/'
PARTITIONED BY (event_date)
COMMENT 'Gold: daily campaign KPIs for BI dashboards'
TBLPROPERTIES (
'pipeline.source' = 'silver.marketing_events',
'pipeline.type' = 'aggregated',
'data.refresh' = 'daily'
)
AS
SELECT
event_date,
campaign_id,
channel,
COUNT(*) AS total_events,
COUNT(DISTINCT user_id) AS unique_users,
SUM(CASE WHEN event_type = 'click' THEN 1 ELSE 0 END) AS clicks,
SUM(CASE WHEN event_type = 'conversion' THEN 1 ELSE 0 END) AS conversions,
SUM(CASE WHEN event_type = 'unsubscribe' THEN 1 ELSE 0 END) AS unsubscribes,
ROUND(SUM(revenue_usd), 2) AS total_revenue,
ROUND(SUM(revenue_usd) / NULLIF(COUNT(DISTINCT user_id), 0), 2) AS revenue_per_user
FROM silver.marketing_events
GROUP BY event_date, campaign_id, channel;
Best Practices для production¶
Правила именования¶
- Используйте
snake_caseдля баз данных, таблиц и колонок - Техничные колонки ETL-слоя (не несущие бизнес-смысл) называйте с префиксом
_:_ingested_at,_source,_batch_id - Имена партиционных колонок должны быть последними в схеме - они физически отделены от данных в Parquet
- TBLPROPERTIES: используйте единый формат ключей:
{namespace}.{key}-pipeline.owner,data.retention_days
Чеклист при создании таблицы¶
DDL_CHECKLIST = """
При создании каждой production-таблицы проверяем:
[ ] CREATE TABLE IF NOT EXISTS - идемпотентность
[ ] USING DELTA (или ICEBERG) - явный формат
[ ] LOCATION - external таблица, данные выживут DROP
[ ] COMMENT на таблице - что хранит, откуда данные
[ ] COMMENT на каждой колонке - что означает поле
[ ] PARTITIONED BY - если > 10 GB и есть временной или категорийный ключ
[ ] TBLPROPERTIES: pipeline.owner, pipeline.version, data.retention_days
[ ] NOT NULL на обязательных колонках (Delta / Iceberg поддерживают constraints)
"""
print(DDL_CHECKLIST)
Metadata governance в TBLPROPERTIES¶
-- Шаблон стандартных TBLPROPERTIES для production
TBLPROPERTIES (
-- Владение и описание
'pipeline.owner' = 'team-name',
'pipeline.slack_channel' = '#team-data-alerts',
-- Версионирование
'pipeline.version' = '1.0.0',
'pipeline.git_commit' = 'abc1234',
'pipeline.created_date' = '2024-01-01',
-- Расписание
'pipeline.schedule' = '0 */6 * * *',
'pipeline.sla_minutes' = '30',
-- Данные
'data.retention_days' = '365',
'data.pii_columns' = 'email,phone',
'data.classification' = 'internal',
-- Техническое
'delta.autoOptimize.optimizeWrite' = 'true',
'delta.autoOptimize.autoCompact' = 'true'
)
Автоматизация DDL в ETL-пайплайне¶
class TableManager:
"""Управляет жизненным циклом таблицы в Spark Catalog."""
def __init__(self, spark: SparkSession):
self.spark = spark
def create_if_not_exists(self, ddl: str) -> bool:
"""Выполняет CREATE TABLE IF NOT EXISTS и возвращает True если создана."""
before = set(t.name for t in self.spark.catalog.listTables())
self.spark.sql(ddl)
after = set(t.name for t in self.spark.catalog.listTables())
created = bool(after - before)
return created
def get_version(self, table: str) -> str:
"""Читает pipeline.version из TBLPROPERTIES."""
props = {
row["key"]: row["value"]
for row in self.spark.sql(f"SHOW TBLPROPERTIES {table}").collect()
}
return props.get("pipeline.version", "unknown")
def bump_version(self, table: str, new_version: str) -> None:
"""Обновляет версию пайплайна в метаданных."""
self.spark.sql(f"""
ALTER TABLE {table}
SET TBLPROPERTIES ('pipeline.version' = '{new_version}')
""")
print(f"Updated {table} version to {new_version}")
def audit(self, table: str) -> dict:
"""Возвращает сводку здоровья таблицы."""
detail = self.spark.sql(f"DESCRIBE DETAIL {table}").collect()[0].asDict()
num_files = detail.get("numFiles", 0)
size_bytes = detail.get("sizeInBytes", 0)
avg_mb = size_bytes / max(num_files, 1) / 1024**2
return {
"table": table,
"format": detail.get("format"),
"size_gb": round(size_bytes / 1024**3, 2),
"num_files": num_files,
"avg_file_mb": round(avg_mb, 1),
"needs_optimize": avg_mb < 64,
}
# Использование
mgr = TableManager(spark)
audit = mgr.audit("silver.marketing_events")
print(audit)
# {'table': 'silver.marketing_events', 'format': 'delta', 'size_gb': 0.05,
# 'num_files': 12, 'avg_file_mb': 4.3, 'needs_optimize': True}
Anti-patterns¶
Anti-pattern 1: Работать только с путями, игнорировать каталог¶
# НЕПРАВИЛЬНО: вся аналитика через spark.read.parquet(...)
df = spark.read.parquet("s3a://datalake/silver/events/2024/01/15/")
# Проблемы: схема неизвестна, нет документации, нельзя делать SQL,
# путь хардкожен, при переезде данных ломается код
# ПРАВИЛЬНО: зарегистрировать таблицу, работать через каталог
df = spark.table("silver.events")
# или
df = spark.sql("SELECT * FROM silver.events WHERE event_date = '2024-01-15'")
Anti-pattern 2: Managed-таблицы в production¶
# ОПАСНО: managed таблица (без LOCATION)
spark.sql("CREATE TABLE silver.orders (...) USING DELTA")
# DROP TABLE silver.orders → ДАННЫЕ УДАЛЕНЫ с диска!
# БЕЗОПАСНО: external таблица
spark.sql("""
CREATE TABLE silver.orders (...)
USING DELTA
LOCATION 's3a://datalake/silver/orders/'
""")
# DROP TABLE silver.orders → только metadata, данные живут
Anti-pattern 3: DDL без IF NOT EXISTS¶
# СЛОМАЕТСЯ при повторном запуске пайплайна:
spark.sql("CREATE TABLE silver.events (...)")
# AnalysisException: Table `silver`.`events` already exists
# ИДЕМПОТЕНТНО:
spark.sql("CREATE TABLE IF NOT EXISTS silver.events (...)")
Anti-pattern 4: Таблицы без TBLPROPERTIES и COMMENT¶
-- ПЛОХО: таблица-загадка
CREATE TABLE analytics.results USING PARQUET AS SELECT ...;
-- Что хранит? Кто создал? Откуда данные? Когда обновляется?
-- ХОРОШО: самодокументируемая таблица
CREATE TABLE analytics.results
USING DELTA
LOCATION '...'
COMMENT 'Daily campaign KPIs aggregated from silver.events. Source: ETL v2.1'
TBLPROPERTIES ('pipeline.owner' = 'analytics-team', 'pipeline.version' = '2.1.0')
AS SELECT ...;
Домашнее задание¶
Написать DDL-скрипт для Medallion-архитектуры интернет-магазина:
-
Создать три базы данных:
bronze_shop,silver_shop,gold_shop -
Создать таблицы:
bronze_shop.raw_orders- внешняя Delta-таблица, сырые заказы из Kafka (order_id, payload STRING, received_at, _source)silver_shop.orders- нормализованные заказы (order_id, customer_id, status, total_amount, order_date) с партиционированием поorder_datesilver_shop.order_items- позиции заказа (item_id, order_id, sku, qty, price, order_date)-
gold_shop.daily_revenue- CTAS-витрина: выручка по категориям и дням -
Для каждой таблицы указать COMMENT и TBLPROPERTIES:
pipeline.owner,pipeline.version,data.retention_days -
Написать блок ALTER TABLE:
- Добавить колонку
promo_code STRINGвsilver_shop.orders - Переименовать
total_amount→order_totalвsilver_shop.orders -
Обновить
pipeline.versionдо1.1.0во всех трёх silver-таблицах -
Выполнить
DESCRIBE DETAILдля каждой таблицы и написать текстовый аудит: формат, количество файлов, средний размер файла, нужна ли оптимизация (OPTIMIZE)
Чеклист¶
- [ ] Понимаю разницу между Managed и External таблицами - знаю что происходит при DROP TABLE в каждом случае
- [ ] Всегда создаю External таблицы с LOCATION в production
- [ ] Знаю полный синтаксис CREATE TABLE: IF NOT EXISTS, USING, PARTITIONED BY, LOCATION, COMMENT, TBLPROPERTIES
- [ ] Использую
USING DELTAилиUSING ICEBERGявно, не оставляю default формат - [ ] Умею создавать таблицы через CTAS (CREATE TABLE ... AS SELECT)
- [ ] Знаю ALTER TABLE: ADD COLUMN, RENAME COLUMN, CHANGE TYPE (с ограничениями безопасных cast), SET TBLPROPERTIES
- [ ] Понимаю разницу DESCRIBE TABLE / DESCRIBE EXTENDED / DESCRIBE DETAIL - когда что использовать
- [ ] Умею читать DESCRIBE DETAIL: numFiles, sizeInBytes, диагностика small file problem
- [ ] Знаю SHOW DATABASES, SHOW TABLES, SHOW PARTITIONS, SHOW TBLPROPERTIES
- [ ] Умею использовать spark.catalog API: listTables, tableExists, refreshTable
- [ ] Понимаю роль ANALYZE TABLE для статистики оптимизатора
- [ ] Знаю REFRESH TABLE vs MSCK REPAIR TABLE - когда что нужно
- [ ] Пишу COMMENT на каждую таблицу и колонку в production-схемах
- [ ] Использую структурированные TBLPROPERTIES: pipeline.owner, pipeline.version, data.retention_days