Hive DDL: STORED AS, ROW FORMAT SERDE, MSCK REPAIR TABLE и StorageDescriptor в метасторе

Глубокий разбор Hive DDL для on-premise Data Lake: анатомия StorageDescriptor в PostgreSQL HMS, трансляция STORED AS в Input/OutputFormat классы, кастомный парсинг через ROW FORMAT SERDE и RegexSerDe, восстановление рассинхронизации через MSCK REPAIR TABLE, Schema Evolution с ALTER TABLE и лабораторный аудит метастора изнутри.

storage platform

1. Архитектурный контекст: Hive DDL как стандарт описания данных

Hive DDL сегодня - это не язык Apache Hive как вычислительного движка. Это декларативный стандарт описания структуры данных в on-premise Data Lake. Когда мы пишем CREATE TABLE ... STORED AS PARQUET LOCATION 'hdfs://...', мы создаём не таблицу в реляционной БД - мы создаём запись в Hive Metastore (HMS), которая описывает физическое расположение файлов и способ их интерпретации.

Это разделение критически важно для понимания всей архитектуры:

Схема показывает ключевой принцип: DDL работает с HMS (слой метаданных), а не напрямую с HDFS. Spark SQL при выполнении запроса сначала обращается в HMS за метаданными, а затем читает данные из HDFS напрямую, минуя HMS.

Зачем Spark понимает Hive DDL

Apache Spark 2.x и выше включает полноценный Hive DDL-парсер. Это означает:

  • Команды CREATE TABLE, ALTER TABLE, DROP TABLE в Spark SQL транслируются напрямую в HMS через Thrift API
  • Spark Catalog API (spark.catalog.createTable(), spark.catalog.listTables()) под капотом использует те же Thrift-вызовы
  • Таблицы, созданные через Spark, видны в Hive CLI, и наоборот

Смена движка выполнения без смены метаданных. Это важное свойство: если ваши DDL-скрипты написаны под Hive, они работают в Spark SQL без изменений. Можно мигрировать с Hive на Spark только на уровне spark.master = yarn без переписывания DDL.

Почему важно понимать физику HMS изнутри

На практике Data Engineer регулярно сталкивается с проблемами:

  • Table not found при правильно написанном SQL - таблица есть в HDFS, но не зарегистрирована в HMS
  • SerDeException: null key при чтении данных - у SerDe неверный разделитель
  • ClassNotFoundException: org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe - JAR с SerDe недоступен в classpath Spark
  • Новые партиции есть на HDFS, но SELECT WHERE date='2024-01-15' возвращает 0 строк

Все эти проблемы решаются только если понимать, как HMS хранит метаданные на уровне PostgreSQL и как Spark интерпретирует StorageDescriptor.


2. Под капотом HMS: анатомия StorageDescriptor

Hive Metastore использует PostgreSQL (или MySQL) как реляционное хранилище для метаданных. Давайте разберём физическую структуру этой БД - это фундамент понимания всего Hive DDL.

Основные таблицы в схеме HMS

Схема раскрывает важную деталь: каждая партиция имеет свою собственную StorageDescriptor. Это означает, что партиции одной таблицы могут физически лежать в разных директориях HDFS и даже в разных форматах файлов. Это мощная, но и опасная возможность.

Разбор StorageDescriptor (SDS): физическое сердце таблицы

StorageDescriptor - самый важный объект HMS. Он описывает:

-- Содержимое таблицы SDS в PostgreSQL HMS

SELECT
    s.SD_ID,
    s.LOCATION,
    s.INPUT_FORMAT,
    s.OUTPUT_FORMAT,
    s.IS_COMPRESSED,
    s.NUM_BUCKETS,
    ser.SLIB as serde_class
FROM SDS s
JOIN SERDES ser ON s.SERDE_ID = ser.SERDE_ID
WHERE s.SD_ID IN (
    SELECT SD_ID FROM TBLS WHERE TBL_NAME = 'user_events'
);

-- Пример результата для Parquet таблицы:
-- SD_ID | LOCATION                                    | INPUT_FORMAT
-- ------+---------------------------------------------+----------------------------------------
-- 42    | hdfs://cluster/data/silver/user_events      | org.apache.hadoop.hive.ql.io.parquet
--       |                                             |   .MapredParquetInputFormat
--
-- OUTPUT_FORMAT                                           | SERDE_CLASS
-- -------------------------------------------------------+------------------------------------------
-- org.apache.hadoop.hive.ql.io.parquet.MapredParquetOut  | org.apache.hadoop.hive.ql.io.parquet
-- putFormat                                               |   .serde.ParquetHiveSerDe

Поле LOCATION - физический URI директории данных на HDFS. Это НЕ путь к конкретному файлу, а к директории. Spark читает все файлы из этой директории (и поддиректорий партиций).

Поле INPUT_FORMAT - Java-класс, реализующий InputFormat interface. Этот класс отвечает за разбивку файлов на InputSplit'ы (которые становятся Spark-партициями) и создание RecordReader'ов для чтения каждого сплита. Именно здесь определяется, умеет ли система читать Parquet, ORC, CSV или бинарный формат.

Поле OUTPUT_FORMAT - Java-класс для записи данных. При INSERT INTO table SELECT... Spark использует OutputFormat для создания файлов на HDFS.

Поле SLIB в SERDES (SerDe Library) - Java-класс, реализующий SerDe interface. Отвечает за десериализацию: преобразование байт из файла (прочитанных через RecordReader) в Java-объект строки таблицы. В отличие от InputFormat (который читает блоки байт), SerDe понимает структуру данных внутри файла.

Параметры SerDe в SERDE_PARAMS

-- Параметры SerDe для CSV-таблицы
SELECT param_key, param_value
FROM SERDE_PARAMS
WHERE SERDE_ID = (
    SELECT SERDE_ID FROM SERDES
    WHERE SLIB = 'org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe'
);

