Hive Metastore: HMS как Thrift-сервер метаданных, enableHiveSupport() и spark.catalog API
Глубокий разбор Hive Metastore для Spark Data Engineer: архитектура Thrift-сервера и RDBMS бэкенда, режимы развёртывания, Partition Explosion как metadata bottleneck, enableHiveSupport() и hive-site.xml, полный spark.catalog API, Managed vs External таблицы, MSCK REPAIR TABLE, тюнинг производительности и траблшутинг.
1. Архитектура каталогов данных: зачем нужен Hive Metastore¶
Когда в Data Lake накапливаются сотни таблиц, сотни партиций на каждую и десятки аналитиков, которые пишут SQL-запросы - возникает фундаментальный вопрос: откуда Spark знает, что SELECT * FROM sales.orders WHERE date = '2024-01-15' нужно читать именно из /data/warehouse/sales/orders/date=2024-01-15/? Кто хранит маппинг «логическое имя → физический путь»? Кто знает схему таблицы (названия колонок, типы данных)? Кто хранит список партиций?
Ответ: Hive Metastore (HMS). Это центральный каталог метаданных всего on-premise Data Lake.
Разделение данных и метаданных: фундаментальная идея¶
Ключевой архитектурный принцип HMS - разделение данных и их описания. Физические данные лежат в HDFS (или S3, MinIO, Ceph) в виде Parquet/ORC/Avro файлов. HMS хранит только описание этих данных:
- Где физически лежат файлы таблицы (HDFS path)
- Схема таблицы: названия и типы колонок
- Формат файлов (Parquet, ORC, текст)
- Список партиций и пути к ним
- Свойства таблицы (параметры сжатия, SerDe класс, custom metadata)
- Права доступа и владельцы
Благодаря этому разделению одни и те же физические данные может читать Spark, Hive, Trino, Presto, Flink - каждый через свой execution engine, но через единый источник правды о метаданных.
Анатомия HMS: три компонента¶
Схема раскрывает архитектуру: HMS - это сервис-посредник между клиентами (Spark, Hive, Trino) и реляционной базой данных метаданных. Клиенты никогда не ходят напрямую в PostgreSQL - только через Thrift API HMS. Сами данные клиенты читают напрямую из HDFS или S3.
Режимы развёртывания HMS¶
Embedded (встроенный) - HMS запускается внутри того же JVM-процесса что и Hive Server или Spark. Использует Derby как базу данных. Подходит только для разработки: каждая сессия имеет изолированное хранилище метаданных, которое исчезает при завершении.
Local (локальный) - HMS запускается в том же процессе, но использует внешнюю базу (PostgreSQL). Не имеет Thrift-сервера, к нему нельзя подключиться извне. Один клиент - одно соединение к СУБД. Не подходит для продакшна из-за отсутствия connection pooling.
Remote (удалённый) - production стандарт. Выделенный HMS процесс с Thrift-сервером на порту 9083. Все клиенты подключаются через Thrift. HMS управляет пулом соединений к PostgreSQL. Несколько HMS инстансов можно запустить для HA.
<!-- hive-site.xml: Remote mode конфигурация -->
<property>
<name>hive.metastore.uris</name>
<!-- Список HMS серверов через запятую для HA -->
<value>thrift://hms1.example.com:9083,thrift://hms2.example.com:9083</value>
</property>
<property>
<name>javax.jdo.option.ConnectionURL</name>
<!-- JDBC URL к PostgreSQL -->
<value>jdbc:postgresql://postgres.example.com:5432/hive_metastore?createDatabaseIfNotExist=true</value>
</property>
<property>
<name>javax.jdo.option.ConnectionDriverName</name>
<value>org.postgresql.Driver</value>
</property>
<property>
<name>javax.jdo.option.ConnectionUserName</name>
<value>hive</value>
</property>
<property>
<name>javax.jdo.option.ConnectionPassword</name>
<value>secure_password_here</value>
</property>
2. Под капотом: HMS как Thrift-сервер и сетевой оверхед¶
Понимание протокола взаимодействия между Spark и HMS критически важно для диагностики проблем производительности - особенно при работе с таблицами, имеющими тысячи партиций.
Протокол Apache Thrift: как работает RPC¶
Apache Thrift - это фреймворк для межпроцессного взаимодействия (RPC), созданный в Facebook. HMS использует Thrift для предоставления API: клиент вызывает метод как обычную функцию, Thrift сериализует параметры в бинарный формат, передаёт по TCP и десериализует ответ.
Thrift IDL (Interface Definition Language) определяет интерфейс HMS:
// Упрощённый фрагмент Thrift IDL HMS (из метастора Hive)
// Описывает что клиент может делать с таблицами
struct Table {
1: string tableName,
2: string dbName,
3: string owner,
4: StorageDescriptor sd, // путь, формат, SerDe
5: list<FieldSchema> partitionKeys,
6: map<string, string> parameters,
}
service ThriftHiveMetastore {
Table get_table(1:string dbname, 2:string tbl_name),
list<Partition> get_partitions(1:string db_name,
2:string tbl_name,
3:i16 max_parts),
void add_partition(1:Partition new_part),
// ... и ещё ~150 методов
}
Пошаговый путь запроса: от SELECT до данных¶
Эта диаграмма показывает, что HMS участвует в каждом SQL-запросе дважды: сначала для получения метаданных таблицы, затем для получения списка партиций, удовлетворяющих фильтру. Только после этого Spark переходит к реальному чтению данных. При большом числе партиций время шага 2 может доминировать над временем шага 3.
Partition Explosion: почему это убивает HMS¶
Partition Explosion - ситуация когда таблица имеет сотни тысяч партиций. Это типично для:
- Таблиц с почасовым партиционированием за несколько лет:
date=2024-01-01/hour=00→ 24 × 365 × 3 = 26 280 партиций - Таблиц с многоуровневым партиционированием:
country/date/event_type→ 50 × 365 × 10 = 182 500 партиций
Что происходит при запросе к такой таблице:
-- Если нет фильтра по ключевым партициям - Spark запрашивает ВСЕ партиции
SELECT SUM(amount) FROM events WHERE event_type = 'purchase'
HMS вынужден вернуть полный список партиций (182 500 строк из PostgreSQL), передать их через Thrift (десятки MB), Spark Driver должен их десериализовать и сохранить в памяти.
Время выполнения при 182 500 партициях:
- SQL запрос к PostgreSQL (get_partitions): 3-8 секунд
- Thrift serialization (182 500 партиций → binary): 1-2 секунды
- Spark Driver deserialization: 2-4 секунды
- Итого только metadata фаза: 6-14 секунд!
- И всё это ДО начала чтения реальных данных
В Spark UI это выражается как аномально высокий Scheduler Delay на первом Stage (Driver ждёт HMS, Executor'ы простаивают). Если HMS не отвечает вовремя - возникает MetaException: Got exception: java.net.SocketTimeoutException.
3. Активация интеграции: enableHiveSupport() и hive-site.xml¶
Без enableHiveSupport(): встроенный in-memory каталог¶
По умолчанию Spark создаёт SparkSession с in-memory каталогом. Таблицы, созданные в одной сессии, исчезают при её завершении. Нет Thrift-соединения, нет PostgreSQL - всё в памяти.
# Без enableHiveSupport: in-memory каталог
spark_no_hive = SparkSession.builder \
.appName("no-hive-demo") \
.getOrCreate()
# Создаём таблицу
spark_no_hive.range(100).write.saveAsTable("temp_numbers")
# Проверяем - таблица видна в текущей сессии
spark_no_hive.sql("SHOW TABLES").show()
# +--------+-----------+-----------+
# |database| tableName|isTemporary|
# +--------+-----------+-----------+
# | default|temp_numbers| false|
# +--------+-----------+-----------+
# Перезапускаем сессию
spark_no_hive.stop()
spark_no_hive2 = SparkSession.builder.appName("no-hive-demo-2").getOrCreate()
# Таблица исчезла!
spark_no_hive2.sql("SHOW TABLES").show()
# +--------+---------+-----------+
# |database|tableName|isTemporary|
# +--------+---------+-----------+
# +--------+---------+-----------+
# (пустой результат)
С enableHiveSupport(): постоянный каталог через HMS¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("hms-integration") \
# enableHiveSupport() делает несколько вещей:
# 1. Загружает hive-site.xml из classpath для конфигурации
# 2. Подключает HiveExternalCatalog вместо InMemoryCatalog
# 3. Включает поддержку HiveQL диалекта
# 4. Добавляет Hive UDF в функциональность Spark SQL
.enableHiveSupport() \
# Явная конфигурация HMS (если hive-site.xml не в classpath)
.config("hive.metastore.uris", "thrift://hms.example.com:9083") \
# Где хранить managed таблицы (HDFS путь)
.config("spark.sql.warehouse.dir", "hdfs://cluster/user/hive/warehouse") \
# Порог для CBO: Spark запрашивает статистику колонок из HMS
# при планировании join'ов для определения стратегии (broadcast vs merge)
.config("spark.sql.statistics.fallBackToHdfs", "true") \
.getOrCreate()
# Теперь Spark видит все таблицы, зарегистрированные в HMS
spark.sql("SHOW DATABASES").show()
# +------------+
# | namespace|
# +------------+
# | default|
# | sales|
# | analytics|
# +------------+
hive-site.xml: конфигурационный файл для Spark¶
hive-site.xml - это XML файл Hive конфигурации, который Spark автоматически читает из classpath при включении enableHiveSupport(). Это основной механизм настройки интеграции.
# Spark ищет hive-site.xml в следующем порядке:
# 1. $SPARK_HOME/conf/hive-site.xml
# 2. Любая директория в classpath через --files
# 3. Программно через .config() в SparkSession.builder
# На production кластере: обычно уже в $SPARK_HOME/conf/
ls /opt/spark/conf/hive-site.xml
# Для передачи через spark-submit:
spark-submit \
--files /etc/hive/conf/hive-site.xml \
--conf "spark.hadoop.hive.metastore.uris=thrift://hms:9083" \
your_job.py
<!-- /opt/spark/conf/hive-site.xml - полная production конфигурация -->
<!-- Адрес HMS Thrift сервера -->
<property>
<name>hive.metastore.uris</name>
<value>thrift://hms1.example.com:9083,thrift://hms2.example.com:9083</value>
<description>Список HMS серверов для failover. Spark пробует их по очереди.</description>
</property>
<!-- Корневая директория warehouse -->
<property>
<name>hive.metastore.warehouse.dir</name>
<value>hdfs://mycluster/user/hive/warehouse</value>
<description>Где HMS создаёт managed таблицы. Должно совпадать с spark.sql.warehouse.dir.</description>
</property>
<!-- Таймаут соединения с HMS (для нестабильных сетей) -->
<property>
<name>hive.metastore.client.socket.timeout</name>
<value>120</value>
<description>Секунды ожидания ответа от HMS. Дефолт 600 сек - слишком много.
При Partition Explosion запрос к HMS может занимать десятки секунд.</description>
</property>
<!-- Число retry при сбое HMS -->
<property>
<name>hive.metastore.failure.retries</name>
<value>3</value>
</property>
<!-- Кеширование метаданных на клиенте (Spark Driver) -->
<property>
<name>hive.metastore.client.cache.enabled</name>
<value>true</value>
<description>Кешировать metadata в памяти Spark Driver.
Снижает число Thrift-запросов при повторных обращениях к HMS.</description>
</property>
<!-- Динамическое партиционирование: разрешить запись в любое число партиций -->
<property>
<name>hive.exec.dynamic.partition.mode</name>
<value>nonstrict</value>
<description>strict требует хотя бы одну статическую партицию. nonstrict - нет ограничений.
Для Spark обычно нужен nonstrict.</description>
</property>
4. Управление каталогом через spark.catalog API¶
spark.catalog - это программный интерфейс Spark для работы с каталогом метаданных. Он предоставляет более безопасный и типизированный способ работы с метаданными по сравнению с сырыми DDL-строками через spark.sql().
Аудит окружения: что есть в каталоге¶
from pyspark.sql import SparkSession
spark = SparkSession.builder.enableHiveSupport() \
.config("hive.metastore.uris", "thrift://hms:9083") \
.getOrCreate()
# ── Просмотр баз данных ───────────────────────────────────────────────
# listDatabases() возвращает DataFrame с колонками: name, description, locationUri
databases = spark.catalog.listDatabases()
databases.show(truncate=False)
# +--------+--------------------+-----------------------------------------------+
# | name| description| locationUri|
# +--------+--------------------+-----------------------------------------------+
# | default| Default Hive db |hdfs://cluster/user/hive/warehouse |
# | sales| Sales analytics |hdfs://cluster/user/hive/warehouse/sales.db |
# | events| Event stream db |hdfs://cluster/data/events |
# +--------+--------------------+-----------------------------------------------+
# Переключить текущую базу данных
spark.catalog.setCurrentDatabase("sales")
# Проверить текущую БД
print(spark.catalog.currentDatabase())
# sales
# ── Просмотр таблиц ───────────────────────────────────────────────────
# listTables() возвращает DataFrame: name, database, description, tableType, isTemporary
tables = spark.catalog.listTables("sales")
tables.show(truncate=False)
# +-----------+--------+--------------------+---------+-----------+
# | name|database| description|tableType|isTemporary|
# +-----------+--------+--------------------+---------+-----------+
# | orders| sales| Customer orders | EXTERNAL| false|
# | products| sales| Product catalog| MANAGED| false|
# | sessions| null|Temp session stats | VIEW| true|
# +-----------+--------+--------------------+---------+-----------+
# tableType: MANAGED (данные управляются HMS) / EXTERNAL (данные внешние)
# isTemporary: true для temp views
# ── Просмотр колонок ────────────────────────────────────────────────
# listColumns() возвращает DataFrame: name, description, dataType, nullable, isPartition, isBucket
columns = spark.catalog.listColumns("orders", "sales")
columns.show(truncate=False)
# +-----------+-----------+---------+--------+-----------+--------+
# | name|description| dataType|nullable|isPartition|isBucket|
# +-----------+-----------+---------+--------+-----------+--------+
# | order_id| | string| true| false| false|
# | user_id| | bigint| true| false| false|
# | amount| | double| true| false| false|
# | date| | string| true| true| false|
# +-----------+-----------+---------+--------+-----------+--------+
# isPartition=true для date означает что это ключ партиционирования
# ── Проверка существования объектов ──────────────────────────────────
print(spark.catalog.tableExists("sales.orders")) # True
print(spark.catalog.tableExists("sales.unknown")) # False
# databaseExists - появился в Spark 3.3+
print(spark.catalog.databaseExists("analytics")) # False
Управление кешем: refreshTable и refreshByPath¶
Кеш метаданных в Spark Driver - одна из самых важных (и часто непонятых) концепций HMS. Spark кеширует метаданные таблиц и партиций локально, чтобы не ходить в HMS при каждом запросе. Но если данные изменились на HDFS (другой ETL добавил партицию), Spark «не знает» об этом.
# ── refreshTable: полное обновление метаданных таблицы ────────────────
# Сбрасывает кеш схемы, статистики и списка партиций
# Используется когда данные изменились внешними процессами
spark.catalog.refreshTable("sales.orders")
# После этого следующий запрос к sales.orders запросит свежие метаданные из HMS
# ── refreshByPath: обновление по физическому пути ─────────────────────
# Полезно когда не знаешь имя таблицы, но знаешь HDFS-путь
spark.catalog.refreshByPath("hdfs://cluster/data/events/date=2024-01-16/")
# ── cacheTable / uncacheTable: кеширование DATA в памяти Executor'ов ──
# Это НЕ кеш метаданных - это кеш самих данных в RAM!
# Полезно для таблиц, которые читаются многократно в одной Job'е
spark.catalog.cacheTable("sales.orders") # данные в RAM Executor'ов
spark.catalog.isCached("sales.orders") # True
# Выгрузить данные из RAM
spark.catalog.uncacheTable("sales.orders")
# Очистить весь кеш данных
spark.catalog.clearCache()
Полноценный аудит каталога: скрипт для Data Quality¶
def audit_hms_catalog(spark: SparkSession) -> dict:
"""
Полный аудит HMS каталога: базы, таблицы, колонки, размеры.
Полезно для Documentation-as-Code, data quality checks,
обнаружения таблиц без данных или с устаревшими метаданными.
"""
from pyspark.sql import functions as F
report = {}
# Собираем все базы
dbs = [row.name for row in spark.catalog.listDatabases().collect()]
report["databases"] = dbs
all_tables = []
for db in dbs:
tables = spark.catalog.listTables(db).collect()
for t in tables:
if t.isTemporary:
continue
table_info = {
"database": db,
"table": t.name,
"type": t.tableType,
}
# Считаем колонки
try:
cols = spark.catalog.listColumns(t.name, db).collect()
table_info["column_count"] = len(cols)
table_info["partition_keys"] = [
c.name for c in cols if c.isPartition
]
except Exception as e:
table_info["error"] = str(e)
all_tables.append(table_info)
continue
# Для external таблиц - проверяем доступность пути
try:
detail = spark.sql(f"DESCRIBE EXTENDED {db}.{t.name}")
location = detail.filter(
F.col("col_name") == "Location"
).select("data_type").first()
if location:
table_info["location"] = location[0]
except Exception:
pass
all_tables.append(table_info)
report["tables"] = all_tables
report["total_tables"] = len(all_tables)
report["managed_tables"] = sum(1 for t in all_tables if t.get("type") == "MANAGED")
report["external_tables"] = sum(1 for t in all_tables if t.get("type") == "EXTERNAL")
print(f"\nHMS Catalog Audit Report")
print(f"{'='*40}")
print(f"Databases: {len(dbs)}")
print(f"Tables (non-temp): {report['total_tables']}")
print(f" Managed: {report['managed_tables']}")
print(f" External: {report['external_tables']}")
return report
5. Физика таблиц: Managed vs External¶
Разница между Managed и External таблицами - одна из первых вещей, которую должен понять каждый Data Engineer, работающий с HMS. Ошибка в этом месте может привести к необратимой потере данных.
Managed (Internal) таблицы: HMS управляет данными¶
Когда Spark создаёт Managed таблицу, HMS берёт на себя полный контроль над жизненным циклом данных. Файлы создаются внутри spark.sql.warehouse.dir (или hive.metastore.warehouse.dir).
# Создание Managed таблицы
spark.sql("CREATE DATABASE IF NOT EXISTS demo")
spark.sql("USE demo")
# Вариант 1: через Spark SQL DDL
spark.sql("""
CREATE TABLE IF NOT EXISTS demo.user_events (
user_id BIGINT,
event_type STRING,
amount DOUBLE,
event_date DATE
)
PARTITIONED BY (event_date)
STORED AS PARQUET
""")
# Вариант 2: через DataFrame API (saveAsTable создаёт Managed по умолчанию)
from pyspark.sql import functions as F
df = spark.range(1000).select(
F.col("id").alias("user_id"),
F.lit("purchase").alias("event_type"),
(F.rand() * 100).alias("amount"),
F.lit("2024-01-15").cast("date").alias("event_date"),
)
df.write.mode("overwrite").saveAsTable("demo.user_events")
# Проверяем где физически лежат данные
spark.sql("DESCRIBE EXTENDED demo.user_events").show(50, truncate=False)
# ...
# | Location | hdfs://cluster/user/hive/warehouse/demo.db/user_events |
# | Type | MANAGED |
# ...
# ВНИМАНИЕ: DROP TABLE на Managed таблице удаляет файлы с HDFS!
spark.sql("DROP TABLE demo.user_events")
# Метаданные удалены из PostgreSQL И файлы удалены с HDFS
# НЕОБРАТИМО без backup!
External таблицы: данные живут независимо от HMS¶
External таблицы - рекомендуемый паттерн для production Data Lake. Spark регистрирует метаданные таблицы в HMS, но физические данные остаются под вашим контролем.
# Создание External таблицы с явным указанием пути
# Сначала кладём данные на HDFS независимо от HMS:
df.write \
.mode("overwrite") \
.partitionBy("event_date") \
.parquet("hdfs://cluster/data/warehouse/demo/user_events/")
# Регистрируем External таблицу поверх этих данных:
spark.sql("""
CREATE TABLE IF NOT EXISTS demo.user_events_ext (
user_id BIGINT,
event_type STRING,
amount DOUBLE
)
PARTITIONED BY (event_date STRING)
STORED AS PARQUET
LOCATION 'hdfs://cluster/data/warehouse/demo/user_events/'
-- LOCATION указывает что таблица EXTERNAL
""")
# Проверяем
spark.sql("DESCRIBE EXTENDED demo.user_events_ext") \
.filter("col_name IN ('Location', 'Type')") \
.show(truncate=False)
# | Location | hdfs://cluster/data/warehouse/demo/user_events/ |
# | Type | EXTERNAL |
# DROP TABLE на External таблице: удаляются ТОЛЬКО метаданные!
spark.sql("DROP TABLE demo.user_events_ext")
# Файлы на HDFS ОСТАЮТСЯ нетронутыми
# hdfs dfs -ls hdfs://cluster/data/warehouse/demo/user_events/ покажет файлы
# После DROP можно пересоздать таблицу с теми же данными
spark.sql("""
CREATE TABLE demo.user_events_ext
USING PARQUET
PARTITIONED BY (event_date)
LOCATION 'hdfs://cluster/data/warehouse/demo/user_events/'
""")
Проблема синхронизации партиций: MSCK REPAIR TABLE¶
Когда файлы на HDFS добавляются в обход HMS (например, через hdfs dfs -put или другим ETL-инструментом без Spark), HMS не знает о новых партициях. Это одна из самых распространённых проблем в production Data Lake.
# Сценарий: новые данные были загружены через hdfs dfs напрямую
hdfs dfs -mkdir -p /data/warehouse/demo/user_events/event_date=2024-01-20/
hdfs dfs -put local_data.parquet \
/data/warehouse/demo/user_events/event_date=2024-01-20/
# Spark НЕ видит эти данные:
spark.sql("SELECT * FROM demo.user_events_ext WHERE event_date = '2024-01-20'").show()
# 0 строк - HMS не знает о партиции!
spark.sql("SHOW PARTITIONS demo.user_events_ext").show()
# Не показывает event_date=2024-01-20
Решение 1: MSCK REPAIR TABLE - сканирует HDFS и синхронизирует список партиций с HMS:
# MSCK REPAIR TABLE (MetaStore Check Repair Table)
# Сканирует LOCATION директорию таблицы
# Находит директории которых нет в HMS
# Добавляет их как партиции
spark.sql("MSCK REPAIR TABLE demo.user_events_ext")
# Вывод: Partitions not in metastore: user_events_ext:event_date=2024-01-20
# Repair: Added partition to metastore user_events_ext:event_date=2024-01-20
# Теперь Spark видит данные!
spark.sql("SELECT COUNT(*) FROM demo.user_events_ext WHERE event_date = '2024-01-20'").show()
# +--------+
# |count(1)|
# +--------+
# | 45678|
# +--------+
# MSCK REPAIR может быть медленным при большом числе партиций
# Для incremental добавления используйте ALTER TABLE ADD PARTITION:
spark.sql("""
ALTER TABLE demo.user_events_ext
ADD IF NOT EXISTS PARTITION (event_date='2024-01-20')
LOCATION 'hdfs://cluster/data/warehouse/demo/user_events/event_date=2024-01-20/'
""")
Решение 2: spark.catalog.refreshTable() - мягкое обновление кеша без полного сканирования:
# Быстрее чем MSCK REPAIR, но может не обнаружить все новые партиции
spark.catalog.refreshTable("demo.user_events_ext")
# Или через refreshByPath для конкретной директории
spark.catalog.refreshByPath(
"hdfs://cluster/data/warehouse/demo/user_events/event_date=2024-01-20/"
)
6. Оптимизация взаимодействия Spark и HMS¶
Правильная настройка взаимодействия между Spark и HMS может сократить время планирования запросов от десятков секунд до долей секунды.
Metadata Caching: снижение нагрузки на Thrift¶
Spark хранит метаданные таблиц в памяти Driver'а. Это означает, что при повторных запросах к одной таблице HMS не запрашивается повторно. Но при работе с таблицами с большим числом партиций кеш может занимать сотни MB памяти Driver'а.
spark = SparkSession.builder \
.enableHiveSupport() \
# Кеш метаданных на стороне HMS клиента (в Driver'е)
# true: HMS возвращает список партиций из кеша вместо PostgreSQL
# Снижает нагрузку на СУБД при повторных запросах
.config("hive.metastore.client.cache.enabled", "true") \
# Размер кеша (число объектов метаданных)
.config("hive.metastore.client.cache.maxSize", "5000") \
# TTL кеша в миллисекундах: 5 минут
# После истечения - свежий запрос к HMS
.config("hive.metastore.client.cache.expiry.seconds", "300") \
# Статистика колонок: использовать для CBO оптимизаций
# Spark использует statistics для выбора join strategy
.config("spark.sql.statistics.histogram.enabled", "true") \
# Выбор оптимальной стратегии join на основе стат. из HMS
.config("spark.sql.cbo.enabled", "true") \
.config("spark.sql.cbo.joinReorder.enabled", "true") \
.getOrCreate()
Partition Pruning: метаданные вместо filesystem scan¶
Metastore Partition Pruning - ключевая оптимизация, которая позволяет Spark фильтровать партиции через SQL-запрос к PostgreSQL (HMS) вместо обхода HDFS директорий.
# Без partition pruning (до Spark 2.4):
# Spark делал: LIST всей директории таблицы на HDFS
# Затем: фильтровал партиции по имени директории
# При 100k партиций: 100k HTTP-запросов к HDFS NameNode!
# С Metastore Partition Pruning (Spark 3+):
# Spark делает: SQL запрос к PostgreSQL через HMS
# "SELECT partition WHERE date > '2024-01-01'"
# Получает ТОЛЬКО нужные партиции - без NameNode листинга
spark = SparkSession.builder \
.config("spark.sql.hive.metastorePartitionPruning", "true") \
.enableHiveSupport() \
.getOrCreate()
# Демонстрация: видно в EXPLAIN EXTENDED
df = spark.table("sales.orders")
df.filter("date >= '2024-01-01'").explain(extended=True)
# В Physical Plan появляется:
# PartitionFilters: [isnotnull(date#0), (date#0 >= 2024-01-01)]
# ^ Этот фильтр применяется к PARTITION LIST из HMS, не к HDFS listing
Тюнинг динамического партиционирования¶
При записи через Spark с partitionBy() каждая запись идёт в HMS для регистрации новых партиций. Неправильные настройки могут вызвать перегрузку HMS:
spark = SparkSession.builder \
# Максимальное число динамических партиций за одну операцию INSERT
# Дефолт: 1000. При записи 5000 партиций - ошибка!
# Увеличьте для больших batch-операций:
.config("hive.exec.max.dynamic.partitions", "10000") \
# Максимальное число партиций на один DataNode
.config("hive.exec.max.dynamic.partitions.pernode", "5000") \
# Режим динамического партиционирования
# nonstrict: разрешает писать во все партиции без ограничений
# strict: требует хотя бы один статический ключ (безопаснее)
.config("hive.exec.dynamic.partition.mode", "nonstrict") \
# Spark автоматически обновляет статистику после записи
# Полезно для CBO оптимизатора
.config("spark.sql.statistics.autoUpdate.enabled", "true") \
.getOrCreate()
# При записи через Spark DataFrame API:
df.write \
.mode("overwrite") \
.format("parquet") \
# partitionBy регистрирует партиции в HMS автоматически
.partitionBy("event_date") \
.saveAsTable("demo.user_events")
# Spark автоматически регистрирует все созданные партиции в HMS
7. Практика: интеграция PySpark с HMS в Docker-окружении¶
Docker Compose стенд¶
# docker-compose.yml - минимальный HMS стенд
services:
# PostgreSQL: бэкенд для HMS
postgres:
image: postgres:15
environment:
POSTGRES_DB: hive_metastore
POSTGRES_USER: hive
POSTGRES_PASSWORD: hive_password
ports:
- "5432:5432"
volumes:
- postgres_data:/var/lib/postgresql/data
networks:
- data-platform
# Hive Metastore Service
metastore:
image: apache/hive:4.0.0
depends_on:
- postgres
environment:
SERVICE_NAME: metastore
DB_DRIVER: postgres
SERVICE_OPTS: >
-Djavax.jdo.option.ConnectionURL=jdbc:postgresql://postgres:5432/hive_metastore
-Djavax.jdo.option.ConnectionDriverName=org.postgresql.Driver
-Djavax.jdo.option.ConnectionUserName=hive
-Djavax.jdo.option.ConnectionPassword=hive_password
volumes:
- ./conf/hive-site.xml:/opt/hive/conf/hive-site.xml
ports:
- "9083:9083"
networks:
- data-platform
networks:
data-platform: {}
volumes:
postgres_data: {}
Лабораторные задания¶
Задание 1: In-Memory vs HMS - демонстрация персистентности
# lab_01_inmemory_vs_hms.py
from pyspark.sql import SparkSession
print("=== ДЕМОНСТРАЦИЯ: In-Memory Catalog ===")
# Сессия без HMS
spark_local = SparkSession.builder \
.master("local[2]") \
.appName("lab-inmemory") \
.getOrCreate()
# Создаём таблицу
spark_local.range(100) \
.write.mode("overwrite") \
.saveAsTable("my_temp_table")
print("Таблицы в локальном каталоге:")
spark_local.sql("SHOW TABLES").show()
spark_local.stop()
# Новая сессия - таблица исчезла
spark_local2 = SparkSession.builder \
.master("local[2]") \
.appName("lab-inmemory-2") \
.getOrCreate()
print("После пересоздания сессии:")
spark_local2.sql("SHOW TABLES").show()
# Пустой результат!
spark_local2.stop()
print()
print("=== ДЕМОНСТРАЦИЯ: HMS Catalog ===")
# Сессия с HMS
spark_hms = SparkSession.builder \
.master("local[2]") \
.appName("lab-hms") \
.config("hive.metastore.uris", "thrift://localhost:9083") \
.config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
.enableHiveSupport() \
.getOrCreate()
spark_hms.sql("CREATE DATABASE IF NOT EXISTS lab")
spark_hms.sql("USE lab")
spark_hms.range(100) \
.write.mode("overwrite") \
.saveAsTable("persistent_table")
print("Таблицы в HMS:")
spark_hms.sql("SHOW TABLES IN lab").show()
spark_hms.stop()
# Новая сессия - таблица СОХРАНИЛАСЬ!
spark_hms2 = SparkSession.builder \
.master("local[2]") \
.appName("lab-hms-2") \
.config("hive.metastore.uris", "thrift://localhost:9083") \
.enableHiveSupport() \
.getOrCreate()
print("После пересоздания сессии - таблица сохранилась:")
spark_hms2.sql("SHOW TABLES IN lab").show()
# persistent_table видна!
spark_hms2.stop()
Задание 2: Managed vs External - проверка поведения DROP TABLE
# lab_02_managed_vs_external.py
import subprocess
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder \
.master("local[2]") \
.appName("lab-tables") \
.config("hive.metastore.uris", "thrift://localhost:9083") \
.config("spark.sql.warehouse.dir", "file:///tmp/spark-warehouse") \
.enableHiveSupport() \
.getOrCreate()
spark.sql("CREATE DATABASE IF NOT EXISTS lab_tables")
# Тестовые данные
df = spark.range(1000).select(
F.col("id"),
(F.rand() * 100).alias("value"),
F.date_add(F.lit("2024-01-01"), (F.rand() * 30).cast("int")).alias("dt"),
)
# ── MANAGED TABLE ───────────────────────────────────────────────────
print("=== MANAGED TABLE ===")
df.write.mode("overwrite") \
.partitionBy("dt") \
.saveAsTable("lab_tables.managed_demo")
# Смотрим расположение (HMS контролирует)
spark.sql("DESCRIBE EXTENDED lab_tables.managed_demo") \
.filter("col_name = 'Location'").show(1, truncate=False)
# Путь внутри warehouse.dir
managed_path = spark.sql("DESCRIBE EXTENDED lab_tables.managed_demo") \
.filter("col_name = 'Location'") \
.select("data_type").first()[0]
print(f"Данные Managed table в: {managed_path}")
# DROP MANAGED TABLE
spark.sql("DROP TABLE lab_tables.managed_demo")
print("После DROP TABLE managed:")
# Проверяем что данные удалены
result = subprocess.run(
["ls", managed_path.replace("file://", "")],
capture_output=True, text=True
)
print(f" ls {managed_path}: {result.stderr or 'No such file (удалено!)'}")
# ── EXTERNAL TABLE ──────────────────────────────────────────────────
print("\n=== EXTERNAL TABLE ===")
external_path = "/tmp/lab-external-data"
df.write.mode("overwrite") \
.partitionBy("dt") \
.parquet(external_path)
spark.sql(f"""
CREATE TABLE lab_tables.external_demo (
id BIGINT,
value DOUBLE
)
PARTITIONED BY (dt DATE)
STORED AS PARQUET
LOCATION 'file://{external_path}'
""")
# Регистрируем партиции
spark.sql("MSCK REPAIR TABLE lab_tables.external_demo")
count_before = spark.sql("SELECT COUNT(*) FROM lab_tables.external_demo").first()[0]
print(f"Строк до DROP: {count_before:,}")
# DROP EXTERNAL TABLE
spark.sql("DROP TABLE lab_tables.external_demo")
print("После DROP TABLE external:")
# Данные ОСТАЛИСЬ!
result = subprocess.run(
["ls", external_path],
capture_output=True, text=True
)
print(f" ls {external_path}: {result.stdout or result.stderr}")
print(" Данные НЕ удалены (только метаданные удалены из HMS)!")
# Можем пересоздать таблицу с теми же данными
spark.sql(f"""
CREATE TABLE lab_tables.external_demo
USING PARQUET
PARTITIONED BY (dt)
LOCATION 'file://{external_path}'
""")
spark.sql("MSCK REPAIR TABLE lab_tables.external_demo")
count_after = spark.sql("SELECT COUNT(*) FROM lab_tables.external_demo").first()[0]
print(f"Строк после пересоздания: {count_after:,} (те же данные!)")
spark.stop()
Задание 3: Траблшутинг - рассинхронизация каталога и MSCK REPAIR
# lab_03_msck_repair.py
import subprocess
import os
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder \
.master("local[2]") \
.appName("lab-msck") \
.config("hive.metastore.uris", "thrift://localhost:9083") \
.enableHiveSupport() \
.getOrCreate()
EXTERNAL_PATH = "/tmp/lab-repair-demo"
# Создаём таблицу с несколькими партициями
for date in ["2024-01-10", "2024-01-11", "2024-01-12"]:
df = spark.range(100).select(
F.col("id"),
F.lit(date).alias("event_date")
)
df.coalesce(1).write \
.mode("overwrite") \
.parquet(f"{EXTERNAL_PATH}/event_date={date}/")
spark.sql("DROP TABLE IF EXISTS lab.repair_demo")
spark.sql(f"""
CREATE TABLE lab.repair_demo (id BIGINT)
PARTITIONED BY (event_date STRING)
STORED AS PARQUET
LOCATION 'file://{EXTERNAL_PATH}'
""")
# Регистрируем начальные партиции
spark.sql("MSCK REPAIR TABLE lab.repair_demo")
print("Начальные партиции:")
spark.sql("SHOW PARTITIONS lab.repair_demo").show()
print("Количество строк:")
spark.sql("SELECT COUNT(*) FROM lab.repair_demo").show()
# ── СИМУЛИРУЕМ ДОБАВЛЕНИЕ ДАННЫХ В ОБХОД HMS ────────────────────────
print("\n=== Добавляем партицию в обход HMS ===")
new_date = "2024-01-13"
new_path = f"{EXTERNAL_PATH}/event_date={new_date}"
os.makedirs(new_path, exist_ok=True)
# Создаём parquet файл напрямую (без Spark saveAsTable)
df_new = spark.range(50).select(
F.col("id"),
F.lit(new_date).alias("event_date")
)
df_new.coalesce(1).write.mode("overwrite").parquet(new_path)
print(f"Файлы добавлены в {new_path} напрямую (без HMS)")
# Spark НЕ видит новые данные
print("\nДО refreshTable / MSCK:")
print("Количество строк (новая партиция не видна!):")
spark.sql("SELECT COUNT(*) FROM lab.repair_demo").show()
# 300 строк (только 3 старые партиции)
# ── РЕШЕНИЕ 1: refreshTable ────────────────────────────────────────
print("\n=== Решение 1: spark.catalog.refreshTable() ===")
spark.catalog.refreshTable("lab.repair_demo")
# Может не обнаружить все партиции без MSCK!
print("После refreshTable:")
spark.sql("SHOW PARTITIONS lab.repair_demo").show()
spark.sql("SELECT COUNT(*) FROM lab.repair_demo").show()
# ── РЕШЕНИЕ 2: MSCK REPAIR TABLE ──────────────────────────────────
print("\n=== Решение 2: MSCK REPAIR TABLE ===")
spark.sql("MSCK REPAIR TABLE lab.repair_demo")
print("После MSCK REPAIR:")
spark.sql("SHOW PARTITIONS lab.repair_demo").show()
spark.sql("SELECT COUNT(*) FROM lab.repair_demo").show()
# 350 строк - новая партиция найдена!
print("\n✅ Рассинхронизация устранена через MSCK REPAIR TABLE")
spark.stop()
Смотрим на HMS изнутри: PostgreSQL schema¶
Для глубокого понимания HMS полезно посмотреть на физическую структуру базы данных PostgreSQL, где хранятся все метаданные:
-- Подключаемся к PostgreSQL HMS напрямую
-- psql -h localhost -U hive -d hive_metastore
-- Таблицы в HMS (DBS → TBLS → SDS → COLUMNS_V2)
-- Список баз данных
SELECT DB_ID, NAME, LOCATION_URI FROM DBS;
-- Список таблиц
SELECT T.TBL_ID, T.TBL_NAME, T.TBL_TYPE,
D.NAME as DB_NAME, T.CREATE_TIME
FROM TBLS T
JOIN DBS D ON T.DB_ID = D.DB_ID
ORDER BY T.CREATE_TIME DESC
LIMIT 20;
-- Физические пути и форматы хранения (StorageDescriptor)
SELECT T.TBL_NAME, S.LOCATION, S.INPUT_FORMAT,
S.OUTPUT_FORMAT, S.IS_COMPRESSED
FROM TBLS T
JOIN SDS S ON T.SD_ID = S.SD_ID
WHERE T.TBL_NAME = 'user_events';
-- Список партиций конкретной таблицы
SELECT P.PART_ID, P.PART_NAME, P.CREATE_TIME,
PS.LOCATION
FROM PARTITIONS P
JOIN SDS PS ON P.SD_ID = PS.SD_ID
JOIN TBLS T ON P.TBL_ID = T.TBL_ID
WHERE T.TBL_NAME = 'user_events'
ORDER BY P.CREATE_TIME DESC
LIMIT 50;
-- Типичный вывод:
-- PART_ID | PART_NAME | LOCATION
-- --------|---------------------|--------------------------------------------
-- 1 | event_date=2024-01-15 | hdfs://cluster/data/.../event_date=2024-01-15
-- 2 | event_date=2024-01-16 | hdfs://cluster/data/.../event_date=2024-01-16
-- Статистика таблицы (для CBO)
SELECT PARAM_KEY, PARAM_VALUE
FROM TABLE_PARAMS
WHERE TBL_ID = (SELECT TBL_ID FROM TBLS WHERE TBL_NAME = 'user_events')
ORDER BY PARAM_KEY;
-- Типичный вывод:
-- numFiles | 50
-- numRows | 5000000
-- rawDataSize | 134217728
-- totalSize | 67108864
-- transient_lastDdlTime | 1705334400
Итоги: HMS как центральная нервная система Data Lake¶
Hive Metastore - это не просто «служебная утилита». Это центральная нервная система всего on-premise Data Lake. Каждый SQL-запрос в Spark, каждое обращение Trino, каждое создание таблицы Iceberg проходит через HMS.
Что должен знать Data Engineer:
- Архитектура: HMS = Thrift Server + PostgreSQL + HMS API. Клиенты (Spark, Trino) общаются через Thrift RPC, никогда напрямую с PostgreSQL.
- enableHiveSupport() переключает Spark с in-memory InMemoryCatalog на HiveExternalCatalog. Без него таблицы исчезают при перезапуске.
- Partition Explosion - главный killer производительности. 100k+ партиций = секунды только на metadata-запрос к HMS. Решения: partition pruning, правильная партиционная схема, избегать мультиуровневых high-cardinality разбивок.
- Managed vs External: Managed - данные удаляются с DROP TABLE. External - только метаданные. В production Data Lake почти всегда используйте External.
- MSCK REPAIR TABLE - обязательный инструмент при добавлении данных в обход HMS. Синхронизирует список партиций между HDFS и PostgreSQL.
- spark.catalog API - предпочтительнее сырых
spark.sql("DDL...")для автоматизации. Типобезопасный интерфейс для listDatabases, listTables, refreshTable.