Hive Metastore: external vs managed tables и partitioned tables DDL
Архитектура HMS, lifecycle managed vs external tables, DDL партиционирования, MSCK REPAIR TABLE, динамические партиции и multi-engine interoperability.
Архитектура 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':
- Catalyst обращается к HMS для получения схемы таблицы
gold.revenue - HMS возвращает список колонок, типы, физический путь (
s3a://datalake/gold/revenue/) - Catalyst проверяет: таблица партиционирована по
dt, значит нужна только директорияdt=2024-01-15/ - HMS возвращает физический путь этой партиции
- 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:
-
Создать базу данных
analyticsс LOCATION, COMMENT и DBPROPERTIES (owner,created_date) -
Создать external партиционированную таблицу
analytics.sales(USING PARQUET): - Колонки:
sale_id BIGINT,product_id BIGINT,region STRING,amount DECIMAL(12,2),quantity INT - Партиционирование по
sale_date DATEиcountry STRING -
LOCATION, COMMENT, TBLPROPERTIES
-
Заполнить три партиции данными через
INSERT INTO ... PARTITION (sale_date, country) -
Добавить новую партицию
sale_date='2024-02-01', country='KZ'вручную: - Создать файлы по пути вне Spark
-
Зарегистрировать через
ALTER TABLE ADD PARTITIONс явным LOCATION -
Приложить вывод
DESCRIBE FORMATTED analytics.salesиSHOW PARTITIONS analytics.sales -
Убедиться через 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 с точечной регистрацией партиций