-- Типичный результат:
-- param_key                        | param_value
-- ---------------------------------+-------------
-- field.delim                      | ,
-- serialization.format             | ,
-- line.delim                       | \n
-- collection.delim                 | \002
-- mapkey.delim                     | \003
-- serialization.null.format        | \N
-- escape.delim                     | \

Эти параметры - именно то, что вы указываете в ROW FORMAT DELIMITED FIELDS TERMINATED BY ','. Они физически хранятся в PostgreSQL и читаются при каждом открытии файла Spark'ом.

Почему рассинхронизация StorageDescriptor - катастрофа

Представьте: вы вручную изменили разделитель в SERDE_PARAMS.field.delim с , на \t через psql, но файлы по-прежнему содержат CSV с запятыми. Spark прочитает файл, попытается разобрать строки по \t, не найдёт разделитель и вернёт весь текст строки как одно значение в первую колонку. Остальные колонки будут NULL. Никаких ошибок - только молчаливо неправильные данные.

Именно поэтому прямые UPDATE/INSERT в PostgreSQL HMS категорически запрещены. Все изменения только через официальный HMS API: Spark SQL DDL или Hive CLI.


3. Механика STORED AS: управление форматами на уровне I/O

Инструкция STORED AS - это декларация того, какие Java-классы использовать для чтения и записи файлов в HDFS. Разберём её механику подробно.

Трансляция STORED AS в Input/OutputFormat классы

Понимание этой схемы объясняет, почему STORED AS TEXTFILE работает медленно на больших данных: TextInputFormat создаёт один InputSplit на строку файла (по умолчанию), что создаёт тысячи мелких Spark-задач при чтении большого CSV. MapredParquetInputFormat, напротив, создаёт сплиты по Row Group (~128 MB), что идеально совпадает с блоком HDFS.

Полный DDL: все варианты STORED AS

-- ── PARQUET ─────────────────────────────────────────────────────────
-- Основной формат для Silver/Gold слоёв Data Lake.
-- Колоночное хранение, поддержка сложных типов (ARRAY, MAP, STRUCT).
CREATE EXTERNAL TABLE IF NOT EXISTS silver.user_events (
    event_id    STRING        COMMENT 'Уникальный идентификатор события',
    user_id     BIGINT        COMMENT 'Идентификатор пользователя',
    event_type  STRING        COMMENT 'Тип события: purchase, view, click',
    amount      DECIMAL(18,2) COMMENT 'Сумма транзакции',
    metadata    MAP<STRING, STRING> COMMENT 'Дополнительные атрибуты'
)
COMMENT 'Очищенные события пользователей (Silver слой)'
PARTITIONED BY (event_date DATE COMMENT 'Дата события для партиционирования')
STORED AS PARQUET
LOCATION 'hdfs://cluster/data/silver/user_events'
TBLPROPERTIES (
    'created_by' = 'data-platform-team',
    'data_quality' = 'validated',
    'compression' = 'snappy'
);

-- ── ORC ──────────────────────────────────────────────────────────────
-- Оптимизирован для Hive ACID и Trino.
-- Встроенные Bloom Filters, более компактные числовые типы.
CREATE EXTERNAL TABLE IF NOT EXISTS silver.transactions (
    txn_id      BIGINT,
    merchant_id INT,
    amount      DOUBLE,
    currency    STRING
)
PARTITIONED BY (txn_date STRING)
STORED AS ORC
LOCATION 'hdfs://cluster/data/silver/transactions'
TBLPROPERTIES (
    'orc.compress' = 'ZLIB',         -- Сжатие на уровне ORC
    'orc.stripe.size' = '67108864'   -- Размер Stripe = 64 MB
);

-- ── AVRO ─────────────────────────────────────────────────────────────
-- Для Bronze слоя (CDC, Kafka consumer output).
-- Схема хранится внутри файла - Schema Evolution из коробки.
CREATE EXTERNAL TABLE IF NOT EXISTS bronze.cdc_orders (
    op          STRING COMMENT 'Операция: c=create, u=update, d=delete',
    before      STRUCT<id:BIGINT, status:STRING, amount:DOUBLE>,
    after       STRUCT<id:BIGINT, status:STRING, amount:DOUBLE>,
    source      STRUCT<db:STRING, table:STRING, ts_ms:BIGINT>
)
PARTITIONED BY (ingestion_date STRING)
STORED AS AVRO
LOCATION 'hdfs://cluster/data/bronze/cdc_orders'
TBLPROPERTIES (
    'avro.schema.url' = 'hdfs://cluster/schemas/cdc_orders.avsc'
);
-- avro.schema.url позволяет использовать внешнюю Avro-схему
-- вместо указания схемы в DDL.

-- ── TEXTFILE с сжатием ───────────────────────────────────────────────
-- Только для staging или интеграции с внешними системами.
-- НЕ использовать для аналитики.
CREATE EXTERNAL TABLE IF NOT EXISTS staging.raw_csv_logs (
    raw_line STRING COMMENT 'Полная строка лога как одна колонка'
)
STORED AS TEXTFILE
LOCATION 'hdfs://cluster/data/staging/raw_logs'
TBLPROPERTIES ('skip.header.line.count' = '1');
-- skip.header.line.count: пропустить первую строку (заголовок CSV)

Миграция текстовой таблицы в Parquet через Spark SQL

Типичная задача в enterprise: есть legacy таблица в CSV/TextFile, нужно переконвертировать в Parquet без потери данных и с минимальным downtime.

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("csv-to-parquet-migration") \
    .config("hive.metastore.uris", "thrift://hms:9083") \
    .enableHiveSupport() \
    .getOrCreate()

# Шаг 1: Проверяем что исходная таблица видна
spark.sql("DESCRIBE EXTENDED legacy.sales_csv").show(50, truncate=False)
# Смотрим: Location, InputFormat (должен быть TextInputFormat)

