Dynamic Partition Overwrite: partitionOverwriteMode=dynamic - идемпотентная перезапись разделов

Полный разбор Dynamic Partition Overwrite в Spark: проблема статической перезаписи и потеря истории, концепция точечной перезаписи только затронутых партиций, конфигурация partitionOverwriteMode, физика FileOutputCommitter, идемпотентность пайплайнов, Partition Explosion и антипаттерны, практическая реализация Safe Retry в Airflow.

storage platform

1. Проблема статической перезаписи партиций в Spark (Static Mode)

Понять Dynamic Partition Overwrite невозможно без детального разбора того, что происходит при стандартной записи с mode("overwrite"). Это поведение является одним из самых контринтуитивных в Spark и регулярно приводит к катастрофическим потерям исторических данных в production.

Поведение по умолчанию: что делает Static Overwrite

Когда разработчик пишет следующий код:

df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/data/silver/events/")

интуитивное ожидание: «Spark перезапишет только те партиции, которые присутствуют в df». Но по умолчанию происходит нечто другое.

Spark в режиме static (дефолтном) перед записью новых данных удаляет всю корневую директорию /data/silver/events/ целиком, а затем записывает данные из df. Если df содержит данные только за 2024-01-15, то все исторические партиции за прошлые дни, недели и годы будут уничтожены.

Схема показывает масштаб катастрофы: при попытке «обновить один день» разработчик уничтожает историю за 4 дня. На production кластере это может означать потерю данных за месяцы или годы.

Почему Spark ведёт себя именно так

Это не баг, а намеренное дизайн-решение, унаследованное из Hive. В Hive SQL команда INSERT OVERWRITE TABLE t PARTITION (date='2024-01-14') явно указывает одну партицию для перезаписи. В Spark DataFrame API нет явного указания партиций - только общий mode("overwrite"). Spark не знает заранее какие партиции находятся в DataFrame до выполнения Job'а, поэтому по умолчанию выбирает «безопасный» вариант: удалить всё и записать заново.

Кроме того, поведение зависит от того, используется ли saveAsTable (запись в HMS таблицу) или прямая запись на путь:

  • df.write.mode("overwrite").parquet("/path/") - удаляет корневой путь
  • df.write.mode("overwrite").saveAsTable("db.table") - удаляет данные всех партиций управляемой таблицы (аналогично)

Исторические workaround'ы: как инженеры решали проблему вручную

До появления Dynamic Partition Overwrite разработчики использовали несколько опасных обходных путей.

Workaround 1: Ручное удаление партиций через HDFS CLI перед записью

# Опасный антипаттерн!
import subprocess

# Перед записью удаляем только нужные партиции
affected_dates = df.select("event_date").distinct().collect()
for row in affected_dates:
    date = row["event_date"]
    path = f"hdfs://cluster/data/silver/events/event_date={date}/"
    subprocess.run(["hdfs", "dfs", "-rm", "-r", path])

# Теперь пишем с append (не overwrite!)
df.write.mode("append").partitionBy("event_date").parquet(path_root)

Это решение имеет критические проблемы:

  • Нет атомарности: между удалением партиции и записью новых данных есть временной промежуток, когда данные недоступны. Параллельно работающий аналитический запрос получит пустой результат.
  • Data Race: если два экземпляра пайплайна запущены одновременно (например, при retry в Airflow), они могут конкурировать за удаление/запись одних и тех же партиций.
  • Ошибка в логике: если df содержит некорректные данные в ключе партиционирования (например, null или будущие даты), subprocess удалит неожиданные директории.

Workaround 2: INSERT OVERWRITE через Spark SQL с явным указанием партиции

# Немного лучше, но требует знания конкретной даты заранее
# и не масштабируется при многодневных пересчётах
spark.sql("""
  INSERT OVERWRITE TABLE silver.events
  PARTITION (event_date='2024-01-14')
  SELECT * FROM temp_increment
  WHERE event_date = '2024-01-14'
""")

Это работает для одной известной партиции, но неудобно при пересчёте нескольких дат. Нужно генерировать динамический SQL для каждой партиции - ненадёжно и медленно.


2. Концепция Dynamic Partition Overwrite

Dynamic Partition Overwrite решает проблему элегантно и безопасно: вместо удаления всей директории или ручного управления партициями, Spark анализирует входящий DataFrame, определяет какие конкретно партиционные значения в нём присутствуют, и удаляет только соответствующие директории.

