Hive Metastore: external vs managed tables и partitioned tables DDL

Архитектура HMS, lifecycle managed vs external tables, DDL партиционирования, MSCK REPAIR TABLE, динамические партиции и multi-engine interoperability.

core

Архитектура Hive Metastore

Spark - вычислительный движок. Он умеет читать Parquet, трансформировать данные, выполнять join. Но сам по себе Spark не знает о существовании таблицы gold.revenue - ему нужно где-то хранить информацию: какие таблицы существуют, какова их схема, где физически лежат файлы, на какие партиции разбиты данные.

Hive Metastore (HMS) - это сервис, отвечающий за хранение всех этих метаданных. Исторически HMS появился как часть Apache Hive - первой SQL-системы поверх Hadoop. Затем его стандарт подхватили все остальные движки экосистемы: Spark, Trino, Presto, Impala, Flink. Сегодня HMS - де-факто стандартный каталог метаданных в open-source Big Data мире.

Что именно хранит HMS

Metastore - это реляционная база данных (обычно PostgreSQL или MySQL) с фиксированной схемой. Ключевые таблицы этой базы:

  • DBS - базы данных (databases): имя, описание, путь к warehouse-директории
  • TBLS - таблицы: имя, тип (MANAGED/EXTERNAL), ссылка на базу данных, время создания
  • COLUMNS_V2 - колонки каждой таблицы: имя, тип, порядковый номер, комментарий
  • PARTITION_KEYS - описание partition-колонок (имя, тип, порядок)
  • PARTITIONS - зарегистрированные партиции: значения ключей + физический путь
  • TABLE_PARAMS и PARTITION_PARAMS - дополнительные свойства (TBLPROPERTIES)
  • TAB_COL_STATS и PART_COL_STATS - статистика для Cost-Based Optimizer