# Шаг 2: Создаём новую Parquet таблицу (External!)
# Важно: указываем НОВЫЙ путь, не тот же что у CSV
spark.sql("""
    CREATE EXTERNAL TABLE IF NOT EXISTS silver.sales_parquet
    LIKE legacy.sales_csv  -- Копируем схему из исходной таблицы
    STORED AS PARQUET      -- Но меняем формат на Parquet
    LOCATION 'hdfs://cluster/data/silver/sales_parquet'
""")

# Шаг 3: Читаем CSV и пишем как Parquet партиция за партицией
# Это позволяет отслеживать прогресс и перезапускать при сбоях
available_dates = spark.sql(
    "SHOW PARTITIONS legacy.sales_csv"
).collect()

for row in available_dates:
    # row.partition = 'sale_date=2024-01-15'
    partition_value = row.partition.split('=')[1]

    df = spark.sql(f"""
        SELECT * FROM legacy.sales_csv
        WHERE sale_date = '{partition_value}'
    """)

    # Пишем напрямую на HDFS (не через saveAsTable во избежание
    # конфликтов при частичных записях)
    df.write \
        .mode("overwrite") \
        .parquet(f"hdfs://cluster/data/silver/sales_parquet/sale_date={partition_value}/")

    # Регистрируем партицию в HMS
    spark.sql(f"""
        ALTER TABLE silver.sales_parquet
        ADD IF NOT EXISTS PARTITION (sale_date='{partition_value}')
        LOCATION 'hdfs://cluster/data/silver/sales_parquet/sale_date={partition_value}/'
    """)
    print(f"Мигрирована партиция: {partition_value}")

# Шаг 4: Верификация
original_count = spark.sql("SELECT COUNT(*) FROM legacy.sales_csv").first()[0]
migrated_count = spark.sql("SELECT COUNT(*) FROM silver.sales_parquet").first()[0]
assert original_count == migrated_count, f"Несоответствие строк: {original_count} vs {migrated_count}"
print(f"Миграция завершена: {migrated_count:,} строк")

4. Кастомный парсинг: ROW FORMAT SERDE и разделители

ROW FORMAT - это инструкция, которая определяет, как SerDe должен интерпретировать содержимое файлов. Существует два варианта.

ROW FORMAT DELIMITED: простые текстовые файлы

ROW FORMAT DELIMITED настраивает LazySimpleSerDe - встроенный SerDe для простых текстовых файлов с разделителями.

-- Полный синтаксис ROW FORMAT DELIMITED
CREATE EXTERNAL TABLE staging.tab_separated_log (
    log_time   STRING,
    client_ip  STRING,
    method     STRING,
    uri        STRING,
    status     INT,
    bytes_sent INT,
    referer    STRING,
    user_agent STRING
)
ROW FORMAT DELIMITED
    FIELDS TERMINATED BY '\t'          -- Разделитель колонок: табуляция
    COLLECTION ITEMS TERMINATED BY ','  -- Разделитель элементов в ARRAY/LIST
    MAP KEYS TERMINATED BY '='          -- Разделитель ключ=значение в MAP
    LINES TERMINATED BY '\n'            -- Разделитель строк (обычно не нужен)
    NULL DEFINED AS 'NULL'              -- Строка, представляющая NULL
STORED AS TEXTFILE
LOCATION 'hdfs://cluster/data/staging/tab_logs';

-- Для CSV с кавычками (quote char) - LazySimpleSerDe НЕ поддерживает.
-- Используйте OpenCSVSerDe:
CREATE EXTERNAL TABLE staging.quoted_csv (
    id       INT,
    name     STRING,
    address  STRING  -- поле может содержать запятые внутри кавычек
)
ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.OpenCSVSerde'
WITH SERDEPROPERTIES (
    'separatorChar' = ',',
    'quoteChar'     = '"',
    'escapeChar'    = '\\'
)
STORED AS TEXTFILE
LOCATION 'hdfs://cluster/data/staging/quoted_csv';

Важное ограничение LazySimpleSerDe: этот SerDe «ленивый» - он не парсит строку до момента первого обращения к колонке. Это ускоряет запросы с Column Pruning, но означает, что ошибки парсинга (например, строка с неверным числом полей) обнаруживаются только при чтении конкретного файла, а не при создании таблицы.

ROW FORMAT SERDE: продвинутый парсинг сложных форматов

ROW FORMAT SERDE позволяет использовать произвольный Java-класс SerDe вместо встроенного LazySimpleSerDe. Это открывает возможности для парсинга форматов, которые невозможно описать простыми разделителями.

RegexSerDe: парсинг структурированных логов

org.apache.hadoop.hive.serde2.RegexSerDe - SerDe на основе регулярных выражений. Позволяет разбирать сложные неструктурированные логи прямо на уровне Hive DDL, без написания UDF.

-- Парсинг Combined Log Format (Nginx/Apache access log)
-- Пример строки лога:
-- 192.168.1.100 - john [15/Jan/2024:10:00:01 +0300] "GET /api/v1/users HTTP/1.1" 200 4521 "https://example.com" "Mozilla/5.0..."

CREATE EXTERNAL TABLE IF NOT EXISTS bronze.nginx_access_log (
    client_ip    STRING COMMENT 'IP-адрес клиента',
    ident        STRING COMMENT 'Идентификатор клиента (обычно -)',
    auth_user    STRING COMMENT 'Аутентифицированный пользователь',
    log_time     STRING COMMENT 'Время запроса: [DD/Mon/YYYY:HH:MM:SS +ZZZZ]',
    request      STRING COMMENT 'HTTP-запрос: METHOD URI PROTOCOL',
    status       INT    COMMENT 'HTTP статус-код',
    bytes_sent   BIGINT COMMENT 'Размер ответа в байтах',
    referer      STRING COMMENT 'Заголовок Referer',
    user_agent   STRING COMMENT 'Заголовок User-Agent'
)
ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.RegexSerDe'
WITH SERDEPROPERTIES (
    -- Регулярное выражение для Combined Log Format
    -- Каждая захватывающая группа () становится одной колонкой
    'input.regex' =
        '(\\S+)\\s+(\\S+)\\s+(\\S+)\\s+\\[(.+?)\\]\\s+"(.*?)"\\s+(\\d+)\\s+(\\d+|-)\\s+"(.*?)"\\s+"(.*?)"',
    -- Формат вывода (обычно %1$s %2$s... - позиционные ссылки)
    'output.format.string' = '%1$s %2$s %3$s [%4$s] "%5$s" %6$s %7$s "%8$s" "%9$s"'
)
STORED AS TEXTFILE
LOCATION 'hdfs://cluster/data/bronze/nginx_access_log'
TBLPROPERTIES (
    'comment' = 'Nginx Combined Log Format - Bronze Layer'
);