Суть паттерна: изоляция изменений

Ключевая гарантия Dynamic Partition Overwrite формулируется так: «Spark затронет только те физические директории в HDFS, которые соответствуют партициям, присутствующим во входящем DataFrame».

Безопасность для истории: математическая гарантия

Если ваш инкремент содержит данные за даты {D1, D2, D3}, то Dynamic Partition Overwrite гарантирует:

  • Директории D1, D2, D3 будут полностью заменены новыми данными из инкремента
  • Все остальные директории (за любые другие даты) остаются нетронутыми
  • Нет риска случайного удаления партиции, которую вы не собирались трогать

Идемпотентность пайплайна: почему это важно для Airflow

Идемпотентность - свойство операции давать одинаковый результат при многократном применении. Для ETL-пайплайна это означает: сколько бы раз ни перезапустить пайплайн за конкретный день, результат всегда будет одинаковым - никаких дубликатов, никаких расхождений.

Dynamic Partition Overwrite делает пайплайн идемпотентным из коробки:

Первый запуск за 2024-01-14:
- Удаляет event_date=2024-01-14/ (если есть старые данные)
- Записывает новые данные за 2024-01-14
- Результат: 1 000 000 строк за эту дату

Повторный запуск за 2024-01-14 (после сбоя):
- Снова удаляет event_date=2024-01-14/
- Снова записывает те же данные
- Результат: 1 000 000 строк (те же, без дублей!)

N-й повторный запуск:
- Тот же результат, что и после первого запуска

Это критически важно в контексте Airflow: при сбое Task Airflow автоматически перезапускает её. Без идемпотентности каждый повторный запуск может создавать дублирующиеся данные (в режиме append) или стирать историю (в режиме static overwrite). Dynamic Partition Overwrite решает обе проблемы.


3. Активация и конфигурация динамического режима

Конфигурационный параметр: spark.sql.sources.partitionOverwriteMode

Единственный параметр, управляющий режимом перезаписи партиций:

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
# или
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")  # дефолт

Два допустимых значения:

  • static (по умолчанию): удалить всю директорию/таблицу перед записью
  • dynamic: удалить только те партиции, которые присутствуют во входящем DataFrame

Три способа установить конфигурацию

Способ 1: Глобально при создании SparkSession

Это лучший подход для production пайплайнов - конфигурация применяется ко всем операциям записи в сессии:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("daily-etl-pipeline") \
    .config("hive.metastore.uris", "thrift://hms:9083") \

    # Устанавливаем глобально: все write.mode("overwrite") в этой сессии
    # будут использовать dynamic режим
    .config("spark.sql.sources.partitionOverwriteMode", "dynamic") \

    .enableHiveSupport() \
    .getOrCreate()

Способ 2: Локально через spark.conf.set (наиболее гибкий)

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

# Для конкретной операции: включаем dynamic
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df_increment.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("silver.events")

# Для другой операции: возвращаем static
# (если нужно полностью пересоздать таблицу)
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")
df_full_reload.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("gold.aggregates")

Способ 3: Через spark-defaults.conf (кластерная конфигурация)

Устанавливается администратором кластера как дефолт для всех приложений:

# /opt/spark/conf/spark-defaults.conf
spark.sql.sources.partitionOverwriteMode dynamic

После этого все пайплайны на кластере будут использовать dynamic режим по умолчанию. Это снижает вероятность случайной потери истории, но требует от разработчиков явно переключаться в static когда нужна полная перезапись.

Важное ограничение: работает только с partitionBy

Dynamic Partition Overwrite работает только при наличии partitionBy(). Для неразбитых таблиц dynamic и static ведут себя одинаково - удаляют всё содержимое:

# НЕ партиционированная таблица:
# dynamic и static работают одинаково = удаляется вся таблица
df.write.mode("overwrite").parquet("/data/no-partitions/")

# Партиционированная таблица (partitionBy обязателен):
# dynamic = удаляет только затронутые партиции
df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \  # без этого dynamic не работает!
    .parquet("/data/events/")

4. Реализация паттерна на PySpark и Spark SQL

DataFrame API: стандартный паттерн для batch ETL

from pyspark.sql import SparkSession, functions as F
from datetime import date, timedelta

