Dead Letter Queue: паттерн Quarantine для плохих записей

Изоляция некорректных записей без остановки пайплайна: архитектура DLQ в Medallion, Single-Pass Splitter, схема карантинной таблицы, rescued_data, try_cast, replay-пайплайны и мониторинг quarantine-rate

streaming

Введение: плохие данные неизбежны - вопрос только в том, как с ними работать

Представьте: пайплайн ingestion обрабатывает 10 миллионов событий из Kafka за ночной батч. В 3:47 AM одна из строк содержит "amount": "N/A" вместо числа. Наивный пайплайн падает с AnalysisException. Airflow перезапускает задачу - снова падение. Дежурный инженер просыпается в 4:15, смотрит в логи, находит одну битую строку и правит её вручную. 9 999 999 корректных строк ждут.

Это классическая проблема fail-fast в негибком варианте: мы жертвуем 100% данных из-за 0.00001% брака. В финансовых системах с жёсткими гарантиями - такой подход оправдан. В аналитических пайплайнах, обрабатывающих пользовательские события или логи, это катастрофа.

Dead Letter Queue (DLQ) - паттерн изоляции некорректных записей без остановки основного потока обработки. Плохие строки уходят в специальное хранилище-карантин, чистые - продолжают движение по пайплайну. Бизнес-процессы не прерываются, данные не теряются, инженеры видят полную картину проблем.


Часть 1. Происхождение и архитектурная концепция паттерна

1.1 Откуда пришёл Dead Letter Queue

Термин DLQ появился в системах message broker'ов задолго до эпохи Big Data. В RabbitMQ, ActiveMQ и позднее в Apache Kafka - сообщение, которое Consumer не смог обработать после N попыток, автоматически перемещалось в специальный топик (dead letter topic). Там его мог прочитать оператор, исправить и отправить в повторную обработку.

В Apache Kafka эта механика реализована через конфигурацию max.delivery.count и DLQ-топики. Consumer Group при превышении лимита повторов перестаёт пытаться обработать «ядовитое» сообщение и публикует его в orders.DLQ с метаданными об ошибке.

Когда data engineering переехал в Lakehouse-архитектуры, паттерн адаптировался: вместо messaging topic - Delta/Iceberg/Parquet таблица. Но концепция осталась той же: контролируемая изоляция проблемных данных с сохранением полной трассировки.

1.2 Три принципиально разных места для «плохих» данных

Прежде чем строить DLQ, нужно понять различие между тремя подходами:

Retry Queue - временное хранилище для данных, которые могут стать валидными при повторной обработке. Downstream-сервис временно недоступен? JDBC-соединение упало? Запись пойдёт в Retry Queue и будет обработана через 5 минут. Это не DLQ - это временный буфер.

Quarantine Storage (DLQ) - хранилище для данных, которые нарушают бизнес-правила или имеют структурные дефекты. Повторная попытка с тем же кодом не поможет. Нужно либо исправить данные на источнике, либо обновить логику валидации.

Permanent Drop - молчаливое удаление. Это антипаттерн в enterprise data engineering: данные теряются бесследно, аудит невозможен, расследование инцидентов - тоже.

1.3 Бизнес-ценность DLQ

DLQ создаёт три ключевые гарантии:

Непрерывность (SLA continuity) - 99.99% валидных данных обрабатываются в штатное время. Один битый JSON не останавливает витрину выручки.

100% сохранность сырых данных - ни одна строка не теряется. Плохая запись хранится в DLQ со всеми метаданными. После исправления бага в коде её можно переобработать.

Observability - DLQ - это инструмент диагностики. Резкий рост карантинных записей - ранний сигнал деградации upstream-системы. Мониторинг quarantine rate позволяет обнаружить проблему задолго до того, как она затронет BI-дашборды.


Часть 2. Fail-Fast vs Quarantine: когда что применять

2.1 Анатомия ошибок данных

Не все «плохие» данные одинаково опасны. Нужна классификация по двум осям: severity (критичность) и recoverability (возможность восстановления).

2.2 Матрица решений

Тип ошибки Пример Действие
Schema corruption Все колонки переименованы в источнике Fail-Fast
Пустой батч 0 записей из источника Fail-Fast
Null в Primary Key order_id = NULL Fail-Fast если > 1%, иначе DLQ
Invalid enum value status = "CANCELD" (опечатка) DLQ
Битый JSON в одной строке {...malformed...} DLQ
Дата из будущего created_at = 2099-01-01 DLQ + alert
Отрицательная сумма amount = -500 DLQ
Null в необязательном поле promo_code = NULL Pass (высокий null ratio норм)

2.3 Availability vs Correctness

Ключевой trade-off: доступность данных vs корректность данных.

Финансовые витрины (P&L, баланс) требуют Correctness first - лучше временно потерять данные, чем показать неверную выручку. Для таких пайплайнов Fail-Fast с автоматическим алертом предпочтительнее.

Операционные дашборды реального времени (количество онлайн-пользователей, активность сессий) требуют Availability first - 0.1% брака не искажает картину принципиально. DLQ + продолжение обработки - правильный выбор.


Часть 3. Архитектура Quarantine Layer в Medallion

3.1 Где живёт карантин

Quarantine Layer - это технический слой, параллельный Bronze. Он не часть Gold и не часть Silver - он находится между ними как «боковой карман»:

3.2 Классификация типов плохих записей

Каждый тип нарушения требует своей стратегии хранения и восстановления.

Parsing Errors - невозможно десериализовать данные. JSON битый, CSV с лишними запятыми, Avro с неправильным magic bytes. Такие записи надо сохранять в сыром виде (raw bytes), потому что применять схему к ним нельзя.

Schema Violations - данные пришли с неожиданной структурой. Новая колонка, переименованная колонка, смена типа. Spark может автоматически отловить такие случаи через _rescued_data (подробнее в части 5).

Business Rule Violations - структура корректна, но значение нарушает бизнес-инвариант. amount = -100, status = "UNKNWN", delivery_date < created_at. Такие записи сохраняются с полной бизнес-схемой плюс колонка с описанием нарушения.

