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, тюнинг производительности и траблшутинг.

storage platform

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.