spark = SparkSession.builder \
    .appName("incremental-load") \
    .config("hive.metastore.uris", "thrift://hms:9083") \
    .config("spark.sql.sources.partitionOverwriteMode", "dynamic") \
    .enableHiveSupport() \
    .getOrCreate()


def process_daily_increment(
    processing_date: date,
    source_path: str,
    target_table: str,
) -> int:
    """
    Обрабатывает ежедневный инкремент с Dynamic Partition Overwrite.

    Гарантии:
    1. Только партиция processing_date будет перезаписана
    2. Все остальные партиции остаются нетронутыми
    3. Повторный запуск даёт тот же результат (идемпотентность)

    Возвращает: число строк в обработанной партиции
    """
    date_str = processing_date.isoformat()  # "2024-01-15"

    # Читаем только данные за нужную дату (не весь датасет!)
    # Это важно: мы не читаем то что не собираемся перезаписывать
    df = spark.read.parquet(source_path) \
        .filter(F.col("event_date") == date_str)

    # Трансформации
    df_cleaned = df \
        .dropDuplicates(["event_id"]) \
        .filter(F.col("event_type").isNotNull()) \
        .withColumn("event_date", F.col("event_date").cast("date"))

    row_count = df_cleaned.count()

    # Dynamic overwrite: перезапишет ТОЛЬКО partition event_date=processing_date
    # Все остальные партиции таблицы target_table останутся нетронутыми
    df_cleaned.write \
        .mode("overwrite") \
        .format("parquet") \
        .partitionBy("event_date") \
        .saveAsTable(target_table)

    print(f"✅ Обработана дата {date_str}: {row_count:,} строк")
    return row_count


# Пример: обработка вчерашнего дня
yesterday = date.today() - timedelta(days=1)
count = process_daily_increment(
    processing_date=yesterday,
    source_path="hdfs://cluster/data/bronze/events/",
    target_table="silver.events",
)

Многодневный пересчёт: backfill

Dynamic Partition Overwrite отлично подходит для backfill - пересчёта исторических данных за несколько дней:

from datetime import date, timedelta
from pyspark.sql import SparkSession, functions as F


def backfill_date_range(
    spark: SparkSession,
    start_date: date,
    end_date: date,
    source_table: str,
    target_table: str,
) -> dict:
    """
    Пересчитывает данные за диапазон дат.
    Каждая дата обрабатывается отдельно с Dynamic Partition Overwrite.

    Почему не одним большим DataFrame?
    - Контроль прогресса: видно какая дата обрабатывается
    - Управление памятью: небольшие DataFrame по одному дню
    - Retry granularity: при сбое перезапускается только упавшая дата

    Альтернатива: читать все даты сразу в один большой DF
    и делать один write - тоже работает с dynamic overwrite,
    но хуже управляется при сбоях.
    """
    results = {}
    current = start_date

    while current <= end_date:
        date_str = current.isoformat()
        print(f"Обрабатываем дату: {date_str}")

        df = spark.table(source_table) \
            .filter(F.col("event_date") == date_str)

        processed = df.groupBy("user_id", "event_type", "event_date") \
            .agg(
                F.count("*").alias("event_count"),
                F.sum("amount").alias("total_amount"),
            )

        row_count = processed.count()

        # Каждая дата перезаписывает только свою партицию
        processed.write \
            .mode("overwrite") \
            .partitionBy("event_date") \
            .saveAsTable(target_table)

        results[date_str] = row_count
        current += timedelta(days=1)

    return results


# Backfill за последние 7 дней
today = date.today()
stats = backfill_date_range(
    spark=spark,
    start_date=today - timedelta(days=7),
    end_date=today - timedelta(days=1),
    source_table="silver.events",
    target_table="gold.daily_metrics",
)
print("Backfill результаты:", stats)

Spark SQL: INSERT OVERWRITE с динамическими партициями

В Spark SQL команда INSERT OVERWRITE TABLE при включённом dynamic режиме автоматически использует Dynamic Partition Overwrite:

# Убеждаемся что dynamic режим включён
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

# Hive-совместимый синтаксис INSERT OVERWRITE
# В dynamic режиме перезапишет ТОЛЬКО партиции присутствующие в SELECT
spark.sql("""
  INSERT OVERWRITE TABLE silver.events
  PARTITION (event_date)
  SELECT
    event_id,
    user_id,
    event_type,
    amount,
    session_id,
    event_date
  FROM bronze.events_raw
  WHERE event_date = '2024-01-15'
    AND event_type IS NOT NULL
""")