-- Тест: читаем несколько строк
-- SELECT client_ip, status, request FROM bronze.nginx_access_log LIMIT 5;

После создания этой таблицы Spark может читать сырые лог-файлы как структурированную таблицу - без написания Python-кода для парсинга.

# Использование через PySpark
df = spark.sql("""
    SELECT
        client_ip,
        status,
        REGEXP_EXTRACT(request, '^(GET|POST|PUT|DELETE)', 1) as method,
        REGEXP_EXTRACT(request, '^\\S+\\s+(\\S+)', 1) as uri,
        CAST(bytes_sent AS BIGINT) as bytes
    FROM bronze.nginx_access_log
    WHERE status BETWEEN 400 AND 599  -- Только ошибки клиента и сервера
    AND log_time LIKE '15/Jan/2024%'
""")

df.groupBy("status", "method") \
  .count() \
  .orderBy("count", ascending=False) \
  .show()

JsonSerDe: полуструктурированные JSON-файлы

Для JSON-файлов в on-premise контуре существует несколько SerDe. org.apache.hive.hcatalog.data.JsonSerDe - наиболее распространённый:

-- JSON SerDe для событий из внешних систем
-- Пример строки: {"event_id":"evt-001","user_id":12345,"payload":{"action":"buy","amount":99.9}}

CREATE EXTERNAL TABLE IF NOT EXISTS bronze.webhook_events (
    event_id    STRING,
    user_id     BIGINT,
    event_type  STRING,
    payload     STRUCT<
        action: STRING,
        amount: DOUBLE,
        items:  ARRAY<STRUCT<sku:STRING, qty:INT, price:DOUBLE>>
    >,
    timestamp   STRING
)
ROW FORMAT SERDE 'org.apache.hive.hcatalog.data.JsonSerDe'
WITH SERDEPROPERTIES (
    'ignore.malformed.json' = 'TRUE'  -- Пропускать битые JSON строки
)
STORED AS TEXTFILE  -- JSON-файлы хранятся как TextFile (NDJSON)
LOCATION 'hdfs://cluster/data/bronze/webhook_events'
TBLPROPERTIES (
    'serialization.encoding' = 'UTF-8'
);

-- ВАЖНО: JsonSerDe требует JSON Lines (NDJSON) формат:
-- каждая строка файла = отдельный JSON-объект
-- НЕ поддерживает многострочные JSON-массивы!

Ограничение JsonSerDe: он не поддерживает файлы-массивы [{...}, {...}]. Только JSON Lines (одна строка - один JSON-объект). Это стандарт для Event Streaming и Kafka output.


5. Ликвидация Metadata Gap: MSCK REPAIR TABLE

«Metadata Gap» - один из самых частых операционных инцидентов в production Data Lake. Понимание его причин и методов устранения критически важно.

Анатомия проблемы Out-of-Sync

Диаграмма показывает полный lifecycle проблемы: ETL пишет файлы в HDFS, не регистрируя партиции в HMS. Spark видит только то, что прописано в HMS. Команда MSCK REPAIR TABLE синхронизирует HMS с фактическим состоянием HDFS.

Команда MSCK REPAIR TABLE: детальный разбор

-- Базовый синтаксис
MSCK REPAIR TABLE bronze.events;
-- Вывод:
-- Partitions not in metastore:
--   events:date=2024-01-10
--   events:date=2024-01-11
-- Repair: Added partition to metastore events:date=2024-01-10
-- Repair: Added partition to metastore events:date=2024-01-11

-- Только проверка (без изменений в HMS):
MSCK REPAIR TABLE bronze.events CHECK PARTITIONS;

-- Добавить только недостающие партиции:
MSCK REPAIR TABLE bronze.events ADD PARTITIONS;

-- Удалить из HMS партиции которых нет на HDFS:
MSCK REPAIR TABLE bronze.events DROP PARTITIONS;

-- Полная синхронизация (добавить + удалить):
MSCK REPAIR TABLE bronze.events SYNC PARTITIONS;

PySpark: программное восстановление партиций

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("repair-partitions") \
    .config("hive.metastore.uris", "thrift://hms:9083") \
    .enableHiveSupport() \
    .getOrCreate()

# ── Способ 1: MSCK REPAIR TABLE через spark.sql() ────────────────────
spark.sql("MSCK REPAIR TABLE bronze.events")

# ── Способ 2: spark.catalog.recoverPartitions() ───────────────────────
# Аналог MSCK REPAIR TABLE, но через Spark Catalog API.
# Более интегрирован с Spark планировщиком и лучше обрабатывает ошибки.
spark.catalog.recoverPartitions("bronze.events")

# ── Способ 3: Точечное добавление партиций (быстрее MSCK) ─────────────
# Когда точно знаем какие партиции добавились
from datetime import date, timedelta

new_dates = ["2024-01-10", "2024-01-11", "2024-01-12"]

for date_str in new_dates:
    spark.sql(f"""
        ALTER TABLE bronze.events
        ADD IF NOT EXISTS PARTITION (event_date='{date_str}')
        LOCATION 'hdfs://cluster/data/bronze/events/event_date={date_str}/'
    """)
    print(f"Зарегистрирована партиция: {date_str}")

# Верификация
partitions = spark.sql("SHOW PARTITIONS bronze.events").collect()
print(f"Всего партиций в HMS: {len(partitions)}")