Referential Integrity Violations - user_id = 99999 не существует в таблице users. Такие записи появляются при CDC race conditions или при удалении пользователей без каскадной обработки.


Часть 4. Проектирование схемы DLQ-таблицы

4.1 Антипаттерн: одна схема для чистых и грязных данных

Распространённая ошибка - сохранять «плохие» записи в ту же схему, что и «хорошие», просто в другую партицию:

# АНТИПАТТЕРН: нельзя применить схему Silver к битым данным
# Если amount = "N/A" - приведение к DoubleType упадёт
# Если order_id = NULL - нарушит NOT NULL constraint
bad_df.write.format("delta").mode("append").save(silver_path)

Проблема: схема Silver требует, чтобы amount был DoubleType. Но именно потому запись плохая, что amount = "N/A". Применение схемы к карантинному хранилищу создаст второй слой ошибок.

4.2 Правильная структура DLQ-таблицы

Карантинная таблица должна принимать данные любой формы, сохраняя сырой payload и обогащая его метаданными:

from pyspark.sql.types import (
    StructType, StructField,
    StringType, ArrayType, TimestampType, LongType
)

DLQ_SCHEMA = StructType([
    # Сырой payload - сохраняем всё как есть, без потерь
    StructField("raw_payload", StringType(), nullable=True),
    # Массив причин: может быть несколько нарушений одновременно
    StructField("error_reasons", ArrayType(StringType()), nullable=False),
    # Имя pipeline для трассировки
    StructField("pipeline_name", StringType(), nullable=False),
    # run_id из Airflow или уникальный UUID батча
    StructField("run_id", StringType(), nullable=True),
    # Имя источника данных
    StructField("source_system", StringType(), nullable=True),
    # Момент помещения в карантин
    StructField("quarantined_at", TimestampType(), nullable=False),
    # Дата обработки - для партиционирования
    StructField("processing_date", StringType(), nullable=False),
    # Категория ошибки для фильтрации
    StructField("error_category", StringType(), nullable=True),
    # Версия схемы пайплайна - важно для replay
    StructField("pipeline_version", StringType(), nullable=True),
    # Признак: была ли предпринята попытка replay
    StructField("replay_attempted", StringType(), nullable=True),
])

Ключевые поля:

raw_payload - исходная строка целиком (JSON, CSV-строка, Avro в Base64). Это единственный способ гарантировать возможность replay: мы восстанавливаем не интерпретацию данных, а сами данные.

error_reasons - массив, а не строка. Одна строка может нарушать сразу несколько правил (["NULL_USER_ID", "NEGATIVE_AMOUNT", "FUTURE_DATE"]). Массив позволяет агрегировать статистику по каждому типу ошибки независимо.

processing_date - партиционирование по дате обработки, а не по дате бизнес-события. Причина: карантин нужно чистить по времени хранения, а не по бизнес-логике. «Удалить карантин старше 90 дней» - простое условие. «Удалить карантин с event_date старше 90 дней» - ловушка при задержках данных.

pipeline_version - версия кода пайплайна в момент изоляции записи. При replay важно знать, какие правила действовали тогда. Иначе невозможно понять: запись была плохой или правила были слишком строгими?

4.3 Партиционирование и формат хранения

DLQ лучше хранить в Delta Lake или Apache Iceberg - это даёт ACID-транзакции при записи и DELETE/UPDATE при операциях replay:

# Структура партиций в DLQ:
# quarantine/
#   processing_date=2025-01-15/
#     error_category=parse_error/
#       part-00000.snappy.parquet
#     error_category=business_violation/
#       part-00001.snappy.parquet
#   processing_date=2025-01-16/
#     ...

DLQ_PARTITION_COLS = ["processing_date", "error_category"]

Двухуровневое партиционирование: сначала по дате (для TTL-очистки), потом по типу ошибки (для быстрого анализа «сколько parse errors за неделю»).


Часть 5. Встроенные механизмы Spark: Rescued Data и try_cast

5.1 _rescued_data: автоматический карантин схемы

Начиная со Spark 3.0 (и в Delta Lake), при чтении JSON/CSV с указанной схемой можно включить rescuedDataColumn. Все поля, которые не вписались в схему, автоматически собираются в специальную колонку _rescued_data в виде JSON-строки:

bronze_df = (
    spark.read
    .schema(BRONZE_ORDERS_SCHEMA)
    .option("rescuedDataColumn", "_rescued_data")  # магическая опция
    .json(bronze_path)
)

# Поле _rescued_data содержит JSON с "лишними" или "несовместимыми" полями
# Если все поля вписались в схему - _rescued_data == null
# Если пришло новое поле "partner_id" которого нет в схеме - оно там

# Смотрим, у кого есть rescued data (schema drift)
schema_drift_df = bronze_df.filter(F.col("_rescued_data").isNotNull())
schema_drift_count = schema_drift_df.count()

if schema_drift_count > 0:
    print(
        f"[DLQ] Schema drift detected: "
        f"{schema_drift_count} rows with unexpected fields"
    )
    # Пример содержимого _rescued_data:
    # {"partner_id": "P001", "new_field": "some_value"}

_rescued_data работает только для новых или несовместимых по типу полей. Если поле есть в схеме, но значение не приводится (например, "amount": "N/A" при DoubleType) - Spark вернёт null для этого поля и не положит его в _rescued_data. Для таких случаев нужен try_cast.

5.2 try_cast: безопасное приведение типов

try_cast - это SQL-функция, которая при невозможности приведения типа возвращает null вместо исключения. Это основа safe parsing в Spark:

from pyspark.sql.functions import expr, col, when, lit, current_timestamp

def safe_parse_bronze(df: DataFrame) -> DataFrame:
    """
    Безопасный парсинг Bronze с try_cast.
    Некорректные типы → null, а не exception.
    """
    return (
        df
        # Безопасное приведение amount: "N/A" -> null, "123.5" -> 123.5
        .withColumn(
            "amount_parsed",
            expr("try_cast(amount_raw as DOUBLE)")
        )
        # Безопасный парсинг даты: "2099-99-99" -> null
        .withColumn(
            "created_at_parsed",
            expr("try_cast(created_at_raw as TIMESTAMP)")
        )
        # Безопасное приведение user_id: "abc" -> null
        .withColumn(
            "user_id_parsed",
            expr("try_cast(user_id_raw as BIGINT)")
        )
    )