# Для нескольких дат одновременно:
spark.sql("""
  INSERT OVERWRITE TABLE silver.events
  PARTITION (event_date)
  SELECT
    event_id, user_id, event_type, amount, session_id, event_date
  FROM bronze.events_raw
  WHERE event_date BETWEEN '2024-01-10' AND '2024-01-15'
""")
# Перезапишет 6 партиций: 2024-01-10 через 2024-01-15
# Все остальные партиции серебряного слоя останутся нетронутыми!

Разница поведения для Managed и External таблиц

Dynamic Partition Overwrite работает одинаково для обоих типов таблиц с точки зрения данных, но есть нюансы на уровне HMS:

# MANAGED TABLE: Dynamic Overwrite
# Spark управляет и файлами и метаданными
# DROP partition = физическое удаление файлов из warehouse.dir
# ADD partition = регистрация в HMS + файлы в warehouse.dir

df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("gold.managed_metrics")
# HMS автоматически обновляет список партиций

# EXTERNAL TABLE: Dynamic Overwrite
# Spark управляет только метаданными
# DROP partition = удаление из HMS + физическое удаление файлов по LOCATION

df.write \
    .mode("overwrite") \
    .option("path", "hdfs://cluster/data/gold/metrics/") \
    .partitionBy("event_date") \
    .saveAsTable("gold.external_metrics")
# HMS обновляет партиции, физические файлы по указанному LOCATION
# Важно: если использовать .parquet("hdfs://...") без saveAsTable,
# партиции в HMS НЕ обновятся - нужен MSCK REPAIR!

5. Под капотом Spark: физика процесса и FileOutputCommitter

Чтобы понять почему Dynamic Partition Overwrite безопасен и почему иногда возникают проблемы, нужно разобраться что именно происходит на уровне Spark Execution Engine и HDFS.

Временные директории: как Spark пишет без race conditions

При запуске Job'а с write.mode("overwrite").partitionBy(...) каждый Executor пишет данные не сразу в целевые директории, а в специальные временные директории _temporary:

/data/silver/events/
├── event_date=2024-01-10/  ← существующие данные
├── event_date=2024-01-11/  ← существующие данные
└── _temporary/
    └── 0/                  ← Job attempt ID
        ├── task_001_00/    ← Task attempt
        │   └── event_date=2024-01-14/
        │       └── part-00001.parquet
        └── task_002_00/
            └── event_date=2024-01-14/
                └── part-00002.parquet

Это означает, что пока Job выполняется, существующие данные в event_date=2024-01-10/ и event_date=2024-01-11/ остаются нетронутыми и читаемыми. Параллельные аналитические запросы получают корректные данные даже в процессе записи инкремента.

Стадия Commit: алгоритм FileOutputCommitter

Когда все Tasks завершены успешно, Spark Driver запускает Job Commit - финальную фазу, которая атомарно публикует результат:

Ключевой момент шага 2: Spark удаляет директорию event_date=2024-01-14/ только потому, что нашёл её в _temporary/. Директории event_date=2024-01-10/ и другие никогда не упоминаются в этом алгоритме - они просто остаются в покое.

Ключевой момент шага 4: rename() в HDFS - это атомарная операция на уровне NameNode (изменение одной записи в пространстве имён). Именно поэтому переход «старые данные удалены, новые ещё не на месте» виден только в течение микросекунд между шагами 2 и 4.

Нагрузка на NameNode при Dynamic Partition Overwrite

Каждая операция Dynamic Partition Overwrite генерирует следующую нагрузку на NameNode:

  1. N операций delete() где N = число затронутых партиций × среднее число файлов на партицию
  2. N операций mkdir() для создания новых директорий партиций
  3. M операций rename() где M = суммарное число новых файлов во всех партициях

Для типичного дневного инкремента (1 партиция × 20 файлов):

  • delete() × 20 (файлы старой партиции)
  • mkdir() × 1
  • rename() × 20

Итого: ~41 операция на NameNode - минимальная нагрузка.

Для большого backfill (30 партиций × 20 файлов):

  • delete() × 600
  • mkdir() × 30
  • rename() × 600