# ── Способ 4: Автоматизация через Airflow ────────────────────────────
# В DAG после каждой загрузки данных:
def repair_partitions_after_load(table_name: str, **kwargs) -> None:
    """
    Вызывается после успешной записи данных в HDFS.
    Регистрирует новые партиции в HMS.
    """
    from pyspark.sql import SparkSession

    spark = SparkSession.builder \
        .appName(f"repair-{table_name}") \
        .enableHiveSupport() \
        .getOrCreate()

    spark.catalog.recoverPartitions(table_name)
    count = spark.sql(f"SHOW PARTITIONS {table_name}").count()
    print(f"Таблица {table_name}: {count} партиций зарегистрировано")
    spark.stop()

Ограничения производительности MSCK REPAIR

MSCK REPAIR TABLE выполняет listStatus() для каждого уровня вложенности директорий на HDFS. Для таблицы с многоуровневым партиционированием и миллионами файлов это:

  • Сотни тысяч HDFS API-вызовов к NameNode
  • Блокирование Thrift-сервера HMS (однопоточная обработка в старых версиях)
  • Таймауты при очень большом числе партиций (> 100k)

Решение для масштаба: Никогда не используйте MSCK REPAIR на таблицах с > 10k партициями. Вместо этого:

def incremental_partition_register(
    spark: SparkSession,
    table_name: str,
    hdfs_base_path: str,
    date_from: str,
    date_to: str,
) -> int:
    """
    Регистрирует только новые партиции за заданный диапазон дат.
    Значительно быстрее MSCK REPAIR для больших таблиц.

    Используем ADD IF NOT EXISTS - идемпотентная операция:
    повторный запуск не вызовет ошибок.
    """
    from datetime import datetime, timedelta

    start = datetime.strptime(date_from, "%Y-%m-%d")
    end = datetime.strptime(date_to, "%Y-%m-%d")

    registered = 0
    current = start
    while current <= end:
        date_str = current.strftime("%Y-%m-%d")
        partition_path = f"{hdfs_base_path}/event_date={date_str}/"

        spark.sql(f"""
            ALTER TABLE {table_name}
            ADD IF NOT EXISTS PARTITION (event_date='{date_str}')
            LOCATION '{partition_path}'
        """)
        registered += 1
        current += timedelta(days=1)

    return registered

6. Schema Evolution в контексте HMS: ALTER TABLE

Schema Evolution - изменение схемы таблицы с сохранением совместимости с уже записанными данными. В контексте HMS это требует особой осторожности.

ALTER TABLE: добавление, изменение и удаление колонок

-- ── Добавление новых колонок ─────────────────────────────────────────
-- Добавление в конец списка колонок (безопасно для Parquet):
ALTER TABLE silver.user_events
ADD COLUMNS (
    platform    STRING COMMENT 'Платформа: web, mobile, api',
    session_id  STRING COMMENT 'Идентификатор сессии'
);

-- ВАЖНО: Добавление колонки меняет StorageDescriptor таблицы В HMS,
-- но НЕ изменяет физические Parquet-файлы!
-- В старых Parquet файлах (до изменения схемы) этих колонок нет.
-- При чтении Spark вернёт NULL для отсутствующих колонок.

-- ── Изменение типа колонки ────────────────────────────────────────────
-- Допустимо только для расширения типа (INT → BIGINT, STRING → stays STRING)
ALTER TABLE silver.user_events
CHANGE COLUMN user_id user_id BIGINT COMMENT 'ID пользователя (расширен до 64bit)';

-- НЕЛЬЗЯ: BIGINT → INT (сужение), INT → STRING (несовместимые типы)

-- ── Переименование колонки ───────────────────────────────────────────
-- Доступно через CHANGE COLUMN (указываем старое и новое имя):
ALTER TABLE silver.user_events
CHANGE COLUMN event_type event_category STRING
COMMENT 'Категория события (переименовано из event_type)';
-- WARNING: старые Parquet файлы всё ещё содержат колонку event_type!
-- Spark найдёт соответствие по индексу колонки, не по имени.
-- Это может привести к тихой ошибке если порядок колонок нарушен.

-- ── Удаление колонки ─────────────────────────────────────────────────
-- В Hive/Spark: удаление через REPLACE COLUMNS (заменяем весь список)
ALTER TABLE silver.user_events
REPLACE COLUMNS (
    event_id    STRING,
    user_id     BIGINT,
    event_type  STRING,
    amount      DECIMAL(18,2)
    -- platform и session_id убраны
);
-- КРИТИЧНО: физические Parquet файлы всё ещё содержат эти колонки!
-- Удаление из HMS только скрывает их от Spark, но не удаляет данные.

-- ── Изменение пути таблицы ────────────────────────────────────────────
-- НИКОГДА не делайте UPDATE в PostgreSQL HMS напрямую!
-- Только через ALTER TABLE SET LOCATION:
ALTER TABLE silver.user_events
SET LOCATION 'hdfs://cluster/data/silver_v2/user_events';

-- Для конкретной партиции:
ALTER TABLE silver.user_events
PARTITION (event_date='2024-01-15')
SET LOCATION 'hdfs://cluster/data/silver_v2/user_events/event_date=2024-01-15';

Проблема «старых партиций» при Schema Evolution

Это тонкая, но важная проблема. При добавлении колонки через ALTER TABLE ADD COLUMNS HMS обновляет StorageDescriptor таблицы. Но каждая партиция имеет собственный StorageDescriptor!

spark = SparkSession.builder.enableHiveSupport().getOrCreate()

# Добавляем новую колонку в таблицу
spark.sql("""
    ALTER TABLE silver.user_events
    ADD COLUMNS (app_version STRING)
""")

# Проверяем статус партиций:
# Новые партиции (созданные ПОСЛЕ ALTER TABLE):
spark.sql("""
    SELECT app_version, COUNT(*) as cnt
    FROM silver.user_events
    WHERE event_date = '2024-02-01'  -- Новая партиция
    GROUP BY app_version
""").show()
# Результат: корректно, включая app_version