Ключевое различие: cast("N/A" as DOUBLE)null в Spark без ошибки (Spark не бросает exception при неудачном cast - это особенность Spark SQL). Но try_cast явно коммуницирует намерение и обеспечивает совместимость с разными версиями.

В Spark SQL есть ещё try_divide для деления, try_to_timestamp и другие try_* функции для безопасных операций.

5.3 columnNameOfCorruptRecord при чтении JSON

При чтении JSON с mode="PERMISSIVE" (режим по умолчанию) Spark помещает строки с синтаксическими ошибками JSON в специальную колонку:

from pyspark.sql.types import StructType, StructField, StringType

# Добавляем специальную колонку для битых JSON строк
SCHEMA_WITH_CORRUPT = BRONZE_ORDERS_SCHEMA.add(
    StructField("_corrupt_record", StringType(), nullable=True)
)

bronze_df = (
    spark.read
    .schema(SCHEMA_WITH_CORRUPT)
    .option("mode", "PERMISSIVE")                     # не падать, а сохранять
    .option("columnNameOfCorruptRecord", "_corrupt_record")  # куда класть битые строки
    .json(bronze_path)
)

# Синтаксически битые JSON-строки теперь в _corrupt_record
# Все остальные поля будут null для таких строк
parse_errors_df = bronze_df.filter(F.col("_corrupt_record").isNotNull())
valid_parse_df = bronze_df.filter(F.col("_corrupt_record").isNull())

print(f"Parse errors: {parse_errors_df.count()}")
print(f"Valid parse: {valid_parse_df.count()}")

Три режима чтения JSON:

  • PERMISSIVE (default): битые строки → _corrupt_record, остальные поля null
  • DROPMALFORMED: битые строки молча удаляются - антипаттерн (данные теряются)
  • FAILFAST: любая битая строка → exception - подходит только если источник обязан давать корректный JSON

Часть 6. Реализация Single-Pass Splitter

6.1 Почему важен Single-Pass

Наивная реализация DLQ выглядит так:

# АНТИПАТТЕРН: два отдельных прохода по данным
clean_df  = df.filter(is_valid_condition)    # Scan 1
invalid_df = df.filter(~is_valid_condition)  # Scan 2

На датасете размером 100GB это означает два полных сканирования: 200GB IO. В production с медленным S3 это может добавить 20-40 минут к времени выполнения батча.

Single-Pass Splitter - добавляем флаг валидности как колонку к тому же DataFrame, кэшируем его, затем делаем два filter по уже материализованным данным:

6.2 Полная реализация Single-Pass Splitter

from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
from pyspark.sql.functions import (
    col, when, lit, array, array_union, size,
    current_timestamp, to_json, struct
)
from dataclasses import dataclass
from typing import Optional
import json


@dataclass
class DLQConfig:
    """Конфигурация DLQ для конкретного пайплайна."""
    table_path: str           # путь к DLQ Delta-таблице
    pipeline_name: str        # имя пайплайна для трассировки
    pipeline_version: str     # версия кода для replay
    source_system: str        # источник данных
    processing_date: str      # дата обработки (партиция)
    run_id: Optional[str] = None  # ID запуска (Airflow run_id)


def build_validation_flags(df: DataFrame) -> DataFrame:
    """
    Шаг 1: добавляем колонки-флаги нарушений.
    Все операции - ленивые (Transformations), без Action.
    Главное: один проход данных позже.
    """
    return (
        df
        # Флаг: NULL в primary key
        .withColumn(
            "_err_null_order_id",
            when(col("order_id").isNull(), lit("NULL_ORDER_ID"))
            .otherwise(lit(None).cast("string"))
        )
        .withColumn(
            "_err_null_user_id",
            when(col("user_id").isNull(), lit("NULL_USER_ID"))
            .otherwise(lit(None).cast("string"))
        )
        # Флаг: некорректное значение amount
        .withColumn(
            "_err_invalid_amount",
            when(col("amount").isNull(), lit("NULL_AMOUNT"))
            .when(col("amount") <= 0, lit("NON_POSITIVE_AMOUNT"))
            .otherwise(lit(None).cast("string"))
        )
        # Флаг: невалидный статус
        .withColumn(
            "_err_invalid_status",
            when(
                ~col("status").isin("paid", "refunded", "pending", "cancelled"),
                lit("INVALID_STATUS")
            ).otherwise(lit(None).cast("string"))
        )
        # Флаг: дата из будущего
        .withColumn(
            "_err_future_date",
            when(
                col("created_at") > current_timestamp(),
                lit("FUTURE_CREATED_AT")
            ).otherwise(lit(None).cast("string"))
        )
        # Флаг: дата доставки раньше создания
        .withColumn(
            "_err_delivery_before_order",
            when(
                col("delivery_date").isNotNull() &
                (col("delivery_date") < col("created_at")),
                lit("DELIVERY_BEFORE_ORDER")
            ).otherwise(lit(None).cast("string"))
        )
    )


def build_error_array(df: DataFrame, error_cols: list) -> DataFrame:
    """
    Шаг 2: собираем все флаги в массив reasons.
    array_compact удаляет null из массива - остаются только реальные ошибки.
    """
    # Формируем array из всех флаг-колонок
    all_flags = array(*[col(c) for c in error_cols])

    return (
        df
        # Компактный массив: только non-null элементы
        .withColumn(
            "error_reasons",
            F.array_compact(all_flags)   # Spark 3.4+
            # Для более ранних версий:
            # F.filter(all_flags, lambda x: x.isNotNull())
        )
        # Булев флаг коррумпированности
        .withColumn(
            "is_corrupted",
            size(col("error_reasons")) > 0
        )
    )