Итого: ~1230 операций - тоже разумно.


6. Архитектурные ловушки и Data Skew

Partition Explosion: когда Dynamic Overwrite опасен

Partition Explosion - ситуация, когда в одном инкременте оказываются данные за огромное число уникальных партиционных значений. Это может произойти из-за:

  • Битых данных: event_date = null или event_date = '1970-01-01' для тысяч записей
  • Высококардинального ключа партиционирования: user_id вместо event_date
  • CDC потока с данными за несколько месяцев в одном батче

Паттерн стабилизации: repartition перед записью

Ключевое решение для предотвращения Partition Explosion - явное репартиционирование по колонке партиционирования перед записью. Это гарантирует что каждый Executor будет писать данные только из тех партиций, которые хэш-функция назначила именно ему:

# Проблемный код: каждый Executor может получить данные из всех партиций
df_increment.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("silver.events")
# При 10k партиций: каждый из 50 Executors открывает до 10k файлов!

# Оптимизированный код с repartition по ключу партиционирования
df_increment \
    .repartition("event_date") \          # Executor получает данные только своих дат
    .write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("silver.events")
# Каждый Executor пишет только в свои партиции
# Максимум открытых файлов: число партиций / число Executors

# Для коротких задач: sortWithinPartitions менее агрессивен по сети
df_increment \
    .sortWithinPartitions("event_date") \  # Сортировка без shuffle
    .write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("silver.events")

# Контролируем максимальное число файлов на партицию:
spark.conf.set("spark.sql.files.maxRecordsPerFile", "1000000")  # 1M строк на файл

Защита от битых данных в ключе партиционирования

from pyspark.sql import functions as F