# Старые партиции (созданные ДО ALTER TABLE):
spark.sql("""
    SELECT app_version, COUNT(*) as cnt
    FROM silver.user_events
    WHERE event_date = '2024-01-01'  -- Старая партиция
    GROUP BY app_version
""").show()
# Результат: app_version = NULL для всех строк
# Это НОРМАЛЬНО: Parquet файлы не содержат колонку app_version

# Исправление: обновить StorageDescriptor для старых партиций
# (редко нужно, но если нужно - через MSCK REPAIR или индивидуально)
spark.sql("""
    ALTER TABLE silver.user_events
    PARTITION (event_date='2024-01-01')
    CHANGE COLUMN app_version app_version STRING
""")

Антипаттерн: прямые SQL-запросы в PostgreSQL HMS

# ❌ КАТЕГОРИЧЕСКИ ЗАПРЕЩЕНО:
import psycopg2

conn = psycopg2.connect("host=hms-db user=hive password=xxx dbname=hive_metastore")
cur = conn.cursor()

# НЕПРАВИЛЬНО: прямой UPDATE в PostgreSQL HMS
cur.execute("""
    UPDATE SDS
    SET LOCATION = 'hdfs://new-cluster/data/silver/events'
    WHERE SD_ID = 42
""")
conn.commit()

# Проблемы этого подхода:
# 1. Обходит валидацию HMS (HMS может хранить кеш в памяти)
# 2. Не обновляет связанные объекты (партиции имеют свои SD_ID!)
# 3. Не сбрасывает кеш Spark Catalog Driver'а
# 4. Не проходит аудит (нет event'а в HMS transaction log)
# 5. Может нарушить foreign key constraints

# ✅ ПРАВИЛЬНО:
spark.sql("""
    ALTER TABLE silver.events
    SET LOCATION 'hdfs://new-cluster/data/silver/events'
""")
# И для партиций - либо через MSCK REPAIR, либо поштучно через
# ALTER TABLE PARTITION SET LOCATION

7. Практика: лабораторный аудит HMS и восстановление таблиц

Задание 1: создание таблицы с RegexSerDe и проверка StorageDescriptor

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .master("local[2]") \
    .appName("hive-ddl-lab") \
    .config("hive.metastore.uris", "thrift://localhost:9083") \
    .config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
    .enableHiveSupport() \
    .getOrCreate()

# Шаг 1: Создаём тестовые данные - сырой лог в HDFS
sample_logs = """192.168.1.100 - alice [15/Jan/2024:10:00:01 +0300] "GET /api/users HTTP/1.1" 200 1234 "-" "curl/7.68.0"
10.0.0.5 - - [15/Jan/2024:10:00:02 +0300] "POST /api/orders HTTP/1.1" 201 456 "https://app.example.com" "Mozilla/5.0"
192.168.1.101 - bob [15/Jan/2024:10:00:03 +0300] "DELETE /api/sessions/xyz HTTP/1.1" 401 89 "-" "Python-requests/2.28"
"""

# Записываем в локальную директорию для теста
import os, tempfile
log_dir = "/tmp/nginx_logs/"
os.makedirs(log_dir, exist_ok=True)
with open(f"{log_dir}/access_2024-01-15.log", "w") as f:
    f.write(sample_logs)

# Шаг 2: Создаём External таблицу с RegexSerDe
spark.sql("CREATE DATABASE IF NOT EXISTS lab")

spark.sql("""
    CREATE EXTERNAL TABLE IF NOT EXISTS lab.nginx_access_log (
        client_ip    STRING,
        ident        STRING,
        auth_user    STRING,
        log_time     STRING,
        request      STRING,
        status       INT,
        bytes_sent   BIGINT,
        referer      STRING,
        user_agent   STRING
    )
    ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.RegexSerDe'
    WITH SERDEPROPERTIES (
        'input.regex' =
            '(\\\\S+)\\\\s+(\\\\S+)\\\\s+(\\\\S+)\\\\s+\\\\[(.+?)\\\\]\\\\s+"(.*?)"\\\\s+(\\\\d+)\\\\s+(\\\\d+|-)\\\\s+"(.*?)"\\\\s+"(.*?)"',
        'output.format.string' = '%1$s %2$s %3$s [%4$s] "%5$s" %6$s %7$s "%8$s" "%9$s"'
    )
    STORED AS TEXTFILE
    LOCATION 'file:///tmp/nginx_logs'
""")

# Шаг 3: Проверяем что данные парсятся корректно
df = spark.sql("SELECT client_ip, status, request FROM lab.nginx_access_log")
df.show(truncate=False)
# Ожидаемый вывод:
# +---------------+------+-------------------------+
# |client_ip      |status|request                  |
# +---------------+------+-------------------------+
# |192.168.1.100  |200   |GET /api/users HTTP/1.1  |
# |10.0.0.5       |201   |POST /api/orders HTTP/1.1|
# |192.168.1.101  |401   |DELETE /api/sessions/xyz |
# +---------------+------+-------------------------+

# Шаг 4: Смотрим StorageDescriptor через DESCRIBE EXTENDED
spark.sql("DESCRIBE EXTENDED lab.nginx_access_log").show(100, truncate=False)
# В выводе ищем:
# | Storage Desc Params: |
# |   input.regex        | (\\S+)\\s+(\\S+)...  |
# | SerDe Library:       | org.apache...RegexSerDe |
# | InputFormat:         | org.apache...TextInputFormat |
# | OutputFormat:        | org.apache...HiveIgnoreKeyTextOutputFormat |

Задание 2: смотрим HMS изнутри через PostgreSQL

import psycopg2
import json

