Managed vs External Tables: lifecycle данных, LOCATION, DROP TABLE поведение и Data Lake паттерн
Фундаментальный разбор двух подходов к владению данными: физика операций DROP TABLE, параметр LOCATION и warehouse.dir, паттерн Medallion Architecture с External/Managed, файлы-сироты, MSCK REPAIR TABLE, полная реализация на PySpark и SQL, траблшутинг рассинхронизации каталога.
1. Архитектурная философия: два подхода к владению данными¶
Разделение таблиц на Managed и External - это не просто техническая деталь SQL DDL. Это фундаментальное архитектурное решение о том, кто несёт ответственность за данные: движок обработки (Spark, Hive) или Data Platform (HDFS, S3, Iceberg). Ошибочный выбор в одну сторону приводит к потере петабайт данных, в другую - к накоплению терабайт «мёртвых» файлов, которые никто не контролирует.
Managed (Internal) Tables: концепция полного владения¶
Когда Spark создаёт Managed таблицу, он берёт на себя полный контроль над жизненным циклом данных. Это означает:
- Spark сам решает куда физически положить файлы (внутрь
warehouse.dir) - Spark создаёт директорию таблицы на HDFS
- Spark отвечает за удаление файлов при
DROP TABLE
Это похоже на то, как операционная система «владеет» временными файлами в /tmp - она их создаёт и она же убирает мусор.
External Tables: концепция слабой связи¶
Когда Spark создаёт External таблицу, он регистрирует в HMS только описание данных: схему, формат, путь. Сами данные существуют независимо, в произвольном месте на HDFS или S3.
Это похоже на каталожную карточку в библиотеке: карточка (метаданные) описывает книгу (данные), но уничтожение карточки не означает уничтожение самой книги. Книга лежит на полке вне зависимости от состояния каталога.
Смена парадигмы в Spark 3.x: свойство PURGE¶
В Spark 3.x поведение DROP TABLE по умолчанию изменилось по сравнению с классическим Hive. Ключевые отличия:
В классическом Hive: DROP TABLE для Managed таблицы перемещает данные в корзину HDFS (.Trash), откуда их можно восстановить. Физического удаления нет немедленно.
В Spark 3.x: DROP TABLE для Managed таблицы по умолчанию также перемещает в корзину.
Ключевое слово PURGE: DROP TABLE tableName PURGE - удаляет файлы немедленно и необратимо, минуя корзину. Это особенно важно понимать при написании DDL-скриптов автоматизации.
-- Стандартное удаление (с перемещением в корзину .Trash)
DROP TABLE sales.orders;
-- Необратимое удаление без корзины (как в классическом Hive PURGE)
DROP TABLE sales.orders PURGE;
-- Проверить что в .Trash осталось:
-- hdfs dfs -ls /user/spark/.Trash/Current/user/hive/warehouse/sales.db/
2. Жизненный цикл метаданных и физических файлов¶
Чтобы понять разницу между Managed и External, нужно разобрать что именно происходит при каждой операции: создании, наполнении данными и удалении.
Создание таблицы: где рождаются файлы¶
Схема показывает принципиальное различие в порядке операций: для Managed таблицы сначала HMS создаёт запись и Spark создаёт директорию, затем в неё пишутся данные. Для External таблицы данные могут уже существовать независимо, а HMS просто регистрирует на них указатель.
Операция DROP TABLE: физика процесса¶
DROP TABLE на Managed таблице:
1. Spark получает запрос DROP TABLE db.orders
2. Spark запрашивает HMS: какой путь у таблицы db.orders?
→ HMS возвращает: /user/hive/warehouse/db.db/orders/
3. HMS удаляет запись о таблице из PostgreSQL
(из таблиц TBLS, SDS, PARTITIONS, COLUMNS_V2...)
4. Spark выполняет hdfs dfs -rm -r /user/hive/warehouse/db.db/orders/
→ Файлы перемещаются в /user/spark/.Trash/Current/...
→ Или удаляются сразу если использован PURGE
5. Данные УНИЧТОЖЕНЫ (или в корзине)
DROP TABLE на External таблице:
1. Spark получает запрос DROP TABLE db.orders
2. Spark проверяет HMS: TBL_TYPE = EXTERNAL_TABLE
3. HMS удаляет ТОЛЬКО запись о таблице из PostgreSQL
4. hdfs dfs -rm НЕ вызывается
5. Файлы в /data/custom/orders/ ОСТАЮТСЯ нетронутыми
6. Данные в безопасности!
Корпоративные риски: реальные кейсы потери данных¶
Путаница между Managed и External приводит к реальным авариям. Типичные сценарии:
Сценарий 1: ошибка в миграционном скрипте. Инженер пишет DDL-скрипт миграции схемы: DROP TABLE IF EXISTS gold.revenue_report; CREATE TABLE gold.revenue_report AS SELECT.... Не зная что gold.revenue_report - Managed таблица, он безвозвратно удаляет данные за весь год. Восстановление только из резервной копии.
Сценарий 2: CI/CD cleanup. В CI/CD пайплайне есть шаг «очистить тестовые таблицы»: DROP TABLE IF EXISTS db.*. Разработчик случайно указывает production базу данных вместо тестовой. Managed таблицы production уничтожены.
Сценарий 3: переименование через DROP+CREATE. Попытка «переименовать» Managed таблицу через DROP + CREATE новой = необратимая потеря данных. Правильный путь: ALTER TABLE old_name RENAME TO new_name.
3. Магия параметра LOCATION и дефолтный Warehouse¶
Где живёт Managed слой: warehouse.dir¶
Каждая Managed таблица создаётся внутри warehouse директории. По умолчанию это /user/hive/warehouse/ в HDFS. Структура следует строгой иерархии:
warehouse.dir/
├── database1.db/
│ ├── table1/
│ │ ├── part-0.parquet
│ │ └── part-1.parquet
│ └── partitioned_table/
│ ├── date=2024-01-15/
│ │ └── part-0.parquet
│ └── date=2024-01-16/
│ └── part-0.parquet
└── database2.db/
└── another_table/
Настройка warehouse.dir в SparkSession и hive-site.xml:
spark = SparkSession.builder \
# Это warehouse.dir для Spark встроенного каталога
.config("spark.sql.warehouse.dir", "hdfs://cluster/user/hive/warehouse") \
# Это для HMS (должны совпадать в production)
.config("hive.metastore.warehouse.dir", "hdfs://cluster/user/hive/warehouse") \
.enableHiveSupport() \
.getOrCreate()
# Проверить warehouse.dir на running кластере
hdfs dfs -ls /user/hive/warehouse/
# drwxrwxr-x - hive hadoop 0 2024-01-15 10:00 /user/hive/warehouse/default.db
# drwxrwxr-x - hive hadoop 0 2024-01-10 12:00 /user/hive/warehouse/sales.db
# drwxrwxr-x - hive hadoop 0 2024-01-12 08:00 /user/hive/warehouse/analytics.db
Параметр LOCATION: явное управление путями¶
LOCATION - ключевое слово DDL, которое превращает Managed таблицу в External и задаёт точный физический путь данных:
# Через DataFrame API: LOCATION через .option("path", ...)
df.write \
.mode("overwrite") \
.option("path", "hdfs://cluster/data/production/sales/orders") \
.saveAsTable("sales.orders")
# Это создаёт External таблицу, потому что путь задан явно!
# Через Spark SQL DDL:
spark.sql("""
CREATE TABLE IF NOT EXISTS sales.orders (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
order_date DATE
)
USING PARQUET
LOCATION 'hdfs://cluster/data/production/sales/orders'
""")
Свобода путей: несколько таблиц на один датасет¶
Одна из мощных возможностей External таблиц - несколько логических таблиц могут указывать на один и тот же физический путь HDFS. Это позволяет предоставить разные "представления" одних данных разным командам без дублирования файлов:
# Физические данные лежат в одном месте
hdfs_path = "hdfs://cluster/data/events/raw/"
# Таблица для команды аналитики (все колонки)
spark.sql(f"""
CREATE EXTERNAL TABLE IF NOT EXISTS analytics.raw_events (
event_id STRING, user_id BIGINT, event_type STRING,
payload STRING, ts TIMESTAMP
)
STORED AS PARQUET
LOCATION '{hdfs_path}'
""")
# Таблица для команды DS (только нужные колонки)
spark.sql(f"""
CREATE EXTERNAL TABLE IF NOT EXISTS ml_team.events_features (
user_id BIGINT, event_type STRING, ts TIMESTAMP
-- payload намеренно исключён (конфиденциальные данные)
)
STORED AS PARQUET
LOCATION '{hdfs_path}'
-- Тот же LOCATION! Те же физические файлы!
""")
# При чтении ml_team.events_features Spark читает те же файлы
# но возвращает только колонки из схемы этой таблицы
Это паттерн «Column Security через схему»: разные команды видят разные колонки одного датасета через разные внешние таблицы, но физически данные не дублируются.
4. Реализация паттернов на PySpark и Spark SQL¶
Создание Managed таблиц через DataFrame API¶
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, DateType
spark = SparkSession.builder \
.appName("managed-vs-external-demo") \
.config("hive.metastore.uris", "thrift://hms:9083") \
.config("spark.sql.warehouse.dir", "hdfs://cluster/user/hive/warehouse") \
.enableHiveSupport() \
.getOrCreate()
# Тестовые данные
schema = StructType([
StructField("order_id", StringType(), nullable=False),
StructField("user_id", LongType(), nullable=False),
StructField("amount", DoubleType(), nullable=True),
StructField("product_id", StringType(), nullable=True),
StructField("order_date", DateType(), nullable=False),
])
df = spark.range(10_000).select(
F.concat(F.lit("ORD-"), F.col("id").cast("string")).alias("order_id"),
(F.rand() * 100_000).cast("long").alias("user_id"),
(F.rand() * 500).alias("amount"),
F.concat(F.lit("PROD-"), (F.rand() * 1000).cast("int").cast("string")).alias("product_id"),
F.date_add(F.lit("2024-01-01"), (F.rand() * 90).cast("int")).alias("order_date"),
)
# ── Способ 1: saveAsTable без явного пути = MANAGED ──────────────────
spark.sql("CREATE DATABASE IF NOT EXISTS demo")
spark.sql("USE demo")
df.write \
.mode("overwrite") \
.partitionBy("order_date") \
.format("parquet") \
.saveAsTable("demo.orders_managed")
# ── Способ 2: Spark SQL CREATE TABLE без LOCATION = MANAGED ──────────
spark.sql("""
CREATE TABLE IF NOT EXISTS demo.orders_managed_sql
USING PARQUET
PARTITIONED BY (order_date)
AS
SELECT * FROM demo.orders_managed
""")
Создание External таблиц через DataFrame API¶
# ── Способ 1: .option("path", ...) = EXTERNAL ────────────────────────
EXTERNAL_PATH = "hdfs://cluster/data/demo/orders_external"
df.write \
.mode("overwrite") \
.partitionBy("order_date") \
.format("parquet") \
.option("path", EXTERNAL_PATH) \
.saveAsTable("demo.orders_external")
# ── Способ 2: Явный SQL с LOCATION = EXTERNAL ────────────────────────
# Сначала пишем данные на HDFS
df.write \
.mode("overwrite") \
.partitionBy("order_date") \
.parquet("hdfs://cluster/data/demo/orders_direct")
# Потом регистрируем таблицу поверх этих данных
spark.sql("""
CREATE TABLE IF NOT EXISTS demo.orders_direct_sql (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
product_id STRING
)
USING PARQUET
PARTITIONED BY (order_date DATE)
LOCATION 'hdfs://cluster/data/demo/orders_direct'
""")
# Регистрируем партиции (данные записаны до HMS регистрации!)
spark.sql("MSCK REPAIR TABLE demo.orders_direct_sql")
SQL DDL: полный синтаксис¶
-- ── MANAGED TABLE (различные варианты синтаксиса) ───────────────────
-- Hive-совместимый синтаксис
CREATE TABLE IF NOT EXISTS sales.managed_orders (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
order_date DATE
)
PARTITIONED BY (region STRING)
STORED AS PARQUET
TBLPROPERTIES (
'created_by' = 'data-platform-team',
'created_at' = '2024-01-15'
);
-- Spark SQL синтаксис (USING вместо STORED AS)
CREATE TABLE IF NOT EXISTS sales.managed_orders_spark (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
order_date DATE
)
USING PARQUET
PARTITIONED BY (region);
-- CTAS (Create Table As Select): быстрое создание с данными
CREATE TABLE sales.managed_orders_ctas
USING PARQUET
PARTITIONED BY (region)
AS
SELECT * FROM raw.incoming_orders;
-- ── EXTERNAL TABLE ────────────────────────────────────────────────────
-- Классический Hive-синтаксис: ключевое слово EXTERNAL
CREATE EXTERNAL TABLE IF NOT EXISTS sales.external_orders (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
order_date DATE
)
PARTITIONED BY (region STRING)
ROW FORMAT SERDE 'org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe'
STORED AS
INPUTFORMAT 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat'
OUTPUTFORMAT 'org.apache.hadoop.mapreduce.lib.output.FileOutputFormat'
LOCATION 'hdfs://cluster/data/sales/orders';
-- Spark SQL синтаксис: LOCATION без EXTERNAL тоже создаёт External!
CREATE TABLE IF NOT EXISTS sales.external_orders_spark (
order_id STRING,
user_id BIGINT,
amount DOUBLE,
order_date DATE
)
USING PARQUET
PARTITIONED BY (region)
LOCATION 'hdfs://cluster/data/sales/orders';
-- В Spark: наличие LOCATION = External table
Инспекция типа таблицы: как проверить что создали¶
# ── DESCRIBE EXTENDED: полная информация о таблице ────────────────────
spark.sql("DESCRIBE EXTENDED demo.orders_managed").show(100, truncate=False)
# Ищем строки с Table Type:
# | Table Type | MANAGED | ... |
# | Location | hdfs://cluster/user/hive/warehouse/demo.db/orders_managed |
spark.sql("DESCRIBE EXTENDED demo.orders_external").show(100, truncate=False)
# | Table Type | EXTERNAL | ... |
# | Location | hdfs://cluster/data/demo/orders_external |
# ── Программная проверка через spark.catalog ─────────────────────────
tables = spark.catalog.listTables("demo")
for t in tables.collect():
print(f"{t.name:30s} | tableType: {t.tableType:8s} | isTemporary: {t.isTemporary}")
# Вывод:
# orders_managed | tableType: MANAGED | isTemporary: False
# orders_external | tableType: EXTERNAL | isTemporary: False
# ── SHOW CREATE TABLE: показывает полный DDL как он есть в HMS ────────
spark.sql("SHOW CREATE TABLE demo.orders_managed").show(1, truncate=False)
# CREATE TABLE demo.orders_managed (
# order_id STRING,
# user_id BIGINT,
# ...
# )
# USING parquet
# PARTITIONED BY (order_date)
# LOCATION 'hdfs://cluster/user/hive/warehouse/demo.db/orders_managed'
# TBLPROPERTIES (...)
# Обратите внимание: HMS автоматически добавил LOCATION для Managed!
# Это путь внутри warehouse.dir
spark.sql("SHOW CREATE TABLE demo.orders_external").show(1, truncate=False)
# ...
# LOCATION 'hdfs://cluster/data/demo/orders_external'
# Явно указанный кастомный путь
5. Data Lake паттерн: слойность и Data Governance¶
Medallion Architecture: где какой тип таблицы¶
Правило выбора типа таблицы в Medallion Architecture не произвольное - оно следует из логики владения данными и рисков:
Bronze слой - всегда External. Данные Bronze - это «сырые факты» о произошедших событиях. Их нельзя пересоздать из других источников (они могут быть уже недоступны). Потеря Bronze данных означает потерю истории. External таблица гарантирует: если HMS упадёт, если кто-то случайно дропнет таблицу, данные останутся в HDFS нетронутыми.
Silver слой - External по тем же причинам. Silver - это очищенные и нормализованные данные. Они воспроизводимы из Bronze, но пересчёт занимает часы или дни. External таблица - дополнительный уровень защиты.
Gold слой - Managed допустим, но осторожно. Агрегаты в Gold пересчитываемы из Silver за разумное время. Если Gold таблица дропнута - её можно пересоздать запуском ETL Job'а. Managed таблицы здесь упрощают lifecycle management: устаревшую золотую витрину можно удалить одной командой без ручной очистки HDFS.
Код паттерна: создание слоёв с правильными типами¶
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder \
.appName("medallion-architecture") \
.config("hive.metastore.uris", "thrift://hms:9083") \
.config("spark.sql.warehouse.dir", "hdfs://cluster/user/hive/warehouse") \
.enableHiveSupport() \
.getOrCreate()
# ── Bronze Layer: External таблицы ──────────────────────────────────
spark.sql("CREATE DATABASE IF NOT EXISTS bronze")
# Bronze: данные пишутся из Kafka consumer / CDC
def register_bronze_table(table_name: str, schema_ddl: str, hdfs_path: str) -> None:
"""
Регистрирует Bronze таблицу как External.
Физические данные уже существуют на HDFS (из Kafka consumer / CDC).
HMS только регистрирует указатель на них.
"""
spark.sql(f"""
CREATE EXTERNAL TABLE IF NOT EXISTS bronze.{table_name}
({schema_ddl})
USING PARQUET
PARTITIONED BY (ingestion_date STRING)
LOCATION 'hdfs://cluster/data/bronze/{table_name}/'
""")
# Синхронизируем партиции (данные могли быть загружены без Spark)
spark.sql(f"MSCK REPAIR TABLE bronze.{table_name}")
register_bronze_table(
"events",
"event_id STRING, user_id BIGINT, event_type STRING, payload STRING",
"hdfs://cluster/data/bronze/events/"
)
# ── Silver Layer: External таблицы (из соображений безопасности) ────
spark.sql("CREATE DATABASE IF NOT EXISTS silver")
def process_bronze_to_silver(
source_table: str,
target_table: str,
silver_path: str,
partition_col: str = "event_date",
) -> int:
"""
Трансформирует Bronze → Silver.
Записывает как External таблицу для защиты данных.
"""
df = spark.table(f"bronze.{source_table}")
silver_df = df.dropDuplicates(["event_id"]) \
.filter(F.col("event_type").isNotNull()) \
.withColumn("event_date",
F.to_date(F.col("ingestion_date")))
# Пишем данные напрямую (без saveAsTable)
silver_df.write \
.mode("overwrite") \
.partitionBy(partition_col) \
.parquet(silver_path)
row_count = silver_df.count()
# Регистрируем / обновляем External таблицу
spark.sql(f"""
CREATE TABLE IF NOT EXISTS silver.{target_table}
USING PARQUET
PARTITIONED BY ({partition_col})
LOCATION '{silver_path}'
""")
spark.sql(f"MSCK REPAIR TABLE silver.{target_table}")
return row_count
count = process_bronze_to_silver(
"events",
"clean_events",
"hdfs://cluster/data/silver/events/"
)
print(f"Silver: {count:,} строк")
# ── Gold Layer: Managed таблицы (допустимо) ───────────────────────────
spark.sql("CREATE DATABASE IF NOT EXISTS gold")
def create_gold_aggregation(
source_table: str,
target_table: str,
agg_sql: str,
) -> None:
"""
Создаёт Gold витрину как Managed таблицу.
Managed подходит для Gold: агрегаты воспроизводимы из Silver,
DROP TABLE + пересчёт = безопасная операция обновления витрины.
"""
spark.sql(f"""
CREATE OR REPLACE TABLE gold.{target_table}
USING PARQUET
AS
{agg_sql}
""")
# CREATE OR REPLACE TABLE создаёт Managed таблицу
print(f"Gold: таблица gold.{target_table} создана")
print("Type: MANAGED (файлы в warehouse.dir)")
create_gold_aggregation(
"silver.clean_events",
"daily_event_stats",
"""
SELECT
event_date,
event_type,
COUNT(*) as event_count,
COUNT(DISTINCT user_id) as unique_users
FROM silver.clean_events
GROUP BY event_date, event_type
"""
)
Data Access Control: Ranger / Sentry поверх Managed/External¶
В enterprise-среде разграничение прав часто накладывается на концепцию Managed/External:
- External таблицы с конкретными LOCATION путями легко интегрируются с политиками Ranger: права на директорию HDFS и права на таблицу в HMS настраиваются независимо, что даёт гибкий контроль.
- Managed таблицы проще контролировать через HMS-уровень (DROP требует
ALLправ на таблицу), но труднее настроить «ручной» доступ к файлам, минуя HMS.
6. Проблема файлов-сирот (Orphan Files) и синхронизация партиций¶
Что такое файлы-сироты¶
Orphan Files - это физические файлы на HDFS, на которые нет указателя в HMS. Они существуют, занимают место на дисках, но ни одна таблица их «не видит».
Способы появления файлов-сирот:
- DROP EXTERNAL TABLE - данные остались, HMS запись удалена
- Прямая запись через hdfs dfs -put без регистрации в HMS
- Переименование LOCATION - старые данные остались на старом пути
- Незавершённые Spark-задачи - оставляют
_temporary/директории - HDFS distcp без последующей регистрации таблицы
# Скрипт поиска файлов-сирот в Data Lake
# Сравниваем физические пути с LOCATION'ами всех External таблиц
# Шаг 1: собираем список всех LOCATION из HMS
hdfs dfsadmin -report | grep "Name:" > /tmp/active_paths.txt
# Шаг 2: листинг всего Data Lake
hdfs dfs -ls -R /data/ | grep -v "^d" | awk '{print $8}' \
| sed 's|/[^/]*$||' | sort -u > /tmp/all_hdfs_dirs.txt
# Шаг 3: через SparkSQL собираем пути зарегистрированных таблиц
spark.sql("""
SELECT table_schema, table_name, table_type, location
FROM information_schema.tables
WHERE table_type = 'EXTERNAL'
""").write.csv("/tmp/registered_locations/")
# Шаг 4: Python-скрипт сравнения
# (comm -23 all_paths registered_paths | xargs hdfs dfs -du -h)
Рассинхронизация каталога: The Metadata Gap¶
Самая распространённая проблема при работе с External таблицами - несоответствие между физическим состоянием HDFS и тем, что знает HMS:
Методы лечения рассинхронизации¶
Метод 1: MSCK REPAIR TABLE - полная синхронизация
# MSCK REPAIR TABLE (MetaStore Check Repair Table)
# Сканирует весь LOCATION путь таблицы в HDFS
# Находит директории которых нет в HMS → добавляет как партиции
# Находит партиции в HMS которых нет на HDFS → удаляет из HMS
spark.sql("MSCK REPAIR TABLE silver.events")
# Вывод: Partitions not in metastore: events:event_date=2024-01-13
# Repair: Added partition to metastore events:event_date=2024-01-13
# MSCK REPAIR TABLE может быть медленным при тысячах партиций!
# Он делает NameNode listStatus() для каждой директории в LOCATION
# При 100k+ партиций → сотни секунд ожидания
# Более быстрый вариант: SYNC PARTITIONS (Spark 3.1+)
spark.sql("ALTER TABLE silver.events RECOVER PARTITIONS")
# Аналог MSCK REPAIR TABLE, оптимизированный для Spark
# Через Catalog API (Spark 3.0+)
spark.catalog.recoverPartitions("silver.events")
Метод 2: ALTER TABLE ADD PARTITION - точечное добавление
# Если вы знаете конкретные новые партиции - быстрее добавить явно
spark.sql("""
ALTER TABLE silver.events
ADD IF NOT EXISTS PARTITION (event_date='2024-01-13')
LOCATION 'hdfs://cluster/data/silver/events/event_date=2024-01-13/'
""")
# ADD IF NOT EXISTS: не вызывает ошибку если партиция уже есть
# Батчевое добавление нескольких партиций за один запрос (эффективнее!)
spark.sql("""
ALTER TABLE silver.events ADD IF NOT EXISTS
PARTITION (event_date='2024-01-13')
LOCATION 'hdfs://cluster/data/silver/events/event_date=2024-01-13/'
PARTITION (event_date='2024-01-14')
LOCATION 'hdfs://cluster/data/silver/events/event_date=2024-01-14/'
PARTITION (event_date='2024-01-15')
LOCATION 'hdfs://cluster/data/silver/events/event_date=2024-01-15/'
""")
Метод 3: refreshTable() - обновление кеша Spark Driver'а
# Если данные были добавлены Spark'ом но кеш Driver'а устарел
# refreshTable НЕ сканирует HDFS - только сбрасывает кеш
spark.catalog.refreshTable("silver.events")
# refreshByPath - обновить конкретный путь
spark.catalog.refreshByPath("hdfs://cluster/data/silver/events/event_date=2024-01-13/")
Удаление партиций которых нет на HDFS¶
Обратная проблема: файлы удалены с HDFS, но HMS всё ещё знает о партициях. Запрос к такой партиции вернёт ошибку или пустой результат.
# Найти «призрачные» партиции (есть в HMS, нет на HDFS)
def find_ghost_partitions(spark, table_name: str, hdfs_base_path: str) -> list[str]:
"""
Сравнивает партиции HMS с реальными директориями на HDFS.
Возвращает список «призрачных» партиций.
"""
import subprocess
# Партиции по мнению HMS
hms_partitions = {
row.partition
for row in spark.sql(f"SHOW PARTITIONS {table_name}").collect()
}
# Реальные директории в HDFS
result = subprocess.run(
["hdfs", "dfs", "-ls", hdfs_base_path],
capture_output=True, text=True
)
hdfs_dirs = set()
for line in result.stdout.splitlines():
if line.startswith("d"):
dir_name = line.split()[-1].split("/")[-1]
hdfs_dirs.add(dir_name)
# Призрачные: в HMS но нет на HDFS
ghosts = hms_partitions - hdfs_dirs
return list(ghosts)
ghosts = find_ghost_partitions(
spark,
"silver.events",
"hdfs://cluster/data/silver/events/"
)
if ghosts:
print(f"Найдено призрачных партиций: {len(ghosts)}")
for g in ghosts:
spark.sql(f"ALTER TABLE silver.events DROP PARTITION ({g})")
print(f" Удалена: {g}")
7. Практика: лабораторный аудит и симуляция катастрофы¶
Шаг 1: Подготовка данных и создание обоих типов таблиц¶
# lab_managed_vs_external.py
# Полная лабораторная работа
import subprocess
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder \
.master("local[2]") \
.appName("lab-managed-external") \
.config("hive.metastore.uris", "thrift://localhost:9083") \
.config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
.enableHiveSupport() \
.getOrCreate()
spark.sql("CREATE DATABASE IF NOT EXISTS lab")
# Тестовые данные
df = spark.range(5000).select(
F.concat(F.lit("TXN-"), F.col("id").cast("string")).alias("txn_id"),
(F.rand() * 10_000).cast("long").alias("user_id"),
(F.rand() * 1000).alias("amount"),
F.array(F.lit("RUB"), F.lit("USD"), F.lit("EUR")).getItem(
(F.rand() * 3).cast("int")
).alias("currency"),
F.date_add(F.lit("2024-01-01"), (F.rand() * 30).cast("int")).alias("txn_date"),
)
# ── Создаём Managed таблицу ───────────────────────────────────────────
print("=== Создание Managed таблицы ===")
df.write \
.mode("overwrite") \
.partitionBy("txn_date") \
.saveAsTable("lab.transactions_managed")
# ── Создаём External таблицу ──────────────────────────────────────────
EXTERNAL_PATH = "/tmp/lab-external-transactions"
print("\n=== Создание External таблицы ===")
df.write \
.mode("overwrite") \
.partitionBy("txn_date") \
.option("path", EXTERNAL_PATH) \
.saveAsTable("lab.transactions_external")
Шаг 2: Инспекция физических путей и HMS метаданных¶
# ── Показываем где физически лежат данные ────────────────────────────
print("\n=== Физические пути ===")
def get_table_location(spark, table_name: str) -> str:
"""Извлекает LOCATION из DESCRIBE EXTENDED."""
result = spark.sql(f"DESCRIBE EXTENDED {table_name}")
location_rows = result.filter("col_name = 'Location'").collect()
return location_rows[0]["data_type"] if location_rows else "UNKNOWN"
def get_table_type(spark, table_name: str) -> str:
"""Извлекает Type (MANAGED/EXTERNAL) из DESCRIBE EXTENDED."""
result = spark.sql(f"DESCRIBE EXTENDED {table_name}")
type_rows = result.filter("col_name = 'Type'").collect()
return type_rows[0]["data_type"] if type_rows else "UNKNOWN"
managed_location = get_table_location(spark, "lab.transactions_managed")
external_location = get_table_location(spark, "lab.transactions_external")
print(f"Managed location: {managed_location}")
print(f"External location: {external_location}")
print(f"Managed type: {get_table_type(spark, 'lab.transactions_managed')}")
print(f"External type: {get_table_type(spark, 'lab.transactions_external')}")
# Проверяем файлы на диске
for path, label in [(managed_location.replace("file://", ""), "Managed"),
(external_location, "External")]:
result = subprocess.run(["ls", "-la", path], capture_output=True, text=True)
print(f"\n{label} файлы ({path}):")
print(result.stdout[:500])
# ── Смотрим в PostgreSQL HMS ──────────────────────────────────────────
# (Для Docker стенда с HMS)
print("\n=== HMS PostgreSQL ===")
spark.sql("""
SELECT TBL_NAME, TBL_TYPE, CREATE_TIME
FROM TBLS
WHERE TBL_NAME IN ('transactions_managed', 'transactions_external')
""")
# (Этот запрос работает при наличии прямого JDBC соединения к PostgreSQL)
Шаг 3: Симуляция катастрофы - DROP TABLE для обоих типов¶
# ── СИМУЛЯЦИЯ: DROP TABLE ─────────────────────────────────────────────
print("\n=== СИМУЛЯЦИЯ DROP TABLE ===")
# Считаем строки до удаления
managed_count = spark.sql("SELECT COUNT(*) FROM lab.transactions_managed").first()[0]
external_count = spark.sql("SELECT COUNT(*) FROM lab.transactions_external").first()[0]
print(f"До DROP: Managed={managed_count:,}, External={external_count:,}")
# DROP обеих таблиц
print("Выполняем DROP TABLE для обеих таблиц...")
spark.sql("DROP TABLE lab.transactions_managed")
spark.sql("DROP TABLE lab.transactions_external")
# Проверяем что HMS таблицы ушли
remaining = spark.catalog.listTables("lab").collect()
print(f"Таблиц в lab базе после DROP: {len(remaining)}")
# Проверяем файлы на диске
print("\nПроверка физических файлов:")
# Managed: файлы должны исчезнуть
managed_check = subprocess.run(
["ls", managed_location.replace("file://", "")],
capture_output=True, text=True
)
if managed_check.returncode != 0:
print(f"✅ Managed файлы УДАЛЕНЫ: {managed_check.stderr.strip()}")
else:
print(f"⚠️ Managed файлы ОСТАЛИСЬ (в корзине): {managed_check.stdout[:100]}")
# External: файлы должны остаться!
external_check = subprocess.run(
["ls", "-la", external_location],
capture_output=True, text=True
)
if external_check.returncode == 0:
lines = external_check.stdout.strip().split("\n")
print(f"✅ External файлы ОСТАЛИСЬ ({len(lines)-1} объектов):")
print(f" {external_location}")
else:
print(f"❌ Ошибка проверки External файлов: {external_check.stderr}")
Шаг 4: Воссоздание таблицы из оставшихся файлов¶
# ── ВОССОЗДАНИЕ External таблицы из оставшихся файлов ─────────────────
print("\n=== ВОССОЗДАНИЕ ТАБЛИЦЫ ИЗ ФАЙЛОВ ===")
# Шаг 4.1: Создаём новую External таблицу поверх оставшихся файлов
spark.sql(f"""
CREATE TABLE lab.transactions_recovered (
txn_id STRING,
user_id BIGINT,
amount DOUBLE,
currency STRING
)
USING PARQUET
PARTITIONED BY (txn_date DATE)
LOCATION '{external_location}'
""")
# Шаг 4.2: Сразу после создания - таблица пустая!
# HMS не знает о партициях на HDFS
empty_count = spark.sql("SELECT COUNT(*) FROM lab.transactions_recovered").first()[0]
print(f"Строк сразу после создания (без MSCK): {empty_count}")
# Вывод: 0 (таблица «видит» только то, что зарегистрировано в HMS)
spark.sql("SHOW PARTITIONS lab.transactions_recovered").show()
# +------ (пусто!)
# Шаг 4.3: MSCK REPAIR TABLE - синхронизируем HMS с HDFS
print("\nЗапускаем MSCK REPAIR TABLE...")
spark.sql("MSCK REPAIR TABLE lab.transactions_recovered")
# Теперь партиции видны!
print("Партиции после MSCK REPAIR:")
spark.sql("SHOW PARTITIONS lab.transactions_recovered").show(5)
# И данные доступны!
recovered_count = spark.sql("SELECT COUNT(*) FROM lab.transactions_recovered").first()[0]
print(f"\n✅ ДАННЫЕ ВОССТАНОВЛЕНЫ: {recovered_count:,} строк")
assert recovered_count == external_count, \
f"Ожидалось {external_count:,} строк, получено {recovered_count:,}"
print("\n=== ИТОГ ЛАБОРАТОРНОЙ РАБОТЫ ===")
print(f"Managed таблица после DROP: ДАННЫЕ ПОТЕРЯНЫ")
print(f"External таблица после DROP: ДАННЫЕ СОХРАНЕНЫ ({recovered_count:,} строк)")
print(f"Восстановление через CREATE EXTERNAL + MSCK REPAIR: УСПЕШНО")
spark.stop()
Полная таблица сравнения: шпаргалка для принятия решений¶
| Критерий | Managed Table | External Table |
|---|---|---|
| DROP TABLE удаляет файлы | ДА (перемещает в .Trash) | НЕТ (только HMS запись) |
| DROP TABLE PURGE | Немедленное удаление | Только HMS запись |
| Физический путь | Внутри warehouse.dir | Любой (задаётся LOCATION) |
| Несколько таблиц → один путь | НЕТ | ДА |
| HMS знает об изменениях вне Spark | Автоматически | MSCK REPAIR нужен |
| Рекомендуется для Bronze | НЕТ (риск потери) | ДА |
| Рекомендуется для Gold | Допустимо | Предпочтительно |
| Просмотр типа | DESCRIBE EXTENDED | DESCRIBE EXTENDED |
| Воссоздание после DROP | Только из backup | CREATE + MSCK REPAIR |
| Lifecycle Management | Автоматический | Ручной |
| Multi-engine Access | Рисковано | Безопасно |
Итоги: правила выбора типа таблицы¶
Ключевой вопрос при создании любой таблицы: «Что должно произойти с данными если эта таблица будет удалена?»
- Если данные должны исчезнуть → Managed
- Если данные должны остаться → External
В современной практике Data Engineering преобладает подход «External для всего»: это даёт максимальную защиту данных, независимость от конкретного HMS, возможность переключаться между движками (Spark, Trino, Flink) и полный контроль над lifecycle данных на уровне платформы, а не отдельного SQL-движка.
Единственный аргумент в пользу Managed - упрощённый lifecycle для вычислимых артефактов (агрегаты, отчёты, временные промежуточные результаты), которые можно легко пересоздать и которые не нужно хранить вечно.