Dead Letter Queue: паттерн Quarantine для плохих записей
Изоляция некорректных записей без остановки пайплайна: архитектура DLQ в Medallion, Single-Pass Splitter, схема карантинной таблицы, rescued_data, try_cast, replay-пайплайны и мониторинг quarantine-rate
Введение: плохие данные неизбежны - вопрос только в том, как с ними работать¶
Представьте: пайплайн 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, остальные поля nullDROPMALFORMED: битые строки молча удаляются - антипаттерн (данные теряются)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% данных.