def split_and_write(
    df: DataFrame,
    dlq_config: DLQConfig,
    silver_path: str,
    write_mode: str = "append",
) -> dict:
    """
    Шаг 3 + 4 + 5: кэшируем, разделяем, пишем.
    Возвращает статистику для мониторинга.
    """
    # Определяем колонки с флагами ошибок
    error_flag_cols = [
        "_err_null_order_id",
        "_err_null_user_id",
        "_err_invalid_amount",
        "_err_invalid_status",
        "_err_future_date",
        "_err_delivery_before_order",
    ]

    # Шаг 1+2: добавляем все флаги и собираем в массив
    flagged_df = build_error_array(
        build_validation_flags(df),
        error_flag_cols,
    )

    # КЛЮЧЕВОЙ МОМЕНТ: кэшируем после добавления флагов.
    # Данные материализуются ОДИН РАЗ.
    # Оба последующих filter идут из кэша, без повторного сканирования.
    flagged_df.cache()

    # Считаем статистику (один Action по кэшу)
    total_count = flagged_df.count()
    corrupted_count = flagged_df.filter(col("is_corrupted")).count()
    clean_count = total_count - corrupted_count

    print(
        f"[DLQ/{dlq_config.pipeline_name}] "
        f"Total: {total_count:,} | "
        f"Clean: {clean_count:,} | "
        f"DLQ: {corrupted_count:,} ({corrupted_count/total_count:.2%})"
    )

    # Колонки для записи в Silver (без служебных флагов)
    internal_cols = error_flag_cols + ["error_reasons", "is_corrupted"]

    # ── Запись чистых данных в Silver ──
    clean_df = (
        flagged_df
        .filter(~col("is_corrupted"))
        .drop(*internal_cols)
    )
    (
        clean_df
        .write
        .mode(write_mode)
        .format("delta")
        .partitionBy("processing_date")
        .save(silver_path)
    )

    # ── Запись плохих данных в DLQ ──
    if corrupted_count > 0:
        dlq_df = (
            flagged_df
            .filter(col("is_corrupted"))
            # Сохраняем raw_payload как JSON-строку исходной записи
            .withColumn(
                "raw_payload",
                to_json(struct(*[c for c in df.columns]))
            )
            # Добавляем метаданные DLQ
            .withColumn("pipeline_name", lit(dlq_config.pipeline_name))
            .withColumn("pipeline_version", lit(dlq_config.pipeline_version))
            .withColumn("source_system", lit(dlq_config.source_system))
            .withColumn("run_id", lit(dlq_config.run_id))
            .withColumn("quarantined_at", current_timestamp())
            .withColumn("processing_date", lit(dlq_config.processing_date))
            .withColumn(
                "error_category",
                # Определяем категорию по первой ошибке в массиве
                when(
                    F.array_contains(col("error_reasons"), "NULL_ORDER_ID") |
                    F.array_contains(col("error_reasons"), "NULL_USER_ID"),
                    lit("null_primary_key")
                ).when(
                    F.array_contains(col("error_reasons"), "NON_POSITIVE_AMOUNT") |
                    F.array_contains(col("error_reasons"), "NULL_AMOUNT"),
                    lit("amount_violation")
                ).when(
                    F.array_contains(col("error_reasons"), "INVALID_STATUS"),
                    lit("status_violation")
                ).otherwise(lit("date_violation"))
            )
            .withColumn("replay_attempted", lit("false"))
            # Выбираем только DLQ-схему (без бизнес-колонок)
            .select(
                "raw_payload",
                "error_reasons",
                "pipeline_name",
                "pipeline_version",
                "source_system",
                "run_id",
                "quarantined_at",
                "processing_date",
                "error_category",
                "replay_attempted",
            )
        )
        (
            dlq_df
            .write
            .mode("append")
            .format("delta")
            .partitionBy("processing_date", "error_category")
            .save(dlq_config.table_path)
        )
        print(f"[DLQ] Written {corrupted_count:,} records to quarantine")

    # Освобождаем кэш
    flagged_df.unpersist()

    return {
        "total": total_count,
        "clean": clean_count,
        "quarantined": corrupted_count,
        "quarantine_rate": corrupted_count / total_count if total_count > 0 else 0.0,
    }

6.3 Использование в пайплайне

def run_bronze_to_silver(spark: SparkSession, batch_date: str, run_id: str) -> None:
    """Полный пайплайн Bronze -> Silver с DLQ."""

    bronze_df = spark.read.format("delta").load(
        f"s3a://datalake/bronze/orders/date={batch_date}/"
    )

    dlq_config = DLQConfig(
        table_path="s3a://datalake/quarantine/orders/",
        pipeline_name="bronze_to_silver_orders",
        pipeline_version="2.4.1",
        source_system="ecommerce_backend",
        processing_date=batch_date,
        run_id=run_id,
    )

    stats = split_and_write(
        df=bronze_df,
        dlq_config=dlq_config,
        silver_path="s3a://datalake/silver/orders/",
    )

    # Оповещаем при высоком уровне карантина
    if stats["quarantine_rate"] > 0.05:  # более 5% - алерт
        raise ValueError(
            f"[ALERT] High quarantine rate: {stats['quarantine_rate']:.1%}. "
            f"DLQ: {stats['quarantined']:,}/{stats['total']:,}. "
            "Check upstream data quality."
        )

Часть 7. Обработка Parse Errors: битый JSON и CSV

7.1 Специфика parsing-ошибок

Parsing errors - особый класс: запись нельзя десериализовать вообще. Невозможно применить схему, нельзя вычислить бизнес-правила. Нужна отдельная стратегия.