def inspect_table_in_hms(table_name: str, db_name: str = "lab") -> None:
    """
    Читает StorageDescriptor прямо из PostgreSQL HMS.
    Только для аудита - никогда не изменяем напрямую!
    """
    conn = psycopg2.connect(
        host="localhost",
        port=5432,
        user="hive",
        password="hive_password",
        dbname="hive_metastore"
    )
    cur = conn.cursor()

    # Получаем ID таблицы
    cur.execute("""
        SELECT t.TBL_ID, t.TBL_NAME, t.TBL_TYPE, d.NAME as db_name
        FROM TBLS t
        JOIN DBS d ON t.DB_ID = d.DB_ID
        WHERE t.TBL_NAME = %s AND d.NAME = %s
    """, (table_name, db_name))
    table = cur.fetchone()
    if not table:
        print(f"Таблица {db_name}.{table_name} не найдена в HMS")
        return
    tbl_id = table[0]
    print(f"Таблица: {table[1]} (ID={tbl_id}), Type={table[2]}")

    # Получаем StorageDescriptor
    cur.execute("""
        SELECT s.SD_ID, s.LOCATION, s.INPUT_FORMAT, s.OUTPUT_FORMAT,
               s.IS_COMPRESSED, ser.NAME as serde_name, ser.SLIB as serde_lib
        FROM TBLS t
        JOIN SDS s ON t.SD_ID = s.SD_ID
        JOIN SERDES ser ON s.SERDE_ID = ser.SERDE_ID
        WHERE t.TBL_ID = %s
    """, (tbl_id,))
    sd = cur.fetchone()
    print(f"\nStorageDescriptor (SD_ID={sd[0]}):")
    print(f"  LOCATION:      {sd[1]}")
    print(f"  INPUT_FORMAT:  {sd[2]}")
    print(f"  OUTPUT_FORMAT: {sd[3]}")
    print(f"  IS_COMPRESSED: {sd[4]}")
    print(f"  SERDE:         {sd[6]}")

    # Параметры SerDe
    cur.execute("""
        SELECT sp.PARAM_KEY, sp.PARAM_VALUE
        FROM TBLS t
        JOIN SDS s ON t.SD_ID = s.SD_ID
        JOIN SERDE_PARAMS sp ON s.SERDE_ID = sp.SERDE_ID
        WHERE t.TBL_ID = %s
        ORDER BY sp.PARAM_KEY
    """, (tbl_id,))
    params = cur.fetchall()
    print(f"\nSerDe параметры:")
    for key, val in params:
        print(f"  {key} = {val[:80] if val else 'NULL'}")

    # Колонки
    cur.execute("""
        SELECT cv.COLUMN_NAME, cv.TYPE_NAME, cv.INTEGER_IDX, cv.COMMENT
        FROM TBLS t
        JOIN SDS s ON t.SD_ID = s.SD_ID
        JOIN CDS c ON s.CD_ID = c.CD_ID
        JOIN COLUMNS_V2 cv ON c.CD_ID = cv.CD_ID
        WHERE t.TBL_ID = %s
        ORDER BY cv.INTEGER_IDX
    """, (tbl_id,))
    columns = cur.fetchall()
    print(f"\nКолонки ({len(columns)}):")
    for col in columns:
        print(f"  [{col[2]}] {col[0]} ({col[1]}) -- {col[3] or ''}")

    conn.close()


# Запускаем аудит
inspect_table_in_hms("nginx_access_log", "lab")

Задание 3: симуляция рассинхронизации и восстановление

import os
import subprocess
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .master("local[2]") \
    .appName("repair-lab") \
    .enableHiveSupport() \
    .getOrCreate()

# Создаём тестовую таблицу с партиционированием
spark.sql("DROP TABLE IF EXISTS lab.partitioned_events")
spark.sql("""
    CREATE EXTERNAL TABLE lab.partitioned_events (
        event_id   STRING,
        user_id    BIGINT,
        amount     DOUBLE
    )
    PARTITIONED BY (event_date STRING)
    STORED AS PARQUET
    LOCATION 'file:///tmp/partitioned_events'
""")

# Записываем данные за 3 дня через Spark (они попадут в HMS)
for date in ["2024-01-10", "2024-01-11", "2024-01-12"]:
    df = spark.range(100).select(
        F.concat(F.lit("evt-"), F.col("id").cast("string")).alias("event_id"),
        (F.rand() * 1000).cast("long").alias("user_id"),
        (F.rand() * 500).alias("amount"),
    )
    df.write \
        .mode("overwrite") \
        .parquet(f"/tmp/partitioned_events/event_date={date}/")

    spark.sql(f"""
        ALTER TABLE lab.partitioned_events
        ADD PARTITION (event_date='{date}')
        LOCATION 'file:///tmp/partitioned_events/event_date={date}/'
    """)

print("Зарегистрированные партиции:")
spark.sql("SHOW PARTITIONS lab.partitioned_events").show()

# ── СИМУЛИРУЕМ ДОБАВЛЕНИЕ ДАННЫХ В ОБХОД HMS ─────────────────────────
print("\n=== Симулируем рассинхронизацию ===")
for date in ["2024-01-13", "2024-01-14"]:
    os.makedirs(f"/tmp/partitioned_events/event_date={date}", exist_ok=True)

    # Создаём данные напрямую через Spark write (минуя HMS регистрацию)
    df = spark.range(50).select(
        F.concat(F.lit("evt-"), F.col("id").cast("string")).alias("event_id"),
        (F.rand() * 1000).cast("long").alias("user_id"),
        (F.rand() * 500).alias("amount"),
    )
    df.coalesce(1).write \
        .mode("overwrite") \
        .parquet(f"/tmp/partitioned_events/event_date={date}/")

print("Файлы добавлены напрямую в HDFS (без HMS регистрации)")

# Проверяем - Spark НЕ видит новые партиции
print("\nЗапрос ПОСЛЕ рассинхронизации:")
count_before = spark.sql("""
    SELECT event_date, COUNT(*) as cnt
    FROM lab.partitioned_events
    GROUP BY event_date ORDER BY event_date
""").collect()
for row in count_before:
    print(f"  {row.event_date}: {row.cnt}")
print("Партиции 2024-01-13 и 2024-01-14 НЕ ВИДНЫ!")