Когда Spark выполняет SELECT * FROM gold.revenue WHERE dt = '2024-01-15':

  1. Catalyst обращается к HMS для получения схемы таблицы gold.revenue
  2. HMS возвращает список колонок, типы, физический путь (s3a://datalake/gold/revenue/)
  3. Catalyst проверяет: таблица партиционирована по dt, значит нужна только директория dt=2024-01-15/
  4. HMS возвращает физический путь этой партиции
  5. Spark читает только Parquet-файлы из этой директории

Multi-engine interoperability

Главная ценность HMS - разделение каталога между движками. Таблица, созданная через Spark, немедленно доступна в Trino, Presto, Flink и Apache Hive. Все они смотрят в один HMS, видят одну схему, одни партиции.

Это фундамент open lakehouse architecture: данные лежат в объектном хранилище (S3, MinIO, HDFS), метаданные - в HMS, вычисления - в любом совместимом движке по вкусу. Нет vendor lock-in: если завтра решите заменить Spark на Trino для интерактивных запросов - таблицы переезжать не нужно.

Подключение Spark к HMS

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("hive-demo") \
    .config("spark.sql.warehouse.dir", "s3a://datalake/warehouse/") \
    .config("hive.metastore.uris", "thrift://hms-host:9083") \
    .enableHiveSupport() \
    .getOrCreate()

# Проверяем подключение
spark.sql("SHOW DATABASES").show()

enableHiveSupport() - ключевой метод. Без него Spark использует in-memory session catalog, который не персистентен и не виден другим движкам. С ним Spark подключается к HMS через Thrift-протокол на порту 9083.

spark.sql.warehouse.dir - путь по умолчанию для MANAGED-таблиц. Все managed-таблицы без явного LOCATION создадутся внутри этой директории.

В embedded режиме (без внешнего HMS, например в тестах) Spark поднимает локальный Derby-based Metastore автоматически при enableHiveSupport(). Derby - однопользовательская встроенная БД, подходит только для разработки.


Managed vs External таблицы: великое противостояние

Это архитектурное решение с необратимыми последствиями. Неправильный выбор в production может привести к безвозвратной потере петабайт данных.

Managed (Internal) таблицы

При создании managed-таблицы HMS берёт под контроль и метаданные, и данные. Данные хранятся в warehouse-директории:

-- MANAGED таблица: нет EXTERNAL, нет LOCATION
CREATE TABLE silver.transactions (
    txn_id      BIGINT,
    user_id     BIGINT,
    amount      DECIMAL(12,2),
    txn_date    DATE,
    status      STRING
)
USING PARQUET
PARTITIONED BY (txn_date);
-- Данные будут созданы в: s3a://datalake/warehouse/silver.db/transactions/

Физически файлы появляются в <warehouse.dir>/<database>.db/<table>/. При записи данных Spark знает куда писать - берёт путь из метаданных.

Что происходит при DROP TABLE:

DROP TABLE silver.transactions;
-- Шаг 1: HMS удаляет запись из TBLS, COLUMNS_V2, PARTITIONS
-- Шаг 2: Spark инициирует удаление ДИРЕКТОРИИ s3a://datalake/warehouse/silver.db/transactions/
-- Все Parquet-файлы внутри УДАЛЕНЫ ФИЗИЧЕСКИ и безвозвратно

Это не предупреждает, не спрашивает подтверждения. Один DROP TABLE - и терабайты данных исчезли.

External таблицы

External-таблица - это только метаданные. Данные лежат по пути, который вы указываете, и HMS никогда не управляет их жизненным циклом:

-- EXTERNAL таблица: явный EXTERNAL и LOCATION
CREATE EXTERNAL TABLE silver.transactions (
    txn_id      BIGINT,
    user_id     BIGINT,
    amount      DECIMAL(12,2),
    txn_date    DATE,
    status      STRING
)
USING PARQUET
PARTITIONED BY (txn_date)
LOCATION 's3a://datalake/silver/transactions/';
-- Данные лежат по указанному пути, HMS только регистрирует путь

Что происходит при DROP TABLE:

DROP TABLE silver.transactions;
-- Шаг 1: HMS удаляет запись из TBLS, COLUMNS_V2, PARTITIONS
-- Шаг 2: НИЧЕГО БОЛЬШЕ. Файлы в s3a://datalake/silver/transactions/ НЕ ТРОНУТЫ

Данные остаются в хранилище. Таблицу можно восстановить повторной регистрацией:

-- Данные не потеряны - заново регистрируем
CREATE EXTERNAL TABLE silver.transactions (...)
USING PARQUET
LOCATION 's3a://datalake/silver/transactions/';
-- Voilà: таблица восстановлена, все данные доступны

В чём разница с точки зрения DDL

В Spark SQL синтаксис немного отличается от классического Hive DDL:

-- Hive DDL style (работает с enableHiveSupport):
CREATE EXTERNAL TABLE silver.orders (...) STORED AS PARQUET LOCATION '...';

-- Spark SQL style (рекомендуется):
CREATE TABLE silver.orders (...) USING PARQUET LOCATION '...';
-- Наличие LOCATION делает таблицу External автоматически!

-- Явно EXTERNAL:
CREATE EXTERNAL TABLE silver.orders (...) USING PARQUET LOCATION '...';

В Spark SQL: если указан LOCATION - таблица внешняя, даже без слова EXTERNAL. Если LOCATION нет - таблица managed. Это можно проверить:

DESCRIBE FORMATTED silver.orders;
-- ...
-- Type: EXTERNAL   ← если есть LOCATION
-- Type: MANAGED    ← если LOCATION нет

Таблица решений: когда что использовать

Сценарий Managed External
Данные создаются и управляются только Spark
Данные уже существуют в S3/HDFS
Несколько движков (Spark + Trino + Flink) ✗ рискованно
DROP TABLE не должен удалять данные
Production Bronze/Silver/Gold layers
Временные артефакты вычислений
Data Scientists' scratch space

Золотое правило production lakehouse: все Bronze, Silver и Gold таблицы - external. Managed таблицы только для scratch-пространства и временных вычислений.


Партиционирование: физическая организация данных

Зачем нужно партиционирование

Без партиционирования запрос WHERE order_date = '2024-01-15' читает весь датасет - все файлы, все строки - и потом фильтрует. При 10 TB данных за год это значит читать 10 TB ради получения 27 GB (данных за один день).

С партиционированием по order_date данные физически разложены по директориям:

s3a://datalake/silver/orders/
  order_date=2024-01-01/
    part-00001-abc.parquet
    part-00002-abc.parquet
  order_date=2024-01-02/
    part-00001-def.parquet
  ...
  order_date=2024-12-31/
    part-00001-xyz.parquet

Запрос WHERE order_date = '2024-01-15' читает только директорию order_date=2024-01-15/. Остальные 364 директории игнорируются полностью. Это называется partition pruning - Spark просто не открывает файлы из других партиций.

Экономия: 10 TB → 27 GB. Ускорение: ~365×.

Синтаксис: PARTITIONED BY

-- Одна колонка партиционирования
CREATE EXTERNAL TABLE silver.orders (
    order_id    BIGINT,
    customer_id BIGINT,
    total       DECIMAL(12,2),
    status      STRING
    -- order_date НЕТ в основной схеме! Она вынесена в PARTITIONED BY
)
USING PARQUET
PARTITIONED BY (order_date DATE)
LOCATION 's3a://datalake/silver/orders/';

-- Несколько колонок партиционирования (иерархия директорий)
CREATE EXTERNAL TABLE silver.events (
    event_id    STRING,
    user_id     BIGINT,
    event_type  STRING,
    payload     STRING
)
USING PARQUET
PARTITIONED BY (event_date DATE, country STRING)
LOCATION 's3a://datalake/silver/events/';
-- Структура: event_date=2024-01-15/country=RU/part-0001.parquet

Критически важно: колонки в PARTITIONED BY не включаются в основную схему таблицы (список колонок в начале DDL). Они определяются только в блоке PARTITIONED BY. При чтении Spark добавляет их из имён директорий - order_date будет в результирующем DataFrame автоматически.

Если добавить order_date и в список колонок, и в PARTITIONED BY - ошибка: AnalysisException: Found duplicate column(s) in the table definition.

Hive-style partition layout

Физическая структура директорий - key=value на каждом уровне:

s3a://datalake/silver/events/
  event_date=2024-01-01/
    country=RU/
      part-0001.parquet   ← строки за 2024-01-01, страна RU
      part-0002.parquet
    country=US/
      part-0001.parquet   ← строки за 2024-01-01, страна US
  event_date=2024-01-02/
    country=RU/
      part-0001.parquet
    country=DE/
      part-0001.parquet

Это и есть Hive partition layout. Spark умеет читать такую структуру автоматически, извлекая значения партиций из имён директорий.


Запись данных: статические и динамические партиции

Динамические партиции через df.write.partitionBy()

Самый удобный способ - позволить Spark самому создать партиции при записи:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date, current_timestamp

spark = SparkSession.builder \
    .appName("partitions-demo") \
    .enableHiveSupport() \
    .getOrCreate()

# Допустим, orders_df содержит колонку order_date
orders_df = spark.read.format("delta").load("s3a://datalake/bronze/orders/")

# Spark автоматически создаст директории order_date=.../
orders_df.write \
    .format("parquet") \
    .mode("overwrite") \
    .partitionBy("order_date") \
    .option("path", "s3a://datalake/silver/orders/") \
    .saveAsTable("silver.orders")
    # Spark: создаёт директории, регистрирует партиции в HMS

При saveAsTable Spark не только пишет файлы, но и регистрирует созданные партиции в HMS. После завершения SHOW PARTITIONS silver.orders покажет все новые партиции.

Для append-режима (добавление новых партиций):

new_orders = spark.read.format("delta") \
    .load("s3a://datalake/bronze/orders/") \
    .filter(col("order_date") == "2024-01-15")

new_orders.write \
    .format("parquet") \
    .mode("append") \
    .partitionBy("order_date") \
    .option("path", "s3a://datalake/silver/orders/") \
    .saveAsTable("silver.orders")
    # Добавляет только новую партицию order_date=2024-01-15

INSERT INTO / INSERT OVERWRITE через SQL

-- Динамический INSERT: Spark сам определяет значения партиций из данных
INSERT INTO silver.orders PARTITION (order_date)
SELECT order_id, customer_id, total, status, order_date
FROM bronze.raw_orders
WHERE order_date = '2024-01-15';

-- INSERT OVERWRITE: перезаписывает конкретную партицию
INSERT OVERWRITE silver.orders PARTITION (order_date = '2024-01-15')
SELECT order_id, customer_id, total, status
FROM bronze.raw_orders
WHERE order_date = '2024-01-15';
-- ВАЖНО: при OVERWRITE указываем статическое значение → только эта партиция перезаписывается
-- Остальные партиции НЕ ТРОНУТЫ

Режим перезаписи: partitionOverwriteMode

При INSERT OVERWRITE без явного значения партиции поведение зависит от настройки:

# Режим STATIC (по умолчанию): перезаписывает ВСЕ партиции
# Опасен: если вставляете только новый день - удалите все остальные дни!
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")

# Режим DYNAMIC: перезаписывает только затронутые партиции
# Рекомендуется для инкрементальных загрузок
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

orders_df.write \
    .format("parquet") \
    .mode("overwrite") \
    .partitionBy("order_date") \
    .save("s3a://datalake/silver/orders/")
# С DYNAMIC: только партиции из orders_df будут перезаписаны
# Партиции отсутствующих дней - не тронуты

dynamic - правильный режим для incremental ETL. static - только если намеренно хотите полную перезапись.


Синхронизация партиций: HMS и файловая система

Проблема: данные есть на диске, HMS не знает

Если файлы появились в S3 без участия Spark (aws s3 cp, другой инструмент, Airflow S3Hook), HMS об этих партициях не знает. Запрос через Spark вернёт пустой результат или не найдёт данные.

# Ситуация: данные загружены напрямую в S3
s3a://datalake/silver/orders/
  order_date=2024-01-14/   ← HMS знает (Spark писал)
    part-0001.parquet
  order_date=2024-01-15/   ← HMS НЕ знает (загружено напрямую)
    part-0001.parquet       ← файл есть, но Spark его не видит!

SELECT COUNT(*) FROM silver.orders WHERE order_date = '2024-01-15';
-- Результат: 0 (ноль! хотя данные физически есть)

MSCK REPAIR TABLE: полное сканирование

MSCK REPAIR TABLE (Hive MSCK = "MetaStor Consistency Kit") сканирует весь LOCATION-путь таблицы рекурсивно, находит все директории вида key=value/, и регистрирует отсутствующие в HMS.

MSCK REPAIR TABLE silver.orders;
-- Вывод:
-- Partitions not in metastore: silver.orders:order_date=2024-01-15
-- Repair: 1 partition(s) are recovered.

SHOW PARTITIONS silver.orders;
-- order_date=2024-01-14
-- order_date=2024-01-15   ← теперь видна

Ограничения MSCK REPAIR TABLE:

  • Медленно: рекурсивно сканирует весь путь на S3. При тысячах партиций - LIST-запросы к S3 занимают минуты
  • Только добавляет: регистрирует партиции, которые есть на диске, но нет в HMS. Не удаляет "мёртвые" партиции, которые есть в HMS, но исчезли с диска
  • Нет фильтрации: нельзя запустить только для новых партиций - только полное сканирование
  • S3 eventual consistency: в редких случаях LIST может не видеть только что созданные объекты

Для таблицы с 10 000 партиций (27 лет ежедневных данных) MSCK REPAIR TABLE может занять 30+ минут. Это неприемлемо для продакшна.

ALTER TABLE ADD PARTITION: точечная регистрация

Лучше MSCK для incremental pipeline: явно регистрировать только новую партицию.

-- Регистрируем конкретную партицию
ALTER TABLE silver.orders ADD PARTITION (order_date = '2024-01-15');
-- HMS записывает: order_date=2024-01-15 → s3a://datalake/silver/orders/order_date=2024-01-15/

-- Несколько партиций за раз
ALTER TABLE silver.orders
ADD PARTITION (order_date = '2024-01-15')
ADD PARTITION (order_date = '2024-01-16')
ADD PARTITION (order_date = '2024-01-17');

-- С кастомным путём (если директория не совпадает с Hive naming convention)
ALTER TABLE silver.orders
ADD PARTITION (order_date = '2024-01-15')
LOCATION 's3a://datalake/silver/orders/custom-path-for-jan-15/';
-- Парт. имеет свой path - отличается от стандартного key=value/

-- Если партиция уже есть - IF NOT EXISTS защищает от ошибки
ALTER TABLE silver.orders ADD IF NOT EXISTS PARTITION (order_date = '2024-01-15');

ALTER TABLE DROP PARTITION

-- Удалить партицию
ALTER TABLE silver.orders DROP PARTITION (order_date = '2024-01-01');
-- Для EXTERNAL таблицы: удаляется только запись в HMS, файлы на S3 остаются
-- Для MANAGED таблицы: удаляется запись в HMS И физические файлы

-- С условием (диапазон)
ALTER TABLE silver.orders DROP PARTITION (order_date < '2023-01-01');
-- Удаляет все партиции до 2023 года

-- Предохранитель: IF EXISTS
ALTER TABLE silver.orders DROP IF EXISTS PARTITION (order_date = '2024-01-01');

SHOW PARTITIONS: инспекция

-- Все партиции таблицы
SHOW PARTITIONS silver.orders;
-- order_date=2024-01-01
-- order_date=2024-01-02
-- ...

-- Фильтрация (только конкретное значение)
SHOW PARTITIONS silver.events PARTITION (event_date = '2024-01-15');
-- event_date=2024-01-15/country=DE
-- event_date=2024-01-15/country=RU
-- event_date=2024-01-15/country=US

-- Количество партиций (через count)
spark.sql("SHOW PARTITIONS silver.orders").count()

Инкрементальный pipeline: правильный паттерн

from datetime import date, timedelta
from pyspark.sql.functions import col, current_timestamp

def load_partition(spark, partition_date: date) -> None:
    """Загружает одну партицию инкрементально."""
    date_str = partition_date.isoformat()

    # 1. Читаем данные за дату из Bronze
    daily_data = spark.read.format("delta") \
        .load("s3a://datalake/bronze/orders/") \
        .filter(col("order_date") == date_str)

    if daily_data.isEmpty():
        print(f"No data for {date_str}, skipping")
        return

    # 2. Трансформируем
    silver_data = daily_data.select(
        col("order_id"),
        col("customer_id"),
        col("total_amount"),
        col("status"),
        current_timestamp().alias("_ingested_at"),
    )

    # 3. Пишем в S3
    output_path = f"s3a://datalake/silver/orders/order_date={date_str}/"
    silver_data.write \
        .format("parquet") \
        .mode("overwrite") \
        .save(output_path)

    # 4. Регистрируем партицию в HMS точечно (не MSCK!)
    spark.sql(f"""
        ALTER TABLE silver.orders ADD IF NOT EXISTS
        PARTITION (order_date = '{date_str}')
        LOCATION '{output_path}'
    """)
    print(f"Loaded partition: order_date={date_str}")


# Загружаем последние 3 дня
today = date.today()
for i in range(3):
    load_partition(spark, today - timedelta(days=i))

Этот паттерн быстр (нет MSCK-сканирования), идемпотентен (ADD IF NOT EXISTS), точечен (только нужные партиции).


Partition Pruning: как это работает внутри

Механизм отсечения партиций

Когда Catalyst видит фильтр на колонку партиции, он получает из HMS список партиций и их физические пути, затем выбирает только те, которые соответствуют условию:

result = spark.sql("""
    SELECT SUM(total_amount)
    FROM silver.orders
    WHERE order_date BETWEEN '2024-01-01' AND '2024-01-31'
""")
result.explain(mode="formatted")

В Physical Plan:

FileScan parquet silver.orders [order_id, total_amount, order_date]
  PartitionFilters: [isnotnull(order_date#10),
                     (order_date#10 >= 2024-01-01),
                     (order_date#10 <= 2024-01-31)]
  PushedFilters: []
  ReadSchema: struct<order_id:bigint, total_amount:decimal(12,2)>
  Location: s3a://datalake/silver/orders/
  SelectedFiles: 31 files from 31 partitions  ← только январь!

PartitionFilters - фильтрация на уровне метаданных HMS (не читая файлы). SelectedFiles: 31 - только 31 файл из 365 партиций.

Важный нюанс: фильтр должен быть на partition key

-- Pruning РАБОТАЕТ: order_date - partition key
SELECT * FROM silver.orders WHERE order_date = '2024-01-15';

-- Pruning НЕ РАБОТАЕТ: customer_id - обычная колонка, не partition key
SELECT * FROM silver.orders WHERE customer_id = 12345;
-- Spark прочитает ВСЕ файлы и потом отфильтрует строки

-- Pruning ЧАСТИЧНО работает: один из двух фильтров на partition key
SELECT * FROM silver.orders
WHERE order_date = '2024-01-15' AND customer_id = 12345;
-- Spark прочитает только order_date=2024-01-15/, потом фильтрует customer_id

Partition pruning vs file-level predicate pushdown

Это два разных уровня оптимизации:

  • Partition pruning (HMS level): исключает целые директории - нет LIST, нет OPEN, нет чтения
  • Predicate pushdown (Parquet level): внутри файла пропускает row groups, не соответствующие условию

Оба применяются вместе. Сначала pruning убирает лишние директории, потом внутри выбранных файлов predicate pushdown убирает лишние row groups.


Антипаттерны партиционирования

Антипаттерн 1: партиционирование по high-cardinality колонке

-- КАТАСТРОФА: партиционирование по user_id
CREATE EXTERNAL TABLE silver.events (
    event_id STRING,
    event_type STRING,
    amount DECIMAL
)
USING PARQUET
PARTITIONED BY (user_id BIGINT)  -- ← 50 миллионов уникальных пользователей!
LOCATION 's3a://datalake/silver/events/';

Последствия:

  • 50 миллионов директорий в S3: LIST-операции становятся невозможными
  • 50 миллионов строк в таблице PARTITIONS в HMS-базе: PostgreSQL страдает, запросы к метастору тормозят все движки
  • Миллионы мелких файлов: каждый пользователь создаёт маленький файл - small file problem
  • Catalyst зависает: при планировании запроса без фильтра по user_id Catalyst пытается получить все 50M партиций из HMS

Правило: карdinality partition-ключа не должна превышать 10 000–50 000 значений для стабильной работы HMS.

Антипаттерн 2: слишком гранулярное партиционирование по времени

-- ПЛОХО: партиционирование по минутам (high cardinality по времени)
PARTITIONED BY (event_minute STRING)
-- event_minute=2024-01-15T10:30/ → один файл за минуту → миллион файлов в год

-- ХОРОШО: партиционирование по дате
PARTITIONED BY (event_date DATE)
-- event_date=2024-01-15/ → разумное количество файлов за день

Стриминговые пайплайны часто пишут маленькими батчами каждые N минут. При партиционировании по дате все батчи за день попадают в одну директорию. Затем раз в день запускают OPTIMIZE / компакцию.

Антипаттерн 3: нет партиционирования вообще для большой таблицы

-- ПЛОХО: таблица 10 TB без партиционирования
CREATE EXTERNAL TABLE silver.events (
    event_id STRING, event_date DATE, user_id BIGINT, ...
)
USING PARQUET LOCATION 's3a://datalake/silver/events/';
-- Любой запрос с WHERE event_date = '...' → скан 10 TB

Для таблиц меньше 1 GB партиционирование не нужно - overhead метаданных перевешивает выгоду. Для таблиц больше нескольких GB - почти всегда нужно.

Антипаттерн 4: слишком много уровней партиционирования

-- ИЗБЫТОЧНО: 4 уровня партиционирования
PARTITIONED BY (year INT, month INT, day INT, hour INT)
-- year=2024/month=01/day=15/hour=10/ → 8760 директорий в год
-- Catalyst при запросе без всех 4 фильтров сканирует слишком много метаданных

-- ЛУЧШЕ: один или два уровня
PARTITIONED BY (event_date DATE)
-- Плюс partition pruning по диапазону дат работает отлично
-- Если нужно hourly: event_date=2024-01-15/ и крупные файлы внутри

Оптимальная стратегия партиционирования

Тип данных Рекомендация Пример
Транзакции, события DATE или YYYY-MM event_date DATE
Мультирегиональные данные DATE + REGION event_date DATE, country STRING
Многотенантные системы DATE + TENANT report_date DATE, tenant_id STRING
Таблицы-справочники (< 1 GB) Без партиционирования -
Таблицы с 1-50 категориями По категории product_type STRING
High-cardinality (users, IDs) НЕ партиционировать Используйте Z-ORDER / file-level stats

Идеальный размер партиции: 128 MB – 1 GB. При меньшем размере - слишком много мелких файлов. При большем - partition pruning даёт меньший выигрыш.


Schema Evolution в партиционированных таблицах

Добавление колонки в партиционированную таблицу

-- Добавить новую колонку
ALTER TABLE silver.orders ADD COLUMNS (promo_code STRING);

-- Существующие файлы: promo_code = null при чтении
-- Новые файлы: promo_code содержит значение

-- Важно: нельзя добавить колонку в PARTITIONED BY через ALTER TABLE
-- Новый partition key - только через создание новой таблицы

Изменение типов в партиционированных таблицах

При изменении типа партиционной колонки - только "безопасные" расширения:

-- Безопасно: DATE → STRING (при чтении Hive-style партиций ключ всегда строка)
ALTER TABLE silver.orders CHANGE COLUMN order_date order_date STRING;

-- ОПАСНО: изменение типа затрагивает имена директорий!
-- Директории event_date=2024-01-15 уже созданы как DATE-формат
-- Смена на другой формат может сломать partition resolution

HMS и современные форматы: Delta Lake и Iceberg

Delta Lake + HMS

Delta Lake хранит собственный транзакционный лог (_delta_log/) и не зависит от HMS для трекинга файлов. Но HMS всё равно используется для регистрации таблицы (schema, location):

-- Delta таблица регистрируется в HMS как "delta" provider
CREATE EXTERNAL TABLE silver.orders_delta (
    order_id BIGINT, customer_id BIGINT, total DECIMAL(12,2)
)
USING DELTA
PARTITIONED BY (order_date DATE)
LOCATION 's3a://datalake/silver/orders_delta/';

Для Delta таблицы:

  • Схема и партиции хранятся в _delta_log/, а не в HMS
  • HMS хранит только: имя таблицы, location, тип провайдера (delta)
  • MSCK REPAIR TABLE для Delta не нужен: Delta знает о всех файлах из лога
  • SHOW PARTITIONS для Delta: Spark читает лог, не HMS
-- Для Delta: используйте DESCRIBE DETAIL вместо SHOW PARTITIONS
DESCRIBE DETAIL silver.orders_delta;
-- Покажет numFiles, numPartitions, sizeInBytes

Apache Iceberg + HMS

Iceberg аналогично хранит собственные метаданные (metadata layer: manifests, manifest lists). HMS используется только как "указатель":

CREATE EXTERNAL TABLE silver.orders_iceberg (
    order_id BIGINT, customer_id BIGINT, total DECIMAL(12,2), order_date DATE
)
USING ICEBERG
LOCATION 's3a://datalake/silver/orders_iceberg/';

Iceberg поддерживает скрытое партиционирование (hidden partitioning) - данные физически партиционированы, но из SQL-запросов это прозрачно.


Практика: полный жизненный цикл

Шаг 1: Демонстрация "катастрофы" с Managed таблицей

spark = SparkSession.builder \
    .appName("hive-practice") \
    .config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
    .enableHiveSupport() \
    .getOrCreate()

# Создаём MANAGED таблицу
spark.sql("""
    CREATE TABLE demo_managed (
        id BIGINT, name STRING, amount DECIMAL(10,2)
    )
    USING PARQUET
""")

# Заполняем данными
spark.sql("""
    INSERT INTO demo_managed VALUES
    (1, 'Alice', 1500.00),
    (2, 'Bob',   2200.00),
    (3, 'Carol',  800.00)
""")

spark.sql("SELECT * FROM demo_managed").show()

# Проверяем где физически лежат файлы
detail = spark.sql("DESCRIBE FORMATTED demo_managed").collect()
for row in detail:
    if row["col_name"] in ["Location", "Type"]:
        print(f"{row['col_name']}: {row['data_type']}")
# Type:     MANAGED
# Location: /tmp/spark-warehouse/demo_managed

# КАТАСТРОФА: DROP TABLE удаляет данные физически
spark.sql("DROP TABLE demo_managed")
# Файл /tmp/spark-warehouse/demo_managed/ УДАЛЁН
# Восстановить невозможно!
print("Данные безвозвратно удалены!")

Шаг 2: Правильная External таблица с партиционированием

import os
from pyspark.sql import Row
from pyspark.sql.types import *
from pyspark.sql.functions import col, current_date, lit, to_date
from datetime import date, timedelta
import random

# Создаём базу данных
spark.sql("CREATE DATABASE IF NOT EXISTS silver_demo LOCATION '/tmp/silver_demo'")

# Схема данных
order_schema = StructType([
    StructField("order_id", LongType()),
    StructField("customer_id", LongType()),
    StructField("total_amount", DoubleType()),
    StructField("status", StringType()),
])

# Создаём EXTERNAL партиционированную таблицу
spark.sql("""
    CREATE EXTERNAL TABLE IF NOT EXISTS silver_demo.orders (
        order_id     BIGINT  COMMENT 'Unique order identifier',
        customer_id  BIGINT  COMMENT 'FK to customers',
        total_amount DOUBLE  COMMENT 'Order total in RUB',
        status       STRING  COMMENT 'completed/pending/cancelled'
    )
    USING PARQUET
    PARTITIONED BY (order_date DATE)
    LOCATION '/tmp/silver_demo_orders/'
    COMMENT 'Silver: normalized orders partitioned by date'
    TBLPROPERTIES ('pipeline.owner' = 'data-team')
""")

# Проверяем тип: должно быть EXTERNAL
spark.sql("DESCRIBE FORMATTED silver_demo.orders") \
    .filter(col("col_name") == "Type") \
    .show(truncate=False)

Шаг 3: Запись данных с автоматическим созданием партиций

# Генерируем данные за 3 дня
all_data = []
for day_offset in range(3):
    order_date = date(2024, 1, 15) + timedelta(days=day_offset)
    for i in range(100):
        all_data.append(Row(
            order_id=day_offset * 1000 + i,
            customer_id=random.randint(1, 500),
            total_amount=round(random.uniform(100, 5000), 2),
            status=random.choice(["completed", "pending", "cancelled"]),
            order_date=order_date,
        ))

schema_with_date = StructType([
    StructField("order_id", LongType()),
    StructField("customer_id", LongType()),
    StructField("total_amount", DoubleType()),
    StructField("status", StringType()),
    StructField("order_date", DateType()),
])

orders_df = spark.createDataFrame(all_data, schema_with_date)

# Spark записывает и автоматически регистрирует партиции в HMS
orders_df.write \
    .format("parquet") \
    .mode("append") \
    .partitionBy("order_date") \
    .option("path", "/tmp/silver_demo_orders/") \
    .saveAsTable("silver_demo.orders")

# Проверяем партиции
spark.sql("SHOW PARTITIONS silver_demo.orders").show()
+------------------+
|         partition|
+------------------+
|order_date=2024-01-15|
|order_date=2024-01-16|
|order_date=2024-01-17|
+------------------+

Шаг 4: Имитация "невидимой" партиции и MSCK REPAIR TABLE

import os
import shutil

# Имитируем: данные добавлены напрямую в S3, минуя Spark
new_date = "2024-01-18"
new_path = f"/tmp/silver_demo_orders/order_date={new_date}/"
os.makedirs(new_path, exist_ok=True)

# Создаём "сторонние" данные и кладём как Parquet-файл напрямую
new_orders = spark.createDataFrame(
    [(5000 + i, random.randint(1, 500), round(random.uniform(100, 5000), 2), "completed")
     for i in range(50)],
    StructType([
        StructField("order_id", LongType()),
        StructField("customer_id", LongType()),
        StructField("total_amount", DoubleType()),
        StructField("status", StringType()),
    ])
)
new_orders.write.format("parquet").mode("overwrite").save(new_path)

# Проверяем: HMS не знает об этой партиции
print("SHOW PARTITIONS (до REPAIR):")
spark.sql("SHOW PARTITIONS silver_demo.orders").show()
# Только 3 партиции (2024-01-15, 16, 17) - 2024-01-18 нет!

# Запрос к "невидимой" партиции
count_invisible = spark.sql(
    "SELECT COUNT(*) FROM silver_demo.orders WHERE order_date = '2024-01-18'"
).collect()[0][0]
print(f"Строк за 2024-01-18 (до REPAIR): {count_invisible}")  # → 0

# Исправляем через MSCK REPAIR TABLE
spark.sql("MSCK REPAIR TABLE silver_demo.orders")
print("\nMSCK REPAIR завершён")

# Теперь партиция видна
print("\nSHOW PARTITIONS (после REPAIR):")
spark.sql("SHOW PARTITIONS silver_demo.orders").show()

count_fixed = spark.sql(
    "SELECT COUNT(*) FROM silver_demo.orders WHERE order_date = '2024-01-18'"
).collect()[0][0]
print(f"Строк за 2024-01-18 (после REPAIR): {count_fixed}")  # → 50

Шаг 5: Partition pruning - бенчмарк

# Таблица с данными за 3 дня (300 строк)
# Тест 1: без партиционного фильтра - читает всё
full_scan = spark.sql("SELECT COUNT(*) FROM silver_demo.orders")
full_scan.explain(mode="formatted")
# PartitionFilters: [] ← нет фильтра по партиции

# Тест 2: с партиционным фильтром - читает только один день
filtered = spark.sql("""
    SELECT COUNT(*) FROM silver_demo.orders
    WHERE order_date = '2024-01-15'
""")
filtered.explain(mode="formatted")
# PartitionFilters: [isnotnull(order_date), EqualTo(order_date, 2024-01-15)]
# SelectedFiles: 1 partition only ← только нужная партиция

# Подтверждение: EXPLAIN показывает разный объём чтения

Шаг 6: Точечное добавление партиции через ALTER TABLE

# Вместо MSCK REPAIR - точечная регистрация только новой партиции
new_date_str = "2024-01-19"
new_path_str = f"/tmp/silver_demo_orders/order_date={new_date_str}/"

# Создаём файл напрямую
new_data = spark.createDataFrame(
    [(6000 + i, random.randint(1, 500), round(random.uniform(100, 5000), 2), "completed")
     for i in range(30)],
    ["order_id", "customer_id", "total_amount", "status"]
)
new_data.write.format("parquet").mode("overwrite").save(new_path_str)

# Регистрируем точечно - мгновенная операция, нет сканирования S3!
spark.sql(f"""
    ALTER TABLE silver_demo.orders
    ADD IF NOT EXISTS PARTITION (order_date = '{new_date_str}')
    LOCATION '{new_path_str}'
""")
print("Партиция зарегистрирована мгновенно через ALTER TABLE ADD PARTITION")

spark.sql("SHOW PARTITIONS silver_demo.orders").show()

Шаг 7: DROP TABLE безопасно на External таблице

# Сохраняем location для проверки
location = spark.sql(
    "DESCRIBE FORMATTED silver_demo.orders"
).filter(col("col_name") == "Location") \
 .collect()[0]["data_type"]

print(f"Данные хранятся в: {location}")

# DROP TABLE на EXTERNAL таблице - безопасно
spark.sql("DROP TABLE silver_demo.orders")
print("DROP TABLE выполнен")

# Проверяем: данные остались!
import os
files_exist = os.path.exists(location)
print(f"Файлы остались на диске: {files_exist}")  # → True

# Восстанавливаем регистрацию
spark.sql(f"""
    CREATE EXTERNAL TABLE silver_demo.orders (
        order_id BIGINT, customer_id BIGINT,
        total_amount DOUBLE, status STRING
    )
    USING PARQUET
    PARTITIONED BY (order_date DATE)
    LOCATION '{location}'
""")
spark.sql("MSCK REPAIR TABLE silver_demo.orders")
print("Таблица восстановлена!")
spark.sql("SELECT COUNT(*) FROM silver_demo.orders").show()

SparkSession.catalog: программный доступ к HMS

# Список всех баз данных
for db in spark.catalog.listDatabases():
    print(f"DB: {db.name} @ {db.locationUri}")

# Список таблиц с их типами
for t in spark.catalog.listTables("silver_demo"):
    print(f"  {t.name}: type={t.tableType}, temp={t.isTemporary}")
    # tableType: MANAGED, EXTERNAL, VIEW

# Колонки таблицы (включая partition keys)
for col in spark.catalog.listColumns("silver_demo", "orders"):
    print(f"  {col.name}: {col.dataType} (part={col.isPartition})")

# Проверка существования
print(spark.catalog.tableExists("silver_demo.orders"))

# Refresh после внешних изменений
spark.catalog.refreshTable("silver_demo.orders")

# Аналог MSCK REPAIR через API (Spark 3.x)
spark.sql("MSCK REPAIR TABLE silver_demo.orders")

Мониторинг и инспекция HMS

DESCRIBE FORMATTED: полный дамп метаданных

DESCRIBE FORMATTED silver_demo.orders;

Ключевые строки вывода:

# Detailed Table Information
Database:            silver_demo
Table:               orders
Owner:               spark
Type:                EXTERNAL    ← тип таблицы
Provider:            parquet
Location:            /tmp/silver_demo_orders
Comment:             Silver: normalized orders

# Partition Information
# col_name    data_type
order_date    date

# Storage Information
SerDe Library:       org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe
InputFormat:         org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat
OutputFormat:        org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat
Compressed:          No
Num Buckets:         -1

# Table Properties
pipeline.owner       data-team

Аудиторская функция для HMS-таблиц

def hms_table_audit(spark, full_table_name: str) -> dict:
    """Сводная информация о Hive Metastore таблице."""
    parts = full_table_name.split(".")
    db, table = parts[0], parts[1]

    # Основные метаданные
    meta_rows = spark.sql(f"DESCRIBE FORMATTED {full_table_name}").collect()
    meta = {row["col_name"].strip(): row["data_type"].strip()
            for row in meta_rows if row["col_name"].strip()}

    # Количество партиций
    try:
        num_partitions = spark.sql(f"SHOW PARTITIONS {full_table_name}").count()
    except Exception:
        num_partitions = 0  # таблица не партиционирована

    return {
        "table": full_table_name,
        "type": meta.get("Type", "unknown"),
        "location": meta.get("Location", "unknown"),
        "provider": meta.get("Provider", "unknown"),
        "num_partitions": num_partitions,
        "comment": meta.get("Comment", ""),
    }


audit = hms_table_audit(spark, "silver_demo.orders")
for k, v in audit.items():
    print(f"  {k:20s}: {v}")

Best Practices

Чеклист при создании таблицы

Production Hive Metastore Table Checklist:
  [ ] Всегда EXTERNAL с явным LOCATION - никогда MANAGED в production
  [ ] Формат: USING PARQUET или USING DELTA (не оставлять Hive SerDe)
  [ ] Партиционирование по DATE или по DATE+CATEGORY (не по ID/timestamp)
  [ ] Cardinality partition-ключа < 50 000 уникальных значений
  [ ] COMMENT на таблице и на каждой колонке
  [ ] TBLPROPERTIES: pipeline.owner, pipeline.version, data.retention_days
  [ ] IF NOT EXISTS для идемпотентности DDL
  [ ] ALTER TABLE ADD PARTITION вместо MSCK REPAIR для incremental pipelines
  [ ] MSCK REPAIR только для одноразовой синхронизации, не в регулярных джобах

Когда использовать MSCK REPAIR vs ALTER TABLE ADD PARTITION

MSCK REPAIR TABLE:
  ✓ Одноразовая синхронизация после bulk-загрузки данных
  ✓ Начальная регистрация существующих данных в HMS
  ✓ Таблица имеет < 1000 партиций
  ✗ Не использовать в регулярных ETL-пайплайнах
  ✗ Не использовать на таблицах с 10K+ партиций (слишком медленно)

ALTER TABLE ADD PARTITION:
  ✓ Incremental ETL: добавление новых дней/партиций
  ✓ Высокочастотные пайплайны (hourly, daily)
  ✓ Когда известен точный путь и значение партиции
  ✓ Мгновенная операция - нет сканирования S3

Домашнее задание

Написать DDL-скрипт миграции аналитической витрины продаж в lakehouse:

  1. Создать базу данных analytics с LOCATION, COMMENT и DBPROPERTIES (owner, created_date)

  2. Создать external партиционированную таблицу analytics.sales (USING PARQUET):

  3. Колонки: sale_id BIGINT, product_id BIGINT, region STRING, amount DECIMAL(12,2), quantity INT
  4. Партиционирование по sale_date DATE и country STRING
  5. LOCATION, COMMENT, TBLPROPERTIES

  6. Заполнить три партиции данными через INSERT INTO ... PARTITION (sale_date, country)

  7. Добавить новую партицию sale_date='2024-02-01', country='KZ' вручную:

  8. Создать файлы по пути вне Spark
  9. Зарегистрировать через ALTER TABLE ADD PARTITION с явным LOCATION

  10. Приложить вывод DESCRIBE FORMATTED analytics.sales и SHOW PARTITIONS analytics.sales

  11. Убедиться через EXPLAIN что запрос WHERE sale_date = '2024-02-01' AND country = 'KZ' использует partition pruning (показать PartitionFilters в плане)


Чеклист

  • [ ] Понимаю роль HMS в экосистеме Spark: хранит метаданные, не данные
  • [ ] Знаю что хранится в HMS: DBS, TBLS, COLUMNS_V2, PARTITIONS, TABLE_PARAMS
  • [ ] Понимаю разницу MANAGED vs EXTERNAL: что происходит при DROP TABLE в каждом случае
  • [ ] Всегда создаю EXTERNAL таблицы с явным LOCATION в production
  • [ ] Знаю синтаксис: LOCATION делает таблицу External автоматически, без слова EXTERNAL
  • [ ] Понимаю Hive partition layout: key=value/ структура директорий
  • [ ] Знаю что partition-колонки не включаются в основную схему (только в PARTITIONED BY)
  • [ ] Умею настраивать partitionOverwriteMode = dynamic для безопасного overwrite
  • [ ] Знаю когда использовать MSCK REPAIR TABLE (одноразово) vs ALTER TABLE ADD PARTITION (регулярно)
  • [ ] Понимаю проблему "невидимых партиций" при записи данных в обход Spark
  • [ ] Знаю антипаттерны: партиционирование по user_id, по timestamp, 4+ уровня, без партиций для 10+ GB
  • [ ] Умею читать explain-план: вижу PartitionFilters и SelectedFiles для подтверждения partition pruning
  • [ ] Знаю особенности Delta Lake + HMS: схема в _delta_log, MSCK не нужен
  • [ ] Могу написать полный инкрементальный pipeline с точечной регистрацией партиций