def ingest_with_parse_dlq(
    spark: SparkSession,
    raw_path: str,
    target_schema: StructType,
    dlq_config: DLQConfig,
) -> DataFrame:
    """
    Чтение JSON с перехватом parsing errors.
    Битые строки идут в DLQ, валидные возвращаются для дальнейшей обработки.
    """
    # Добавляем специальную колонку для битых строк в схему
    schema_with_corrupt = target_schema.add(
        StructField("_corrupt_record", StringType(), nullable=True)
    )

    # Читаем в PERMISSIVE режиме - Spark не падает на битых строках
    raw_df = (
        spark.read
        .schema(schema_with_corrupt)
        .option("mode", "PERMISSIVE")
        .option("columnNameOfCorruptRecord", "_corrupt_record")
        .option("rescuedDataColumn", "_rescued_data")  # schema drift
        .json(raw_path)
    )

    raw_df.cache()

    # Три группы строк:
    # 1. Полностью битые (JSON синтаксически невалиден)
    parse_errors = raw_df.filter(col("_corrupt_record").isNotNull())

    # 2. Schema drift (новые/несовместимые поля)
    schema_drift = raw_df.filter(
        col("_corrupt_record").isNull() &
        col("_rescued_data").isNotNull()
    )

    # 3. Успешно распарсенные
    valid_parsed = raw_df.filter(
        col("_corrupt_record").isNull() &
        col("_rescued_data").isNull()
    )

    # Отправляем parse errors в DLQ
    if parse_errors.count() > 0:
        parse_dlq = (
            parse_errors
            .select(
                col("_corrupt_record").alias("raw_payload"),
                F.array(lit("PARSE_ERROR_MALFORMED_JSON")).alias("error_reasons"),
                lit("parse_error").alias("error_category"),
                lit(dlq_config.pipeline_name).alias("pipeline_name"),
                lit(dlq_config.pipeline_version).alias("pipeline_version"),
                lit(dlq_config.source_system).alias("source_system"),
                lit(dlq_config.run_id).alias("run_id"),
                current_timestamp().alias("quarantined_at"),
                lit(dlq_config.processing_date).alias("processing_date"),
                lit("false").alias("replay_attempted"),
            )
        )
        (
            parse_dlq
            .write.mode("append").format("delta")
            .partitionBy("processing_date", "error_category")
            .save(dlq_config.table_path)
        )

    # Schema drift - логируем как предупреждение, но пропускаем данные
    drift_count = schema_drift.count()
    if drift_count > 0:
        print(
            f"[DLQ] Schema drift: {drift_count} rows with unexpected fields. "
            f"Example: {schema_drift.select('_rescued_data').first()['_rescued_data']}"
        )

    raw_df.unpersist()

    # Возвращаем только валидно распарсенные строки (без служебных колонок)
    return valid_parsed.drop("_corrupt_record", "_rescued_data")

Часть 8. Replay Pipeline: жизненный цикл карантинных данных

8.1 Зачем нужен Replay

DLQ без процесса восстановления - это «чёрная дыра». Строки уходят туда и никогда не возвращаются. Со временем карантинная таблица вырастает в терабайты, никто не помнит, что там находится, и она начинает восприниматься как «место для мусора».

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

8.2 Реализация Replay Pipeline

from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F
from delta.tables import DeltaTable


def run_dlq_replay(
    spark: SparkSession,
    dlq_path: str,
    silver_path: str,
    processing_date: str,
    error_categories: list = None,
    dry_run: bool = False,
) -> dict:
    """
    Replay Pipeline: перечитываем DLQ, парсим raw_payload,
    применяем ОБНОВЛЁННУЮ логику валидации и дополняем Silver.

    dry_run=True - только считает, что можно восстановить, без записи.
    """

    # Читаем DLQ за нужную дату
    dlq_filter = F.col("processing_date") == processing_date
    if error_categories:
        dlq_filter = dlq_filter & F.col("error_category").isin(error_categories)

    dlq_df = (
        spark.read.format("delta").load(dlq_path)
        .filter(dlq_filter)
        .filter(F.col("replay_attempted") == "false")  # только необработанные
    )

    replay_count = dlq_df.count()
    print(f"[Replay] Found {replay_count:,} records for replay on {processing_date}")

    if replay_count == 0:
        return {"replayed": 0, "still_invalid": 0}

    # Парсим raw_payload обратно в DataFrame с текущей схемой
    # Используем from_json для безопасного парсинга JSON-строки
    parsed_df = (
        dlq_df
        .withColumn(
            "parsed",
            F.from_json(F.col("raw_payload"), BRONZE_ORDERS_SCHEMA)
        )
        .select("parsed.*", "raw_payload")  # разворачиваем struct
        .filter(F.col("order_id").isNotNull())  # базовая проверка
    )

    # Применяем НОВУЮ логику валидации (с исправленным кодом)
    flagged_df = build_error_array(
        build_validation_flags(parsed_df),
        error_flag_cols=ERROR_FLAG_COLS,
    )
    flagged_df.cache()

    # Что теперь стало валидным?
    now_valid_df = (
        flagged_df
        .filter(~F.col("is_corrupted"))
        .drop("error_reasons", "is_corrupted", "raw_payload",
              *ERROR_FLAG_COLS)
    )

    # Что по-прежнему невалидно?
    still_invalid_df = flagged_df.filter(F.col("is_corrupted"))

    valid_count = now_valid_df.count()
    still_invalid_count = still_invalid_df.count()

    print(
        f"[Replay] Now valid: {valid_count:,} | "
        f"Still invalid: {still_invalid_count:,}"
    )

    if not dry_run and valid_count > 0:
        # Дополняем Silver через MERGE для идемпотентности
        # (защита от дублирования при повторном replay)
        silver_table = DeltaTable.forPath(spark, silver_path)
        (
            silver_table.alias("silver")
            .merge(
                now_valid_df.alias("replayed"),
                "silver.order_id = replayed.order_id"
            )
            .whenNotMatchedInsertAll()  # вставляем только новые
            .execute()
        )
        print(f"[Replay] Merged {valid_count:,} records into Silver")

        # Помечаем успешно обработанные в DLQ
        dlq_table = DeltaTable.forPath(spark, dlq_path)
        (
            dlq_table.update(
                condition=f"processing_date = '{processing_date}' AND replay_attempted = 'false'",
                set={"replay_attempted": "'done'"}
            )
        )
        print(f"[Replay] Marked {valid_count:,} DLQ records as replayed")

    flagged_df.unpersist()

    return {
        "replayed": valid_count,
        "still_invalid": still_invalid_count,
    }

8.3 TTL и очистка DLQ

Карантинные данные нельзя хранить вечно - это дорого и создаёт путаницу. Нужна явная политика TTL:

def cleanup_dlq(
    spark: SparkSession,
    dlq_path: str,
    retention_days: int = 90,
) -> int:
    """
    Удаляем из DLQ записи старше retention_days дней.
    Работает только для записей, прошедших replay (replay_attempted = 'done').
    Необработанные записи не удаляем - сначала нужен replay.
    """
    from datetime import datetime, timedelta

    cutoff_date = (datetime.now() - timedelta(days=retention_days)).strftime("%Y-%m-%d")

    dlq_table = DeltaTable.forPath(spark, dlq_path)

    # Считаем, что будет удалено
    to_delete = (
        spark.read.format("delta").load(dlq_path)
        .filter(
            (F.col("processing_date") < cutoff_date) &
            (F.col("replay_attempted") == "done")
        )
        .count()
    )

    if to_delete > 0:
        dlq_table.delete(
            condition=(
                f"processing_date < '{cutoff_date}' "
                f"AND replay_attempted = 'done'"
            )
        )
        # Физическая очистка файлов (Delta VACUUM)
        dlq_table.vacuum(retentionHours=0)  # immediate cleanup
        print(f"[DLQ Cleanup] Deleted {to_delete:,} records older than {cutoff_date}")

    return to_delete

Часть 9. Мониторинг и Observability DLQ

9.1 Quarantine Rate как ключевая метрика

Самая важная метрика DLQ-системы - quarantine rate: доля карантинных записей от общего объёма батча.

from dataclasses import dataclass

@dataclass
class DLQMetrics:
    """Метрики одного прогона DLQ для записи в monitoring-таблицу."""
    pipeline_name: str
    processing_date: str
    run_id: str
    total_records: int
    clean_records: int
    quarantined_records: int
    quarantine_rate: float
    quarantine_by_category: dict   # {"null_primary_key": 100, "amount_violation": 50}
    run_duration_seconds: float
    recorded_at: str


def write_dlq_metrics(
    spark: SparkSession,
    metrics: DLQMetrics,
    metrics_table: str = "metadata.dlq_metrics",
) -> None:
    """Записываем метрики DLQ в мониторинг-таблицу."""
    import json
    from dataclasses import asdict

    row = asdict(metrics)
    row["quarantine_by_category"] = json.dumps(row["quarantine_by_category"])
    metrics_df = spark.createDataFrame([row])
    (
        metrics_df
        .write
        .mode("append")
        .format("delta")
        .partitionBy("processing_date", "pipeline_name")
        .saveAsTable(metrics_table)
    )


def compute_quarantine_by_category(
    spark: SparkSession,
    dlq_path: str,
    processing_date: str,
) -> dict:
    """Группируем DLQ-записи по категории ошибки для детального мониторинга."""
    rows = (
        spark.read.format("delta").load(dlq_path)
        .filter(F.col("processing_date") == processing_date)
        .groupBy("error_category")
        .count()
        .collect()
    )
    return {row["error_category"]: row["count"] for row in rows}

9.2 Пороги алертинга

def check_dlq_thresholds(
    metrics: DLQMetrics,
    critical_rate: float = 0.05,
    warning_rate: float = 0.01,
) -> str:
    """
    Проверяем quarantine rate против пороговых значений.
    Возвращает уровень алерта: "OK", "WARNING", "CRITICAL".
    """
    rate = metrics.quarantine_rate

    if rate >= critical_rate:
        return (
            f"CRITICAL: quarantine_rate={rate:.1%} "
            f"({metrics.quarantined_records:,}/{metrics.total_records:,} records). "
            f"Pipeline: {metrics.pipeline_name}, date: {metrics.processing_date}. "
            f"Breakdown: {metrics.quarantine_by_category}"
        )

    if rate >= warning_rate:
        return (
            f"WARNING: quarantine_rate={rate:.1%}. "
            f"Monitor upstream data quality."
        )

    return "OK"


def push_to_airflow_xcom(run_id: str, value: dict) -> None:
    """
    В реальном production используем XCom или внешнюю систему алертов.
    Здесь - заглушка с print для демонстрации.
    """
    import json
    print(f"[XCom/{run_id}] {json.dumps(value)}")

9.3 Трендовый анализ quarantine rate

Единичный высокий quarantine rate может быть нормой (праздники, маркетинговые акции). Опасен растущий тренд - это ранний сигнал деградации источника:

def detect_dlq_trend(
    spark: SparkSession,
    pipeline_name: str,
    metrics_table: str = "metadata.dlq_metrics",
    lookback_days: int = 7,
) -> dict:
    """
    Сравниваем сегодняшний quarantine rate со средним за последние N дней.
    Рост на 50%+ - критический алерт.
    """
    rows = (
        spark.read.format("delta").load(f"spark_catalog.{metrics_table}")
        .filter(F.col("pipeline_name") == pipeline_name)
        .orderBy(F.col("processing_date").desc())
        .limit(lookback_days + 1)
        .select("processing_date", "quarantine_rate")
        .collect()
    )

    if len(rows) < 2:
        return {"trend": "insufficient_data"}

    today_rate = rows[0]["quarantine_rate"]
    historical_avg = sum(r["quarantine_rate"] for r in rows[1:]) / (len(rows) - 1)

    ratio = today_rate / historical_avg if historical_avg > 0 else float("inf")

    return {
        "today_rate": today_rate,
        "historical_avg": historical_avg,
        "ratio": ratio,
        "trend": "SPIKE" if ratio > 1.5 else "NORMAL",
    }

Часть 10. Anti-Patterns в DLQ-архитектуре

10.1 Silent Drop: молчаливое удаление

Самый опасный антипаттерн - молча удалять плохие строки без записи в DLQ:

# АНТИПАТТЕРН: данные теряются без следа
df = df.filter(col("amount") > 0)  # отрицательные суммы исчезают
df = df.dropna(subset=["order_id"])  # строки без order_id исчезают

# Нет DLQ, нет метрики, нет алерта, нет возможности аудита.
# Аналитик через месяц: "почему у нас пропали 50K транзакций в январе?"
# Ответа нет.