def safe_dynamic_overwrite(
    spark,
    df,
    target_table: str,
    partition_col: str,
    expected_partitions: list[str] | None = None,
) -> int:
    """
    Безопасная обёртка для Dynamic Partition Overwrite.
    Включает валидацию данных и защиту от Partition Explosion.

    expected_partitions: если задан - проверяем что DF содержит
    только ожидаемые партиции, иначе выбрасываем исключение.
    """
    # Шаг 1: Найти null значения в ключе партиционирования
    null_count = df.filter(F.col(partition_col).isNull()).count()
    if null_count > 0:
        raise ValueError(
            f"Найдено {null_count:,} строк с NULL в ключе партиционирования '{partition_col}'. "
            f"Невозможно записать в партицию 'null' - это может создать миллионы файлов. "
            f"Добавьте фильтр: df.filter(F.col('{partition_col}').isNotNull())"
        )

    # Шаг 2: Проверить число уникальных значений
    actual_partitions = {
        row[partition_col]
        for row in df.select(partition_col).distinct().collect()
    }

    max_partitions = 100  # Разумный предел для batch инкремента
    if len(actual_partitions) > max_partitions:
        raise ValueError(
            f"Обнаружено {len(actual_partitions)} уникальных значений партиционирования. "
            f"Это может вызвать Partition Explosion. "
            f"Максимально допустимое: {max_partitions}. "
            f"Пример значений: {list(actual_partitions)[:5]}..."
        )

    # Шаг 3: Если задан ожидаемый список - проверяем что нет неожиданных
    if expected_partitions:
        unexpected = actual_partitions - set(expected_partitions)
        if unexpected:
            raise ValueError(
                f"Неожиданные партиции в данных: {unexpected}. "
                f"Ожидались: {expected_partitions}"
            )

    print(f"Валидация пройдена: {len(actual_partitions)} партиций → {target_table}")
    print(f"Партиции: {sorted(actual_partitions)}")

    # Шаг 4: Репартиционирование для стабильности
    row_count = df.count()
    files_per_partition = max(1, row_count // (len(actual_partitions) * 1_000_000))
    total_partitions = len(actual_partitions) * files_per_partition

    df_repartitioned = df.repartition(total_partitions, partition_col)

    # Шаг 5: Dynamic Partition Overwrite
    spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
    df_repartitioned.write \
        .mode("overwrite") \
        .partitionBy(partition_col) \
        .saveAsTable(target_table)

    print(f"✅ Успешно записано: {row_count:,} строк в {len(actual_partitions)} партиций")
    return row_count

Многоуровневое партиционирование: дополнительные сложности

При многоуровневом партиционировании (например, partitionBy("year", "month", "day")) Dynamic Partition Overwrite удаляет только листовые партиции, которые встречаются в DataFrame:

# Многоуровневое партиционирование
df.write \
    .mode("overwrite") \
    .partitionBy("year", "month", "day") \
    .saveAsTable("silver.events_multilevel")

# Если df содержит только year=2024/month=01/day=15:
# Удалится: /events_multilevel/year=2024/month=01/day=15/
# НЕ удалится: year=2024/month=01/day=14/ (другой день в том же месяце)
# НЕ удалится: year=2024/month=01/ (родительская директория)
# НЕ удалится: year=2023/ (прошлый год)

# Это именно то поведение которое нам нужно!

7. Практика: лабораторная реализация Safe Retry и траблшутинг

Демонстрация Static vs Dynamic: наглядное сравнение

# lab_partition_overwrite.py
# Демонстрирует разницу между static и dynamic partition overwrite

import subprocess
from datetime import date, timedelta
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .master("local[2]") \
    .appName("lab-dynamic-overwrite") \
    .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_dpo")

def create_week_of_data(table: str) -> None:
    """Создаёт тестовую таблицу с данными за 7 дней."""
    spark.sql(f"DROP TABLE IF EXISTS {table}")

    all_dfs = []
    for i in range(7):
        d = date(2024, 1, 10) + timedelta(days=i)
        df_day = spark.range(1000).select(
            F.col("id"),
            F.lit(str(d)).alias("event_date"),
            (F.rand() * 100).alias("amount"),
        )
        all_dfs.append(df_day)

    from functools import reduce
    from pyspark.sql import DataFrame
    full_df: DataFrame = reduce(DataFrame.union, all_dfs)

    full_df.write \
        .mode("overwrite") \
        .partitionBy("event_date") \
        .saveAsTable(table)

    print(f"Создана таблица {table}: {full_df.count():,} строк за 7 дней")


def count_by_date(table: str) -> dict:
    """Считает строки по каждой дате."""
    rows = spark.sql(f"SELECT event_date, COUNT(*) as cnt FROM {table} GROUP BY event_date ORDER BY event_date").collect()
    return {row["event_date"]: row["cnt"] for row in rows}


# ── Шаг 1: СТАТИЧЕСКИЙ OVERWRITE (катастрофа) ────────────────────────
print("=" * 60)
print("ТЕСТ 1: Static Overwrite (mode по умолчанию)")
print("=" * 60)

create_week_of_data("lab_dpo.sales_static")
print(f"До записи: {count_by_date('lab_dpo.sales_static')}")

# Инкремент только за один день
increment = spark.range(500).select(
    F.col("id"),
    F.lit("2024-01-14").alias("event_date"),
    (F.rand() * 200).alias("amount"),
)

# Static overwrite (дефолт!)
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")
increment.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("lab_dpo.sales_static")

after_static = count_by_date("lab_dpo.sales_static")
print(f"После Static overwrite: {after_static}")
print(f"КАТАСТРОФА: осталось {len(after_static)} дат вместо 7!")
# Вывод: {'2024-01-14': 500}  ← только один день!

# ── Шаг 2: ДИНАМИЧЕСКИЙ OVERWRITE (безопасно) ────────────────────────
print("\n" + "=" * 60)
print("ТЕСТ 2: Dynamic Overwrite (безопасный режим)")
print("=" * 60)

create_week_of_data("lab_dpo.sales_dynamic")
before_dynamic = count_by_date("lab_dpo.sales_dynamic")
print(f"До записи: {before_dynamic}")

# Тот же инкремент за один день
increment = spark.range(500).select(
    F.col("id"),
    F.lit("2024-01-14").alias("event_date"),
    (F.rand() * 200).alias("amount"),
)

# Dynamic overwrite
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
increment.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("lab_dpo.sales_dynamic")

after_dynamic = count_by_date("lab_dpo.sales_dynamic")
print(f"После Dynamic overwrite: {after_dynamic}")
print(f"✅ Сохранено {len(after_dynamic)} дат (все 7!)")
print(f"✅ Дата 2024-01-14 обновлена: {before_dynamic['2024-01-14']}{after_dynamic['2024-01-14']}")

# Проверяем что остальные даты не изменились
for d, cnt in before_dynamic.items():
    if d != "2024-01-14":
        assert after_dynamic[d] == cnt, f"Дата {d} изменилась!"
print("✅ Все исторические даты нетронуты!")

Тест идемпотентности: Safe Retry симуляция

# Продолжение lab_partition_overwrite.py

# ── Шаг 3: ИДЕМПОТЕНТНОСТЬ (Safe Retry) ─────────────────────────────
print("\n" + "=" * 60)
print("ТЕСТ 3: Идемпотентность - симуляция повторного запуска Airflow")
print("=" * 60)

# Создаём стабильный инкремент для тестирования
fixed_increment = spark.createDataFrame([
    (1001, "2024-01-14", 100.0),
    (1002, "2024-01-14", 200.0),
    (1003, "2024-01-14", 150.0),
], ["id", "event_date", "amount"])

def run_increment(run_number: int) -> int:
    """Симулирует один запуск ETL Task в Airflow."""
    print(f"\nЗапуск #{run_number}:")
    spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
    fixed_increment.write \
        .mode("overwrite") \
        .partitionBy("event_date") \
        .saveAsTable("lab_dpo.sales_dynamic")

    count = spark.sql("""
        SELECT COUNT(*) FROM lab_dpo.sales_dynamic
        WHERE event_date = '2024-01-14'
    """).first()[0]

    print(f"  Строк за 2024-01-14: {count}")
    return count


# Запускаем 5 раз (симуляция 5 retry после сбоя)
results = [run_increment(i) for i in range(1, 6)]

print(f"\nРезультаты по запускам: {results}")
assert all(r == results[0] for r in results), "Результаты отличаются - нет идемпотентности!"
print("✅ Все 5 запусков дали одинаковый результат - идемпотентность подтверждена!")

Паттерн для Airflow DAG: эталонный Safe ETL

# dags/daily_silver_load.py
# Airflow DAG с идемпотентной Dynamic Partition Overwrite

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator

DEFAULT_ARGS = {
    "owner": "data-platform",
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "email_on_failure": True,
}


def process_silver_events(ds: str, **context) -> None:
    """
    Обрабатывает Bronze → Silver события за дату ds.

    ds: date string от Airflow (YYYY-MM-DD)
    Идемпотентно: безопасный повтор при любом числе retry.

    Ключевой момент: мы используем ds (execution_date) как фильтр
    и для записи. При retry Airflow передаёт тот же ds,
    поэтому Dynamic Overwrite перезапишет ту же партицию
    теми же данными - дублей нет.
    """
    from pyspark.sql import SparkSession, functions as F

    spark = SparkSession.builder \
        .appName(f"silver-events-{ds}") \
        .config("spark.sql.sources.partitionOverwriteMode", "dynamic") \
        .config("hive.metastore.uris", "thrift://hms:9083") \
        .enableHiveSupport() \
        .getOrCreate()

    try:
        print(f"Обрабатываем дату: {ds}")

        # Читаем только нужную дату из Bronze
        df = spark.table("bronze.events_raw") \
            .filter(F.col("event_date") == ds)

        row_count = df.count()
        if row_count == 0:
            print(f"⚠️  Нет данных за {ds} - пропускаем")
            return

        print(f"Строк в Bronze за {ds}: {row_count:,}")

        # Трансформации Silver-слоя
        silver_df = df \
            .dropDuplicates(["event_id"]) \
            .filter(F.col("event_type").isNotNull()) \
            .withColumn("processed_at", F.current_timestamp()) \
            .repartition("event_date")  # Предотвращение Partition Explosion

        # Dynamic Partition Overwrite: атомарная идемпотентная запись
        silver_df.write \
            .mode("overwrite") \
            .partitionBy("event_date") \
            .saveAsTable("silver.events")

        final_count = spark.sql(f"""
            SELECT COUNT(*) FROM silver.events WHERE event_date = '{ds}'
        """).first()[0]

        print(f"✅ Silver: {final_count:,} строк за {ds}")

    finally:
        spark.stop()


with DAG(
    dag_id="daily_silver_events_load",
    default_args=DEFAULT_ARGS,
    description="Ежедневная загрузка Bronze → Silver",
    schedule_interval="@daily",
    start_date=datetime(2024, 1, 1),
    catchup=True,    # Catchup=True: отработает пропущенные даты при рестарте
    max_active_runs=3,
) as dag:

    process_silver = PythonOperator(
        task_id="process_silver_events",
        python_callable=process_silver_events,
        doc_md="""
        Загружает события из Bronze в Silver с Dynamic Partition Overwrite.

        Идемпотентно: безопасный retry при любом числе попыток.
        Партиции: event_date.
        При catchup обрабатывает каждую дату независимо.
        """,
    )

Траблшутинг: типичные проблемы и решения

# Типичные ошибки и их исправление

# ── ОШИБКА 1: Too many open files ────────────────────────────────────
# Симптом:
# java.io.IOException: Too many open files
#   at sun.nio.ch.FileDispatcherImpl.write0(FileDispatcherImpl.java:62)

# Причина: Partition Explosion или отсутствие repartition

# Диагностика: проверить число уникальных партиций
partition_values = df.select("event_date").distinct().count()
print(f"Уникальных партиций: {partition_values}")
# Если > 100 при небольшом инкременте - подозрительно!

# Исправление: repartition перед записью
df \
    .repartition("event_date") \  # Каждый Executor - свои партиции
    .write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("target")


# ── ОШИБКА 2: Dynamic Overwrite записывает в static режиме ───────────
# Симптом: история удаляется несмотря на dynamic настройку

# Причина 1: конфигурация установлена ПОСЛЕ создания SparkSession
# Конфигурация spark.sql.sources.partitionOverwriteMode должна быть
# установлена ДО первого использования!

# НЕПРАВИЛЬНО:
spark_wrong = SparkSession.builder.getOrCreate()
spark_wrong.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
# Поздно - сессия уже создана со статическим режимом!

# ПРАВИЛЬНО:
spark_right = SparkSession.builder \
    .config("spark.sql.sources.partitionOverwriteMode", "dynamic") \
    .getOrCreate()

# Причина 2: параметр действует только для DataSource v1 (Parquet, ORC)
# Для Delta Lake и Iceberg - используйте их собственные механизмы!


# ── ОШИБКА 3: Партиции не обновляются в HMS ──────────────────────────
# Симптом: файлы на HDFS обновились, но SELECT возвращает старые данные

# Причина: использование df.write.parquet(path) вместо saveAsTable
# При прямой записи на путь HMS не обновляется!

# НЕПРАВИЛЬНО (если нужна HMS синхронизация):
df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/data/silver/events/")  # Только HDFS, не HMS!

spark.sql("SELECT * FROM silver.events WHERE event_date='2024-01-14'").show()
# Может вернуть устаревшие данные (HMS кеш не обновлён)!

# ПРАВИЛЬНО для HMS-таблиц:
df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .saveAsTable("silver.events")  # Обновляет и HDFS и HMS!

# Или: после прямой записи принудительно обновляем HMS
spark.sql("MSCK REPAIR TABLE silver.events")
# или
spark.catalog.refreshTable("silver.events")

Итоги: когда использовать Dynamic Partition Overwrite

Dynamic Partition Overwrite - это правильный выбор для подавляющего большинства batch ETL-пайплайнов в Data Lake:

  • Ежедневные инкременты Bronze → Silver → Gold: каждый день пересчитывается только один или несколько дней
  • Backfill исторических данных: пересчёт нескольких дней/недель без риска задеть другую историю
  • Retry-безопасные пайплайны: повторный запуск Airflow Task не создаёт дублей
  • Инкрементальные обновления: данные за конкретный день изменились (late-arriving events, CDC)

Когда НЕ использовать Dynamic Partition Overwrite:

  • Полная пересборка таблицы (schema migration, полный recalculation): используйте static для явного удаления всех данных
  • Row-level updates внутри партиции: используйте Iceberg/Delta Lake MERGE INTO
  • SCD Type 2 / CDC с историей: используйте Iceberg MERGE INTO или Delta Lake MERGE
  • Streaming: Dynamic Partition Overwrite несовместим с Structured Streaming - используйте foreachBatch с ручным управлением партициями

Правило «установить один раз и забыть»: в production спарке рекомендуется устанавливать spark.sql.sources.partitionOverwriteMode = dynamic глобально в spark-defaults.conf. Это делает поведение предсказуемым и защищает от случайного удаления истории. Для случаев когда нужен static - указывайте явно через spark.conf.set.