Catalog API: CREATE TABLE, ALTER TABLE, DESCRIBE DETAIL

Управление метаданными в Spark: Hive Metastore, managed vs external, полный жизненный цикл таблиц через DDL и programmatic Catalog API для production lakehouse.

core

Зачем нужен каталог

В самой простой архитектуре 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 для конкретной таблицы. После этого следующий запрос прочитает актуальные метаданные из HMS
  • MSCK 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-архитектуры интернет-магазина:

  1. Создать три базы данных: bronze_shop, silver_shop, gold_shop

  2. Создать таблицы:

  3. bronze_shop.raw_orders - внешняя Delta-таблица, сырые заказы из Kafka (order_id, payload STRING, received_at, _source)
  4. silver_shop.orders - нормализованные заказы (order_id, customer_id, status, total_amount, order_date) с партиционированием по order_date
  5. silver_shop.order_items - позиции заказа (item_id, order_id, sku, qty, price, order_date)
  6. gold_shop.daily_revenue - CTAS-витрина: выручка по категориям и дням

  7. Для каждой таблицы указать COMMENT и TBLPROPERTIES: pipeline.owner, pipeline.version, data.retention_days

  8. Написать блок ALTER TABLE:

  9. Добавить колонку promo_code STRING в silver_shop.orders
  10. Переименовать total_amountorder_total в silver_shop.orders
  11. Обновить pipeline.version до 1.1.0 во всех трёх silver-таблицах

  12. Выполнить 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