# ПРАВИЛЬНО: всегда явный DLQ с метаданными
df, dlq_df, stats = split_to_dlq(df, rules, dlq_config)

10.2 DLQ без Replay: чёрная дыра

DLQ, в который данные только добавляются, но никогда не читаются - бессмысленная трата места:

# АНТИПАТТЕРН: DLQ как конечная точка
dlq_df.write.mode("append").parquet(dlq_path)
# И всё. Replay-пайплайна нет. Метрик нет. Алертов нет.
# Через полгода - 5TB карантинных данных, которые никто не знает как обработать.

# ПРАВИЛЬНО: DLQ - это начало recovery-процесса
# 1. Запись в DLQ с метаданными
# 2. Alert при превышении порога
# 3. Запланированный Replay DAG
# 4. TTL и очистка

10.3 Одинаковая схема для DLQ и Silver

# АНТИПАТТЕРН: DLQ использует Silver-схему
bad_records.write.format("delta").save(silver_path + "_dlq")
# Проблема: если запись плохая из-за amount = "N/A" (строка вместо числа),
# то при записи в Silver-схему (amount: DoubleType) - cast упадёт снова.
# DLQ сам стал источником ошибок.

# ПРАВИЛЬНО: DLQ хранит raw_payload как StringType
# Схема DLQ всегда фиксированная и принимает любые данные

10.4 Excessive Tolerance: карантин вместо остановки

# АНТИПАТТЕРН: всё идёт в карантин, пайплайн "успешно" завершается
stats = split_and_write(df, dlq_config, silver_path)
# quarantine_rate = 0.99 (99% записей в карантине)
# Silver практически пустой, Gold показывает нули
# Но пайплайн "зелёный" - он же завершился без ошибок!

# ПРАВИЛЬНО: контролируем максимально допустимый quarantine rate
if stats["quarantine_rate"] > 0.20:  # 20% - Fail-Fast
    raise ValueError(
        f"Quarantine rate {stats['quarantine_rate']:.1%} exceeds 20%. "
        "Data is massively corrupted. Stopping pipeline."
    )

10.5 Quarantine без ownership

Карантин - это не «место для мусора, который разгребут потом». У каждой DLQ-таблицы должен быть явный владелец с обязательствами:

# metadata/dlq_ownership.yaml
tables:
  - path: s3a://datalake/quarantine/orders/
    owner_team: data-engineering
    on_call: @data-eng-oncall
    sla_replay_hours: 24       # replay в течение 24 часов после алерта
    max_retention_days: 90
    alert_threshold_rate: 0.05
    alert_channels:
      - slack: "#data-incidents"
      - pagerduty: "data-platform-critical"

Часть 11. Production Case: полный DLQ-пайплайн

11.1 Постановка задачи

Реальный сценарий: ingestion customer events из S3 (JSON от мобильного приложения). События содержат информацию о заказах. Источник - внешняя команда, контракт данных слабый, schema drift - норма жизни.

Требования:

  • Битые JSON-строки не должны останавливать ingestion
  • Все нарушения фиксируются с полной трассировкой
  • При quarantine rate > 5% - критический алерт
  • Раз в сутки - Replay DAG перебирает DLQ

11.2 Конфигурация пайплайна

# config/dlq_pipeline_config.py

from dataclasses import dataclass
from pyspark.sql.types import (
    StructType, StructField, StringType, LongType,
    DoubleType, TimestampType, IntegerType
)

CUSTOMER_EVENTS_SCHEMA = StructType([
    StructField("event_id", StringType(), nullable=False),
    StructField("order_id", StringType(), nullable=False),
    StructField("user_id", LongType(), nullable=False),
    StructField("event_type", StringType(), nullable=True),
    StructField("amount", DoubleType(), nullable=True),
    StructField("currency", StringType(), nullable=True),
    StructField("status", StringType(), nullable=True),
    StructField("created_at", TimestampType(), nullable=False),
    StructField("device_type", StringType(), nullable=True),
    StructField("session_id", StringType(), nullable=True),
])

CUSTOMER_EVENTS_RULES = [
    # Hard: критические поля PK
    {"name": "NULL_EVENT_ID",   "expr": "event_id IS NOT NULL",  "severity": "critical"},
    {"name": "NULL_ORDER_ID",   "expr": "order_id IS NOT NULL",  "severity": "critical"},
    {"name": "NULL_USER_ID",    "expr": "user_id IS NOT NULL",   "severity": "critical"},
    {"name": "NULL_CREATED_AT", "expr": "created_at IS NOT NULL","severity": "critical"},
    # Soft: бизнес-правила
    {"name": "INVALID_EVENT_TYPE",
     "expr": "event_type IN ('order_placed', 'order_updated', 'order_cancelled', 'payment_confirmed')",
     "severity": "warning"},
    {"name": "NEGATIVE_AMOUNT",
     "expr": "amount IS NULL OR amount >= 0",
     "severity": "warning"},
    {"name": "FUTURE_EVENT",
     "expr": "created_at <= current_timestamp()",
     "severity": "warning"},
]

DLQ_PATHS = {
    "parse_errors": "s3a://datalake/quarantine/customer_events/parse_errors/",
    "business_violations": "s3a://datalake/quarantine/customer_events/business/",
    "silver": "s3a://datalake/silver/customer_events/",
}

11.3 Полный pipeline

# jobs/customer_events_ingestion.py

import sys
import time
from datetime import datetime
from pyspark.sql import SparkSession
from pyspark.sql import functions as F