# ── СПОСОБ 1: MSCK REPAIR TABLE ───────────────────────────────────────
print("\n=== Восстановление через MSCK REPAIR TABLE ===")
spark.sql("MSCK REPAIR TABLE lab.partitioned_events")

count_after_msck = spark.sql("""
    SELECT event_date, COUNT(*) as cnt
    FROM lab.partitioned_events
    GROUP BY event_date ORDER BY event_date
""").collect()
print("После MSCK REPAIR:")
for row in count_after_msck:
    print(f"  {row.event_date}: {row.cnt}")

# ── СПОСОБ 2: catalog.recoverPartitions() ───────────────────────────
print("\n=== Восстановление через spark.catalog.recoverPartitions() ===")
# Сначала сбрасываем партиции обратно для демонстрации
spark.sql("ALTER TABLE lab.partitioned_events DROP PARTITION (event_date='2024-01-13')")
spark.sql("ALTER TABLE lab.partitioned_events DROP PARTITION (event_date='2024-01-14')")

print("После удаления (имитируем рассинхронизацию снова):")
spark.sql("SHOW PARTITIONS lab.partitioned_events").show()

# Восстанавливаем через Catalog API
spark.catalog.recoverPartitions("lab.partitioned_events")

print("После recoverPartitions():")
spark.sql("SHOW PARTITIONS lab.partitioned_events").show()

# Финальная проверка
total = spark.sql("SELECT COUNT(*) FROM lab.partitioned_events").first()[0]
print(f"\n✅ Всего строк в таблице: {total:,}")
print(f"✅ Восстановление успешно!")

spark.stop()

Чек-лист проектирования DDL для on-premise Data Lake

-- ═══════════════════════════════════════════════════════════════════
-- ЧИСТЫЙ ШАБЛОН DDL ДЛЯ PRODUCTION DATA LAKE
-- Соответствует best practices Hive + Spark
-- ═══════════════════════════════════════════════════════════════════

-- 1. Всегда External Table для данных Bronze/Silver/Gold
-- 2. Явный LOCATION - никаких managed-путей
-- 3. STORED AS PARQUET или ORC - никакого TEXTFILE для аналитики
-- 4. Партиционирование по низкокардинальной колонке
-- 5. TBLPROPERTIES для документирования и настройки

CREATE EXTERNAL TABLE IF NOT EXISTS silver.transactions (

    -- Строгая типизация: используйте правильные типы
    txn_id        BIGINT        NOT NULL  COMMENT 'Уникальный ID транзакции',
    user_id       BIGINT                  COMMENT 'ID пользователя',
    merchant_id   INT                     COMMENT 'ID мерчанта',
    amount        DECIMAL(18,2)           COMMENT 'Сумма в базовой валюте',
    currency      CHAR(3)                 COMMENT 'ISO-4217 код валюты',

    -- Сложные типы: Struct, Array, Map
    items         ARRAY<STRUCT<
                      sku:    STRING,
                      qty:    INT,
                      price:  DECIMAL(10,2)
                  >>          COMMENT 'Позиции в транзакции',

    metadata      MAP<STRING, STRING>     COMMENT 'Доп. атрибуты источника',

    -- Аудит: технические поля
    source_system STRING                  COMMENT 'Система-источник',
    created_at    TIMESTAMP               COMMENT 'Время создания в источнике',
    loaded_at     TIMESTAMP               COMMENT 'Время загрузки в Data Lake'
)
COMMENT 'Транзакции Silver слоя: очищенные, дедуплицированные'

-- Партиционирование: только низкокардинальные колонки!
-- txn_date DATE = 365 партиций в год (разумно)
-- НЕ партиционировать по user_id, txn_id (высокая кардинальность!)
PARTITIONED BY (txn_date DATE COMMENT 'Дата транзакции')

-- Бакетирование для join-оптимизации (опционально):
-- CLUSTERED BY (user_id) INTO 50 BUCKETS

STORED AS PARQUET

-- Явный путь (никаких managed!)
LOCATION 'hdfs://cluster/data/silver/transactions'

TBLPROPERTIES (
    -- Документирование
    'owner'            = 'data-platform-team',
    'created_date'     = '2024-01-15',
    'source_systems'   = 'postgres-payments, kafka-events',
    'refresh_cadence'  = 'hourly',
    'sla_freshness_h'  = '2',

    -- Настройки Spark
    'spark.sql.statistics.numRows'    = '500000000',
    'spark.sql.statistics.totalSize'  = '107374182400',

    -- Parquet-специфичные настройки
    'parquet.compression'             = 'snappy',
    'parquet.block.size'              = '134217728',  -- 128 MB = HDFS block

    -- Защита от случайного DROP с удалением данных
    -- (только для HMS 3.0+, игнорируется в старых версиях)
    'external.table.purge'            = 'false'
);

Итоги: ключевые выводы

Hive DDL - это декларативный язык описания того, как Spark должен интерпретировать физические файлы на HDFS. Понимание физики HMS (StorageDescriptor, SerDe, Input/OutputFormat) позволяет:

  • Диагностировать проблемы «таблица есть, данных нет» - это MSCK REPAIR
  • Понять откуда берётся ClassNotFoundException - это отсутствующий JAR SerDe в classpath
  • Правильно спроектировать Schema Evolution - знать, что ALTER TABLE меняет таблицу в HMS, но не физические Parquet-файлы
  • Безопасно мигрировать данные - только через официальный HMS API, никогда через прямой UPDATE в PostgreSQL
  • Оптимально выбрать STORED AS - Parquet для аналитики, Avro для потоков, TEXTFILE только для staging

Три правила production Data Lake DDL:

  1. External Tables everywhere - никогда не создавайте Managed Tables для продуктовых данных
  2. STORED AS PARQUET (или ORC) - никогда не используйте TEXTFILE для аналитических таблиц
  3. Автоматизируйте регистрацию партиций - используйте ALTER TABLE ADD PARTITION сразу после записи данных, не полагайтесь на MSCK REPAIR как основной механизм