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 и лабораторный аудит метастора изнутри.
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, но не зарегистрирована в HMSSerDeException: 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:
- External Tables everywhere - никогда не создавайте Managed Tables для продуктовых данных
- STORED AS PARQUET (или ORC) - никогда не используйте TEXTFILE для аналитических таблиц
- Автоматизируйте регистрацию партиций - используйте
ALTER TABLE ADD PARTITIONсразу после записи данных, не полагайтесь на MSCK REPAIR как основной механизм