def main():
    spark = (
        SparkSession.builder
        .appName("CustomerEventsIngestion")
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .getOrCreate()
    )

    batch_date = sys.argv[1] if len(sys.argv) > 1 else str(datetime.today().date())
    run_id = sys.argv[2] if len(sys.argv) > 2 else f"local_{batch_date}"
    raw_path = f"s3a://datalake/bronze/customer_events/date={batch_date}/"

    print(f"=== Customer Events Ingestion: {batch_date} | run_id={run_id} ===")
    start = time.time()

    dlq_config = DLQConfig(
        table_path=DLQ_PATHS["business_violations"],
        pipeline_name="customer_events_ingestion",
        pipeline_version="3.1.0",
        source_system="mobile_backend",
        processing_date=batch_date,
        run_id=run_id,
    )

    # ─── Шаг 1: Чтение с перехватом parse errors ───
    print("\n[Step 1] Reading Bronze with parse error handling...")
    parse_dlq_config = DLQConfig(
        table_path=DLQ_PATHS["parse_errors"],
        pipeline_name="customer_events_ingestion",
        pipeline_version="3.1.0",
        source_system="mobile_backend",
        processing_date=batch_date,
        run_id=run_id,
    )
    valid_parsed_df = ingest_with_parse_dlq(
        spark, raw_path, CUSTOMER_EVENTS_SCHEMA, parse_dlq_config
    )

    # ─── Шаг 2: Дедупликация (до бизнес-валидации) ───
    print("\n[Step 2] Deduplication...")
    before_dedup = valid_parsed_df.count()
    valid_parsed_df = valid_parsed_df.dropDuplicates(["event_id"])
    after_dedup = valid_parsed_df.count()
    print(f"  Dedup: {before_dedup - after_dedup:,} duplicates removed")

    # ─── Шаг 3: Бизнес-валидация + Single-Pass Splitter ───
    print("\n[Step 3] Business validation + DLQ split...")
    stats = split_and_write(
        df=valid_parsed_df,
        dlq_config=dlq_config,
        silver_path=DLQ_PATHS["silver"],
    )

    # ─── Шаг 4: Запись метрик мониторинга ───
    duration = time.time() - start
    by_category = compute_quarantine_by_category(
        spark, DLQ_PATHS["business_violations"], batch_date
    )

    metrics = DLQMetrics(
        pipeline_name="customer_events_ingestion",
        processing_date=batch_date,
        run_id=run_id,
        total_records=stats["total"],
        clean_records=stats["clean"],
        quarantined_records=stats["quarantined"],
        quarantine_rate=stats["quarantine_rate"],
        quarantine_by_category=by_category,
        run_duration_seconds=duration,
        recorded_at=datetime.utcnow().isoformat(),
    )
    write_dlq_metrics(spark, metrics)

    # ─── Шаг 5: Проверка порогов и алертинг ───
    alert_level = check_dlq_thresholds(metrics, critical_rate=0.05, warning_rate=0.01)
    print(f"\n[Alert] {alert_level}")

    if alert_level.startswith("CRITICAL"):
        raise RuntimeError(f"[CRITICAL DLQ] {alert_level}")

    # ─── Шаг 6: Трендовый анализ ───
    trend = detect_dlq_trend(spark, "customer_events_ingestion")
    if trend.get("trend") == "SPIKE":
        print(
            f"[TREND ALERT] DLQ spike: today={trend['today_rate']:.1%}, "
            f"7d_avg={trend['historical_avg']:.1%}, ratio={trend['ratio']:.1f}x"
        )

    print(f"""
=== Pipeline COMPLETE ===
  Duration:       {duration:.1f}s
  Total events:   {stats['total']:,}
  Clean:          {stats['clean']:,}
  DLQ:            {stats['quarantined']:,} ({stats['quarantine_rate']:.2%})
  DLQ breakdown:  {by_category}
""")


if __name__ == "__main__":
    main()

11.4 Airflow DAG с Replay

# dags/customer_events_pipeline.py

from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
from datetime import timedelta

with DAG(
    dag_id="customer_events_pipeline",
    schedule_interval="0 4 * * *",
    start_date=days_ago(1),
    catchup=False,
    default_args={
        "retries": 2,
        "retry_delay": timedelta(minutes=10),
    },
    tags=["customer_events", "dlq"],
) as dag:

    # Основной ingestion с DLQ
    ingest = SparkSubmitOperator(
        task_id="ingest_customer_events",
        application="jobs/customer_events_ingestion.py",
        application_args=["{{ ds }}", "{{ run_id }}"],
        conf={"spark.executor.memory": "8g"},
        pool="spark_heavy_pool",
    )

    # Ежедневный Replay DLQ предыдущего дня
    replay = SparkSubmitOperator(
        task_id="replay_dlq",
        application="jobs/dlq_replay.py",
        # Replay за вчера - к этому моменту уже известно, что исправлено
        application_args=["{{ macros.ds_add(ds, -1) }}", "{{ run_id }}"],
        conf={"spark.executor.memory": "4g"},
        pool="spark_light_pool",
    )

    # Еженедельная очистка устаревшего карантина (по пятницам)
    cleanup = SparkSubmitOperator(
        task_id="cleanup_old_dlq",
        application="jobs/dlq_cleanup.py",
        application_args=["90"],  # TTL = 90 дней
        pool="spark_light_pool",
        trigger_rule="none_failed_min_one_success",
    )

    ingest >> replay >> cleanup

Итоги

Dead Letter Queue - это не просто «корзина для мусора». Это полноценный архитектурный компонент, определяющий надёжность data platform.

Single-Pass Splitter - фундаментальный принцип. Добавляем флаги нарушений как lazy Transformations, кэшируем DataFrame один раз, затем два filter из кэша. Это вдвое быстрее наивного подхода с двумя полными сканами данных.

Raw payload всегда - карантинная таблица должна хранить исходную строку в нетронутом виде (raw_payload: StringType). Без этого replay невозможен: нельзя восстановить то, что уже интерпретировано через битую логику.

DLQ без Replay - антипаттерн. Карантин имеет смысл только вместе с явным процессом восстановления: Replay DAG, политика TTL и ownership команды.

Quarantine Rate - ведущая метрика здоровья источника. Медленный рост quarantine rate - ранний сигнал деградации upstream-системы, который видно за дни до того, как BI-аналитики начнут получать неверные цифры.

Схема DLQ ≠ схема Silver. DLQ принимает данные любой формы через StringType raw_payload. Silver принимает только корректные данные через строгую бизнес-схему. Смешение этих двух схем ведёт к каскадным ошибкам.

Пороги quarantine rate: до 1% - норма, 1–5% - warning и расследование, больше 5% - critical, больше 20% - Fail-Fast, не пытаемся «выжить» на 80% данных.