Контроль Ingestion: inferSchema, режимы чтения и изоляция брака

DataFrameReader internals, inferSchema vs explicit StructType, PERMISSIVE/DROPMALFORMED/FAILFAST, _corrupt_record, badRecordsPath - надёжный ingestion gate для production pipelines

core optimization

Почему ingestion - самое опасное место в data pipeline

Каждый data pipeline начинается с чтения данных. Именно здесь принимаются решения, которые определяют качество всего последующего потока: какой схеме доверять, что делать с битыми записями, как изолировать брак от основного потока. Ошибки на этапе ingestion - самые дорогостоящие, потому что они скрытые.

Представьте реальный сценарий: внешний поставщик ежедневно присылает CSV-файл с транзакциями. В понедельник в файле появился новый тип значения - строка "N/A" вместо числа в поле amount. Spark в режиме по умолчанию молча заменит эти значения на NULL. Аналитик в BI-системе увидит просадку суммы продаж на 15% и поднимет тревогу - но это не реальное падение продаж, это тихая потеря данных на уровне чтения.

Такие ситуации происходят регулярно. По данным опросов Data Engineering Community, более 70% критических инцидентов в data pipelines связаны с изменением схемы или формата данных источника. И большинство из них обнаруживаются не сразу - а через несколько дней или недель.

В этом уроке мы разберём механизмы, которые позволяют контролировать ingestion:

  • явные схемы как фиксированный контракт с источником
  • режимы обработки ошибок PERMISSIVE, DROPMALFORMED, FAILFAST
  • изоляция брака через badRecordsPath и паттерн Dead Letter Queue
  • опции парсинга для превентивной очистки данных

DataFrameReader: API чтения данных в Spark

Точкой входа для любого чтения данных в Spark является DataFrameReader - объект, доступный через spark.read. Он работает по паттерну Builder: методы цепочки конфигурируют читатель, а финальный вызов (csv(), json(), parquet()) инициирует ленивое построение Logical Plan.

# Базовая структура DataFrameReader
df = (
    spark.read
    .format("csv")           # формат: csv, json, parquet, orc, avro, delta, iceberg
    .option("key", "value")  # опции, специфичные для формата
    .schema(schema)          # явная схема (опционально)
    .load("s3a://bucket/path/")  # путь к данным (файл, папка, глоб-паттерн)
)

DataFrameReader - это lazy API. Вызов .load() не читает данные немедленно. Spark только строит Logical Plan (описание намерений). Реальное IO начинается только при вызове action: count(), show(), write(). Это принципиально отличает Spark от pandas, где pd.read_csv() немедленно загружает данные.

Важно также понимать, что DataFrameReader поддерживает два режима задания схемы: автоматический вывод (inferSchema) и явное объявление (StructType). Выбор между ними - фундаментальное архитектурное решение, которое определяет надёжность и производительность всего pipeline.

inferSchema: автоматический вывод типов

Опция inferSchema=true привлекает простотой: не нужно думать о схеме, Spark сам разберётся. Но за этим удобством скрыта высокая цена в production.

Как работает inferSchema под капотом

Когда Spark видит inferSchema=true, он выполняет полное предварительное сканирование (full scan) данных перед основным чтением. Этот процесс состоит из нескольких шагов.

Шаг 1: Sampling или Full Scan

Для CSV по умолчанию Spark читает все строки файла (параметр samplingRatio равен 1.0). Каждый executor читает свою партицию и для каждой колонки пробует подобрать тип данных. Для JSON поведение аналогично.

Шаг 2: Type Promotion - иерархия типов

Spark пробует разобрать каждое значение по иерархии типов, от самого строгого к самому широкому:

LongType → DoubleType → BooleanType → TimestampType → DateType → StringType

Алгоритм жадный: как только одно значение в колонке не помещается в тип Long, вся колонка «повышается» до Double. Если встречается строка, которую нельзя разобрать как число или дату - колонка становится StringType. Решение принимается по всему корпусу прочитанных строк в партиции.

Шаг 3: Schema Merge на Driver

После того как каждый executor прочитал свою партицию и вывел локальную схему, результаты отправляются на Driver. Driver применяет правила слияния: для каждой колонки берётся «наиболее широкий» тип из всех партиций. Например, если в партиции 1 колонка age имеет тип Long, а в партиции 2 - String (потому что там оказалась строка "unknown"), итоговый тип будет String.

Итог: inferSchema=true - это два полных прохода по данным. После вывода схемы Spark выполняет второй полный проход для фактического чтения.

На диаграмме видно ключевое различие: при inferSchema данные сканируются дважды. Для файла размером 50 ГБ это означает 100 ГБ прочитанных данных вместо 50. При чтении из S3 с оплатой за трафик - это буквально двойная стоимость каждого запуска pipeline.

Проблемы inferSchema в production

Помимо двойного сканирования, inferSchema=true имеет ряд системных проблем.

Нестабильность схемы. Если данные приходят партиями (ежедневные файлы), и в первые несколько дней колонка user_id содержала только числа, Spark выведет тип Long. Но в пятницу поставщик добавил буквенные идентификаторы - и схема стала String. Pipeline, который работал неделю, сломается при попытке сравнить результаты с таблицей, где user_id является Long.

Non-reproducibility. Схема, выведенная из выборки samplingRatio=0.1, может отличаться от схемы полного датасета, если «редкие» типы встречаются только в последних 10% данных. Запуск на разных наборах данных даёт разные схемы - это нарушает воспроизводимость pipeline.

Запрет в Streaming. В Spark Structured Streaming inferSchema вообще не работает - Spark не может делать полный scan бесконечного потока данных. Явная схема обязательна для любого streaming pipeline.

Type promotion surprises. Если в CSV-файле из 100 миллионов строк одна строка содержит значение "true" в числовой колонке, вся колонка станет Boolean. Инженер, который не проверил каждую из 100M строк, получит неожиданный результат.

Производительность на реальных данных:

Метод Время чтения Количество сканирований Примечание
inferSchema=true ~94 сек 2 Полный прогон дважды
samplingRatio=0.1 ~58 сек ~1.1 Риск неверной схемы на «хвостовых» данных
Explicit Schema ~41 сек 1 Прямое чтение, типы известны заранее

Тест на CSV 8 ГБ, 45M строк, кластер 8 executors × 4 CPU. Ускорение в 2.3× только за счёт смены подхода к схеме - без изменения логики обработки.

Explicit Schema: явный контракт

Производственный код всегда должен объявлять схему явно через StructType. Это не просто оптимизация - это фиксация контракта с источником данных: «я ожидаю данные в такой структуре, и если она изменится - я хочу об этом знать немедленно».

Анатомия StructType

StructType - это упорядоченный список StructField. Каждое поле имеет три обязательных параметра: имя, тип и признак nullable.

from pyspark.sql.types import (
    StructType, StructField,
    StringType, IntegerType, LongType,
    DoubleType, BooleanType, TimestampType, DateType,
    ArrayType, MapType
)

# Явная схема для файла транзакций
TRANSACTION_SCHEMA = StructType([
    StructField("transaction_id", StringType(),    nullable=False),
    StructField("user_id",        LongType(),      nullable=False),
    StructField("amount",         DoubleType(),    nullable=True),
    StructField("currency",       StringType(),    nullable=True),
    StructField("event_ts",       TimestampType(), nullable=True),
    StructField("is_fraud",       BooleanType(),   nullable=True),
    StructField("metadata",       MapType(StringType(), StringType()), nullable=True),
])

Разберём каждый параметр StructField подробнее:

  • name - имя колонки. Для JSON чувствительно к регистру (JSON-ключи чувствительны к регистру). Для CSV по умолчанию не чувствительно. Должно точно совпадать с именем в заголовке файла.

  • dataType - ожидаемый тип данных. Spark будет принудительно приводить значения из источника к этому типу при чтении. Если приведение невозможно - поведение зависит от режима чтения (PERMISSIVE, DROPMALFORMED, FAILFAST).

  • nullable - разрешены ли NULL значения. Важно понять: это не SQL-constraint и не вызовет ошибку само по себе при появлении NULL. Это метаданные для Catalyst Optimizer, которые позволяют ему генерировать более эффективный код (например, пропускать проверки на NULL там, где это гарантировано не нужно). Для downstream-валидации можно добавить явные проверки.

Вложенные структуры: StructType рекурсивен

Для JSON с вложенными объектами StructType применяется рекурсивно. Допустим, входной JSON выглядит так:

{
  "order_id": "ORD-12345",
  "total": 299.99,
  "customer": {
    "id": 42,
    "name": "Alice",
    "address": {"city": "Moscow", "country": "RU"}
  },
  "items": [
    {"sku": "SKU-001", "quantity": 2, "price": 149.99}
  ]
}

Соответствующая схема:

ORDER_SCHEMA = StructType([
    StructField("order_id", StringType(), nullable=False),
    StructField("total",    DoubleType(), nullable=True),
    StructField("customer", StructType([
        StructField("id",   LongType(),   nullable=False),
        StructField("name", StringType(), nullable=True),
        StructField("address", StructType([
            StructField("city",    StringType(), nullable=True),
            StructField("country", StringType(), nullable=True),
        ]), nullable=True),
    ]), nullable=True),
    StructField("items", ArrayType(StructType([
        StructField("sku",      StringType(),  nullable=False),
        StructField("quantity", IntegerType(), nullable=True),
        StructField("price",    DoubleType(),  nullable=True),
    ])), nullable=True),
])

При чтении с этой схемой Spark автоматически десериализует вложенные объекты и массивы в соответствующие StructType и ArrayType колонки. К полям вложенных структур обращаются через точечную нотацию: df["customer.address.city"] или через col("customer.address.city").

Схема как DDL-строка

Помимо программного определения через StructType, Spark поддерживает компактную строковую форму DDL:

# Эквивалент StructType через строку DDL
schema = spark.read.schema(
    "transaction_id STRING NOT NULL, user_id LONG NOT NULL, "
    "amount DOUBLE, currency STRING, event_ts TIMESTAMP"
)

# Или через метод _parse_datatype_string
from pyspark.sql.types import _parse_datatype_string
schema = _parse_datatype_string(
    "STRUCT<transaction_id: STRING, user_id: BIGINT, amount: DOUBLE>"
)

DDL-форма удобна для хранения схемы во внешнем конфигурационном файле или в Hive Metastore. Но для сложных вложенных структур StructType в коде читается лучше.

Что происходит при несовпадении типов

Когда Spark читает значение, которое не соответствует объявленному типу в схеме, поведение определяется режимом чтения (Read Mode). Рассмотрим конкретный пример: в CSV-файле колонка amount объявлена как DoubleType, но одна из строк содержит значение "N/A":

transaction_id,user_id,amount
TXN-001,100,299.99
TXN-002,101,N/A         ← проблемная строка
TXN-003,102,150.00

Что произойдёт с этой строкой - зависит исключительно от режима чтения. Spark поддерживает три режима для текстовых форматов (CSV, JSON): PERMISSIVE, DROPMALFORMED и FAILFAST.

Режимы обработки ошибок

На диаграмме видны три принципиально разных поведения. PERMISSIVE сохраняет все строки, но маркирует проблемные. DROPMALFORMED молча удаляет проблемные строки. FAILFAST останавливает весь job при первой ошибке. Выбор режима - это выбор стратегии обработки ошибок данных.

PERMISSIVE - мягкий режим (по умолчанию)

PERMISSIVE - режим по умолчанию. Его философия: «принять как можно больше данных, но не потерять информацию о проблемах».

Поведение PERMISSIVE зависит от типа ошибки. При ошибке типа поля (несовпадение типа, например "N/A" вместо числа): Spark подставляет NULL в это поле. Остальные поля строки разбираются нормально. Колонка _corrupt_record при этом не заполняется - строка просто получает NULL на месте проблемного поля.

При структурной ошибке (битый JSON, неправильное число колонок в CSV, незакрытая кавычка): Spark сохраняет исходную строку целиком в специальную колонку _corrupt_record. Все остальные поля этой строки заполняются NULL.

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

# _corrupt_record ОБЯЗАТЕЛЬНО должен быть в схеме,
# иначе Spark не будет заполнять его при ошибках
schema = StructType([
    StructField("transaction_id",  StringType(),  nullable=True),
    StructField("user_id",         LongType(),    nullable=True),
    StructField("amount",          DoubleType(),  nullable=True),
    StructField("_corrupt_record", StringType(),  nullable=True),
])

df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("mode", "PERMISSIVE")   # значение по умолчанию, можно не указывать
    .schema(schema)
    .load("s3a://data-bucket/transactions/")
)

# Фильтруем битые строки - те, где _corrupt_record заполнен
corrupt_df = df.filter(df["_corrupt_record"].isNotNull())
clean_df   = df.filter(df["_corrupt_record"].isNull())

total_count   = df.cache().count()
corrupt_count = corrupt_df.count()
clean_count   = total_count - corrupt_count

print(f"Всего строк:   {total_count:,}")
print(f"Корректных:    {clean_count:,}")
print(f"Структурных ошибок: {corrupt_count:,}")

# Смотрим на содержимое битых строк
corrupt_df.select("_corrupt_record").show(5, truncate=False)

Есть критически важный нюанс: _corrupt_record заполняется только при структурных ошибках (битый JSON, неправильное число полей в CSV). Ошибки типа поля ("N/A" вместо числа) приводят к NULL в поле, но _corrupt_record при этом не заполняется. Это нужно учитывать при построении мониторинга качества данных.

Также можно переименовать техническую колонку:

# Опция columnNameOfCorruptRecord позволяет изменить имя технической колонки
df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("columnNameOfCorruptRecord", "_raw_error")  # новое имя
    .schema(schema_with_raw_error_field)  # схема должна содержать поле _raw_error
    .load(path)
)

Когда использовать PERMISSIVE:

PERMISSIVE - правильный выбор для Bronze-слоя в архитектуре Lakehouse. Принцип Bronze: «сохранить всё, что пришло, не терять ни байта информации». Битые строки остаются в датасете и могут быть проанализированы и повторно обработаны после исправления источника.

DROPMALFORMED - режим тихого удаления

DROPMALFORMED - режим, при котором Spark молча удаляет любую строку, которую не смог разобрать в соответствии со схемой.

df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("mode", "DROPMALFORMED")
    .schema(schema)
    .load("s3a://data-bucket/transactions/")
)
# Битые строки исчезли без следа - DataFrame содержит только корректные записи

Режим удобен на первый взгляд: чистый DataFrame без NULL в неожиданных местах, не нужно думать об обработке ошибок. Но он несёт серьёзный риск - silent data loss (тихая потеря данных).

Почему DROPMALFORMED опасен:

Представьте: поставщик изменил формат поля amount, и теперь 30% строк не парсятся. В DROPMALFORMED-режиме вы получите DataFrame, в котором тихо отсутствует треть данных. Аналитик посмотрит на таблицу - данные есть, всё выглядит корректно. Никакого предупреждения, никакого исключения. Только через несколько дней при сравнении с данными поставщика выяснится, что потеряно 30% транзакций.

В Spark UI при DROPMALFORMED количество удалённых строк не отображается по умолчанию. Чтобы получить метрику потерь, необходим внешний контроль:

# Паттерн: DROPMALFORMED с обязательным контролем потерь
raw_line_count  = spark.read.text(path).count() - 1  # -1 за строку заголовка в CSV
df_clean = spark.read.option("mode", "DROPMALFORMED").schema(schema).csv(path)
clean_count     = df_clean.count()

drop_rate = (raw_line_count - clean_count) / raw_line_count
print(f"Удалено строк: {raw_line_count - clean_count:,} ({drop_rate:.2%})")

if drop_rate > 0.01:  # более 1% потерь - это проблема
    raise ValueError(
        f"DROPMALFORMED удалил {drop_rate:.1%} строк - возможно изменение схемы источника!"
    )

Этот паттерн превращает DROPMALFORMED из «молчаливого удалятеля» в управляемый инструмент. Но он требует двойного чтения данных (один раз как текст для подсчёта, второй - для парсинга), что удваивает стоимость IO.

Когда DROPMALFORMED допустим:

В ограниченных сценариях DROPMALFORMED имеет смысл:

  • небольшой процент брака (менее 0.01%) и брак ожидаем по природе источника
  • есть внешняя система reconciliation (сравнение итогов с поставщиком)
  • данные некритичны: например, аналитика кликов, где потеря 0.01% событий незначительна
  • есть мониторинг записей «до» и «после» чтения

FAILFAST - режим строгого контроля

FAILFAST - режим нулевой терпимости. При обнаружении первой же некорректной записи Spark немедленно останавливает выполнение и выбрасывает исключение.

df = (
    spark.read
    .format("json")
    .option("mode", "FAILFAST")
    .schema(schema)
    .load("s3a://data-bucket/payments/")
)
# Если в данных есть хоть одна битая строка - здесь будет SparkException

При ошибке Spark выбрасывает:

org.apache.spark.SparkException: Malformed records are detected in record parsing.
Parse Mode: FAILFAST. To process malformed records as null result, try setting
the option 'mode' as 'PERMISSIVE'.

Исключение содержит информацию о том, какая запись вызвала ошибку и в каком файле. Это делает диагностику быстрой: нет нужды искать битые строки по всему датасету.

Важная особенность в распределённом контексте. Spark выполняется на нескольких executor-ах параллельно. Когда один executor обнаруживает ошибку, он завершает свою задачу с исключением. Driver помечает Stage как FAILED. Но до этого момента другие executor-ы могут успеть записать часть данных в sink. Для FAILFAST-сценариев необходимо использовать транзакционные форматы (Iceberg, Delta Lake) с поддержкой atomic writes.

# Паттерн: FAILFAST с retry-логикой и алертом
try:
    df = (
        spark.read
        .format("json")
        .option("mode", "FAILFAST")
        .schema(schema)
        .load(daily_path)
    )
    df.write.format("iceberg").mode("append").saveAsTable("bronze.payments")

except Exception as e:
    # Отправить алерт инженеру данных
    send_alert(f"FAILFAST: Ingestion failed for {daily_path}\nError: {e}")
    # Сохранить сырые данные в карантин для ручного анализа
    spark.read.text(daily_path).write.mode("overwrite").text(f"s3a://quarantine/{today}/")
    raise  # пробросить исключение, чтобы оркестратор пометил job как FAILED

Когда использовать FAILFAST:

FAILFAST - правильный выбор в следующих сценариях:

  • финансовые транзакции и платёжные данные - любая потеря или искажение данных неприемлемы; лучше упасть и разобраться, чем молча исказить данные
  • нормативные справочники (NSI, классификаторы) - если структура справочника сломана, весь последующий JOIN даст неверные результаты
  • первая загрузка новой схемы - при первом запуске нового pipeline FAILFAST быстро выявляет несоответствие схемы контракту
  • регуляторная отчётность - данные с потенциальными ошибками нельзя отправлять в регуляторные системы

Сравнительная таблица режимов

Характеристика PERMISSIVE DROPMALFORMED FAILFAST
Поведение при ошибке NULL + _corrupt_record Тихое удаление Немедленный Exception
Данные в результате Все строки Только валидные Только при 0 ошибок
Видимость потерь Высокая Нулевая Максимальная
Риск потери данных Низкий Высокий Нет
Пригоден для streaming Да Да Нет (останавливает stream)
Применение Bronze Layer, EDA Некритичные данные Финансы, справочники

Dead Letter Queue: badRecordsPath

Режим PERMISSIVE с _corrupt_record хранит брак прямо в основном DataFrame - это смешивает «хорошие» и «плохие» данные в одном потоке. Опция badRecordsPath реализует более архитектурно чистое решение: Dead Letter Queue (DLQ) - автоматический отвод битых записей в отдельное хранилище.

Концепция Dead Letter Queue

DLQ - это паттерн из систем обмена сообщениями (RabbitMQ, Kafka). Смысл: если сообщение не может быть обработано, вместо удаления или блокировки основного потока оно перемещается в «карантинную очередь» для последующего анализа. Основной поток продолжает работу.

В Spark badRecordsPath реализует тот же принцип: корректные записи идут в Bronze-слой Lakehouse, битые - в отдельную папку (DLQ). Инженер может в любой момент проанализировать содержимое DLQ, понять причины ошибок и при необходимости повторно обработать данные.

Диаграмма показывает полный цикл обработки с DLQ: данные из S3 поступают в DataFrameReader, который разделяет корректные и битые записи. Корректные записи идут в Bronze Iceberg table. Битые сохраняются в DLQ, где их подхватывает Airflow-задача мониторинга. Если процент брака превышает порог - отправляется алерт. Если нет - метрика записывается в дашборд качества данных.

Как использовать badRecordsPath

df = (
    spark.read
    .format("json")
    .option("mode", "PERMISSIVE")   # badRecordsPath работает ТОЛЬКО с PERMISSIVE
    .option("badRecordsPath", "s3a://data-lake-dlq/orders/2024-01-15/")
    .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSSZ")
    .schema(order_schema)
    .load("s3a://data-lake-raw/orders/2024-01-15/")
)

Есть несколько ограничений, которые важно понимать:

  • badRecordsPath работает только совместно с режимом PERMISSIVE. С DROPMALFORMED или FAILFAST - игнорируется
  • Поддерживаемые форматы: JSON и CSV. Для Parquet, ORC, Avro - не применимо (эти форматы бинарные и имеют собственные механизмы контроля целостности)
  • В badRecordsPath попадают только структурные ошибки (битый JSON, неправильное число полей в CSV). Ошибки типа поля ("N/A" в числовой колонке) приводят к NULL в поле, но в DLQ не записываются
  • Запись в DLQ создаёт дополнительный IO overhead - для данных с высоким процентом брака (>10%) это может замедлить ingestion

Структура папки брака

Spark создаёт следующую структуру файлов внутри badRecordsPath:

bad_records/
  orders/
    20240115T120000/           ← timestamp запуска job
      bad_records/
        part-00000-abc123.json
        part-00001-def456.json

Каждый файл - это JSON Lines (по одной записи на строку). Каждая запись содержит три поля:

{
  "path": "s3a://data-lake-raw/orders/2024-01-15/orders_batch_042.json",
  "reason": "Malformed JSON: Unexpected character ('{' (code 123)) at position 47",
  "record": "{broken json string here that caused the error..."
}

Значение каждого поля:

  • path - полный путь к исходному файлу, который содержит ошибку. Это критично для диагностики: сразу видно, какой именно файл прислал поставщик с проблемой
  • reason - текстовое описание ошибки от парсера (Jackson для JSON или Univocity для CSV). Позволяет быстро понять тип проблемы: неправильная кодировка, обрыв строки, нарушение структуры
  • record - исходная строка (или JSON-объект), которая вызвала ошибку. Можно смотреть на реальные данные источника и понять, что именно пришло

Анализ брака в downstream-задаче

Папку DLQ можно читать как обычный DataFrame для построения метрик и алертов:

from pyspark.sql.functions import col, count, lit, current_date

today = "2024-01-15"

# Читаем содержимое DLQ - это обычный JSON Lines формат
dlq_df = (
    spark.read
    .format("json")
    .load(f"s3a://data-lake-dlq/orders/{today}/")
)

bad_count  = dlq_df.count()
total_raw  = spark.read.text(f"s3a://data-lake-raw/orders/{today}/").count()
bad_pct    = bad_count / total_raw * 100 if total_raw > 0 else 0

print(f"Брак: {bad_count:,} из {total_raw:,} строк ({bad_pct:.2f}%)")

# Разбивка по типам ошибок - понять, что именно сломалось
print("\nТоп-5 причин ошибок:")
(
    dlq_df
    .groupBy("reason")
    .count()
    .orderBy(col("count").desc())
    .show(5, truncate=80)
)

# Разбивка по исходным файлам - найти проблемные файлы поставщика
print("\nТоп-5 проблемных файлов:")
(
    dlq_df
    .groupBy("path")
    .count()
    .orderBy(col("count").desc())
    .show(5, truncate=120)
)

if bad_pct > 1.0:
    send_alert(
        f"Высокий процент брака: {bad_pct:.1f}% (порог: 1%) "
        f"— требуется проверка схемы источника orders"
    )

Ключевой момент: мониторинг DLQ - это не опциональная часть, а обязательный компонент production ingestion. Без него badRecordsPath теряет смысл: данные изолированы, но никто не знает, что происходит.

Особенности форматов CSV и JSON

CSV и JSON - наиболее проблематичные форматы для ingestion, каждый по-своему. Оба они текстовые и не имеют встроенного механизма валидации структуры.

CSV: мина замедленного действия

CSV выглядит простым - строки, разделители, заголовок. Но на практике это один из самых сложных форматов для надёжного парсинга, потому что стандарта RFC 4180 придерживаются далеко не все системы.

Кавычки и разделители внутри значений. Если поле содержит запятую или кавычку, нужна экранировка. Разные системы используют разные конвенции экранирования. Spark использует Univocity парсер, который по умолчанию следует RFC 4180 (двойные кавычки для экранирования):

"Smith, John","He said ""hello""",100.00

Но реальные CSV-файлы из ERP-систем, банковских выгрузок, 1С часто нарушают стандарт. Опция escape позволяет задать символ экранирования:

.option("quote",  '"')   # символ кавычки (по умолчанию)
.option("escape", '"')   # экранирование внутри кавычек (RFC 4180: двойная кавычка)
# или экранирование обратным слешем (MySQL-стиль):
.option("escape", "\\")

Переносы строк внутри значений. Значение поля может содержать \n, что ломает построчный парсинг:

id,description,amount
1,"Многострочное
описание товара",299.99
2,"Обычное описание",100.00

Для такого CSV необходима опция multiLine=true:

df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("multiLine", "true")
    .option("quote", '"')
    .option("escape", '"')
    .schema(schema)
    .load(path)
)

Предупреждение: при multiLine=true Spark не может читать один файл несколькими executor-ами (файл становится non-splittable). Весь файл читается одним executor-ом последовательно. Для очень больших файлов (>1 ГБ) это может привести к проблемам с памятью executor-а.

JSON: структурная сложность

JSON сложнее тем, что каждая строка может иметь разную структуру (semi-structured data).

Ожидаемый формат: JSON Lines (NDJSON). По умолчанию Spark ожидает один JSON-объект на строку - формат JSON Lines (также называемый NDJSON или Newline Delimited JSON). Это стандарт de facto для потоковых систем (Kafka, Kinesis).

{"order_id": "O1", "amount": 100.00}
{"order_id": "O2", "amount": 200.00}
{"order_id": "O3", "amount": 150.00}

«Красивый» JSON с отступами. Если поставщик присылает файл с форматированием (pretty-printed JSON):

[
  {
    "order_id": "O1",
    "amount": 100.00
  }
]

Нужна опция multiLine=true - с теми же ограничениями по параллелизму.

Тип null vs отсутствие поля. В JSON {"a": null} и {} (поле a вообще отсутствует) - это разные вещи. Spark обрабатывает оба варианта как NULL при явной схеме, что может скрывать разницу между «явным null» и «отсутствующим полем».

Ключевые опции очистки при чтении

Spark предоставляет набор опций для превентивной очистки данных прямо на этапе чтения. Это позволяет избежать дополнительных преобразований в downstream-логике.

nullValue - замена строк на NULL

Опция nullValue задаёт строку, которая при чтении интерпретируется как SQL NULL. Это позволяет превентивно обработать «специальные значения» поставщика до применения типов:

df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("nullValue", "N/A")   # строка "N/A" → NULL для всех колонок
    .schema(schema)
    .load(path)
)

Если в поле amount (DoubleType) стоит "N/A", то без nullValue Spark попытается привести "N/A" к Double, потерпит неудачу, и - в зависимости от режима - запишет NULL (PERMISSIVE) или упадёт (FAILFAST). С nullValue="N/A" строка сначала преобразуется в SQL NULL, и уже NULL корректно присваивается в DoubleType-поле без ошибки парсинга.

Ограничение: nullValue поддерживает только одно значение. Если поставщик использует несколько «пустых» значений ("N/A", "NA", "-", "none", "null"), придётся делать postprocessing:

from pyspark.sql.functions import col, when, trim

# Постобработка: заменить несколько вариантов "пустых" значений на NULL
NULL_SYNONYMS = {"N/A", "NA", "-", "none", "null", ""}

df_clean = df
for column in ["amount", "currency", "category"]:
    df_clean = df_clean.withColumn(
        column,
        when(trim(col(column)).isin(*NULL_SYNONYMS), None).otherwise(col(column))
    )

emptyValue - обработка пустых строк

По умолчанию пустая строка "" в CSV (Alice,,30 - пустой email) становится пустой строкой "" в StringType-колонке, а не NULL. Опция emptyValue меняет это поведение:

# Пустая строка → NULL для строковых колонок
.option("emptyValue", "")    # по умолчанию: пустое значение → ""
# Нет прямого способа задать NULL через emptyValue в текущих версиях Spark
# Для этого используют postprocessing: col("email").eqNullSafe("") → NULL

Для числовых типов пустое значение (,,) всегда становится NULL автоматически - emptyValue для числовых колонок не нужна.

dateFormat и timestampFormat - форматы дат

Даты в CSV и JSON часто приходят в нестандартном формате. Spark по умолчанию ожидает ISO 8601:

  • DateType: yyyy-MM-dd
  • TimestampType: yyyy-MM-dd'T'HH:mm:ss[.SSS][Z]

Для других форматов используйте нотацию java.text.SimpleDateFormat:

df = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("dateFormat",      "dd.MM.yyyy")             # 31.12.2024
    .option("timestampFormat", "dd.MM.yyyy HH:mm:ss")    # 31.12.2024 23:59:59
    .schema(schema)
    .load(path)
)

Если значение не соответствует формату, поведение определяется режимом: PERMISSIVE → NULL в поле даты, FAILFAST → исключение с описанием ошибки.

Полная таблица полезных опций для CSV:

Опция Значение по умолчанию Описание
sep , Разделитель полей
header false Первая строка - заголовок с именами колонок
quote " Символ кавычки для экранирования
escape \ Символ экранирования внутри кавычек
nullValue "" Строка, интерпретируемая как NULL
emptyValue "" Значение для пустых строк в StringType-колонках
dateFormat yyyy-MM-dd Формат для DateType колонок
timestampFormat ISO 8601 Формат для TimestampType колонок
multiLine false Разрешить перенос строк внутри значений (non-splittable)
ignoreLeadingWhiteSpace false Удалять пробелы в начале значения
ignoreTrailingWhiteSpace false Удалять пробелы в конце значения
encoding UTF-8 Кодировка файла (важно для Windows-1251, KOI8-R)
samplingRatio 1.0 Доля строк для вывода схемы (только с inferSchema)
columnNameOfCorruptRecord _corrupt_record Имя колонки для сохранения битых строк
mode PERMISSIVE Режим обработки ошибок
badRecordsPath - Путь для DLQ (только CSV и JSON)

Под капотом: как Catalyst обрабатывает чтение текстовых форматов

За каждым spark.read.csv() стоит многоуровневая обработка через Catalyst Optimizer. Понимание этого механизма позволяет писать более эффективный код и правильно интерпретировать Spark UI.

Физический план чтения

Диаграмма показывает путь данных: от файла в S3 - через Catalyst Optimizer - до executor-ов. Catalyst на этапе оптимизации применяет Column Pruning: если запрос использует только 3 из 20 колонок, Spark всё равно читает весь CSV (у него нет возможности пропустить «лишние» колонки, как в Parquet). Но он исключает ненужные колонки из парсинга - это небольшая, но реальная экономия CPU.

Univocity: парсер CSV

Spark использует библиотеку Univocity - высокопроизводительный Java-парсер CSV. Univocity обрабатывает байт за байтом, соблюдая правила кавычек, разделителей и экранирования согласно конфигурации.

Splittable vs Non-splittable. В режиме без multiLine Univocity является splittable: Spark может разделить один CSV-файл на несколько партиций и читать их параллельно. Executor начинает чтение с произвольной позиции в файле и обрабатывает строки до конца своей партиции. С multiLine=true это невозможно - файл читается целиком одним executor-ом.

Обнаружение ошибок. Когда Univocity встречает незакрытую кавычку, неправильное число колонок или другую структурную ошибку, он передаёт исходную строку в обработчик _corrupt_record. Обработчик - в зависимости от режима - либо записывает строку в _corrupt_record-колонку DataFrame, либо в файл DLQ, либо выбрасывает исключение.

Jackson: парсер JSON

Для JSON Spark использует библиотеку Jackson. Jackson строгий к синтаксису JSON: обязательны двойные кавычки у ключей, запрещены комментарии и завершающие запятые (trailing commas). Spark предоставляет опции для включения расширенного режима:

.option("allowComments",           "true")  # разрешить // и /* */ комментарии
.option("allowUnquotedFieldNames", "true")  # разрешить ключи без кавычек
.option("allowSingleQuotes",       "true")  # разрешить одинарные кавычки
.option("allowNumericLeadingZeros","true")  # разрешить "007" как число

Эти опции расширяют совместимость с нестандартными JSON-генераторами, но немного снижают производительность парсинга из-за дополнительных проверок.

Память при парсинге вложенных JSON. Jackson создаёт дерево объектов в памяти executor-а при парсинге вложенных структур. Для глубоко вложенных JSON (5+ уровней, сотни полей) это может потреблять значительное количество heap. При OOM executor-а следует проверить глубину вложенности входных данных и при необходимости увеличить spark.executor.memory.

Почему нет predicate pushdown для CSV и JSON

В отличие от Parquet (где фильтры могут исключить чтение целых Row Groups на основе статистики min/max), CSV и JSON не поддерживают predicate pushdown на уровне файла или блока данных. Причина: у этих форматов нет встроенной статистики.

Что это означает на практике: если вы читаете CSV с фильтром WHERE event_date = '2024-01-15', Spark обязан прочитать весь файл побайтово, преобразовать каждую строку в InternalRow и только потом применить фильтр. Нет никакого способа «перепрыгнуть» через строки, которые не соответствуют условию.

Именно поэтому Bronze-слой в Lakehouse - это конвертация CSV/JSON в Parquet/Iceberg как можно раньше. Один раз потратить IO на конвертацию - и потом экономить в десятки раз на каждом аналитическом запросе.

Schema Drift: когда схема источника меняется

Schema drift (дрейф схемы) - ситуация, когда поставщик данных изменяет структуру файлов без предупреждения. Это реальная и частая проблема в production, особенно с внешними поставщиками.

Типичные сценарии schema drift

Добавление нового поля. Поставщик добавил поле discount_amount в CSV. При явной схеме без этого поля Spark игнорирует новое поле - оно просто не попадает в DataFrame. Данные читаются корректно, но новое поле теряется. Это может быть проблемой, если поле важно для бизнеса.

Удаление поля. Поле legacy_code исчезло из файла. При явной схеме с nullable=True это поле становится NULL для всех строк. При nullable=False - значение всё равно будет NULL, но Catalyst может генерировать неоптимальный код.

Изменение типа поля. Поле user_id изменилось с Long на String. Spark при PERMISSIVE попытается привести строку к Long; если не получится - NULL в поле. При FAILFAST - немедленное исключение.

Изменение кодировки. Поставщик поменял кодировку с UTF-8 на Windows-1251. Spark прочитает данные, но все строки с кириллицей будут отображаться как «кракозябры». Это особенно коварная ошибка - данные технически читаются, но содержат мусор.

Стратегии обработки schema drift

# Стратегия 1: Читать как строку, затем парсить с явной схемой
# Позволяет сохранить оригинальную строку и получить типизированные поля
from pyspark.sql.functions import from_json, col

raw_df = spark.read.text(source_path)  # каждая строка - StringType

parsed_df = raw_df.select(
    from_json(col("value"), order_schema).alias("data"),
    col("value").alias("_raw")          # сохраняем оригинальную строку
).select("data.*", "_raw")

# Теперь можно фильтровать по ошибкам парсинга:
# если order_id = NULL, но _raw не пустой - значит, структура JSON изменилась
bad_parse = parsed_df.filter(col("order_id").isNull() & col("_raw").isNotNull())


# Стратегия 2: mergeSchema для Parquet - автоматическое слияние схем
# Работает только для Parquet (и ORC), не для CSV/JSON
df = (
    spark.read
    .option("mergeSchema", "true")
    .parquet("s3a://data/partitioned/")
)
# Spark объединит схемы всех Parquet-файлов:
# если в одном файле есть колонка X, а в другом нет - X будет в итоговом DataFrame,
# но со значением NULL для файлов, где её не было


# Стратегия 3: Мониторинг новых полей через inferSchema + diff
current_schema = explicit_schema
inferred_schema = spark.read.option("inferSchema", "true").csv(path).schema

new_fields     = set(inferred_schema.fieldNames()) - set(current_schema.fieldNames())
removed_fields = set(current_schema.fieldNames()) - set(inferred_schema.fieldNames())

if new_fields:
    send_alert(f"Обнаружены новые поля в источнике: {new_fields} - обновите схему!")
if removed_fields:
    send_alert(f"Поля удалены из источника: {removed_fields} - проверьте контракт!")

Стратегия 1 (читать как текст + from_json) - наиболее гибкая для production. Она позволяет сохранить исходную строку (_raw) и одновременно получить типизированные поля. При изменении схемы источника данные не теряются - они остаются в _raw, и инженер может дополнительно обработать их.

Schema Enforcement в Lakehouse

Современные форматы таблиц - Apache Iceberg и Delta Lake - добавляют второй рубеж защиты: schema enforcement на уровне таблицы. Это дополняет контроль на уровне чтения.

Как это работает

При записи DataFrame в Iceberg или Delta таблицу система сравнивает схему DataFrame с зарегистрированной схемой таблицы в метаданных. Если схемы не совпадают - выбрасывается исключение. Данные не записываются.

# Попытка записать DataFrame с лишней колонкой в Iceberg-таблицу
df_with_extra.write.format("iceberg").mode("append").saveAsTable("bronze.orders")
# → AnalysisException: Cannot write incompatible data to table 'bronze.orders'

# Для добавления новой колонки в таблицу Iceberg:
spark.sql("ALTER TABLE bronze.orders ADD COLUMN discount_amount DOUBLE")
# После этого запись DataFrame с discount_amount пройдёт успешно

# Delta Lake: автоматическое добавление новых колонок через mergeSchema
df_with_extra.write \
    .format("delta") \
    .option("mergeSchema", "true") \
    .mode("append") \
    .saveAsTable("bronze.orders")

Schema enforcement в Lakehouse и явные схемы при чтении образуют двойную защиту: первый рубеж - при чтении (что мы принимаем из источника), второй рубеж - при записи (что мы сохраняем в хранилище). Такой подход обеспечивает надёжный data contract на всём протяжении pipeline.

Практика: строим отказоустойчивый Ingestion Gate

Применим полученные знания для построения реального production-grade ingestion pipeline.

Сценарий: внешний поставщик ежедневно кладёт JSON-файлы с заказами в S3-бакет. Файлы иногда содержат проблемы: нарушенная кодировка, изменённые типы, обрезанные строки, неправильный формат дат. Нужно построить надёжный pipeline, который принимает данные «как есть», изолирует брак и оповещает команду при превышении допустимого порога.

Шаг 1: Явная схема - фиксируем контракт

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

# Полная схема заказа - зафиксированный контракт с поставщиком
# Обновляется только при явном изменении контракта (версионирование)
ORDER_SCHEMA = StructType([
    StructField("order_id",    StringType(),    nullable=False),
    StructField("customer_id", LongType(),      nullable=False),
    StructField("total",       DoubleType(),    nullable=True),
    StructField("currency",    StringType(),    nullable=True),
    StructField("status",      StringType(),    nullable=True),
    StructField("created_at",  TimestampType(), nullable=True),
    StructField("is_paid",     BooleanType(),   nullable=True),
    StructField("items_count", IntegerType(),   nullable=True),
    # Техническая колонка для PERMISSIVE-режима - ОБЯЗАТЕЛЬНО в схеме
    # Без неё Spark не будет заполнять _corrupt_record при структурных ошибках
    StructField("_corrupt_record", StringType(), nullable=True),
])

Схема точно отражает контракт с поставщиком. Поле _corrupt_record явно включено - без этого Spark не заполняет его даже в PERMISSIVE-режиме.

Шаг 2: Чтение с изоляцией брака

from datetime import date
from pyspark.sql.functions import col, count, current_timestamp, lit

today      = date.today().strftime("%Y-%m-%d")
source_path = f"s3a://data-lake-raw/orders/{today}/"
dlq_path    = f"s3a://data-lake-dlq/orders/{today}/"

# Чтение с PERMISSIVE и badRecordsPath
# badRecordsPath изолирует структурные ошибки (битый JSON)
# _corrupt_record в схеме ловит те же ошибки внутри DataFrame
raw_df = (
    spark.read
    .format("json")
    .option("mode", "PERMISSIVE")
    .option("badRecordsPath", dlq_path)
    .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSSZ")
    .option("nullValue", "N/A")          # "N/A" → NULL во всех полях
    .schema(ORDER_SCHEMA)
    .load(source_path)
)

# cache() важен: без него последующие count() будут заново читать файлы
raw_df.cache()

total_count   = raw_df.count()
corrupt_df    = raw_df.filter(col("_corrupt_record").isNotNull())
corrupt_count = corrupt_df.count()
null_id_count = raw_df.filter(col("order_id").isNull() & col("_corrupt_record").isNull()).count()
clean_count   = total_count - corrupt_count

corrupt_pct   = corrupt_count / total_count * 100 if total_count > 0 else 0

print(f"[Ingestion Report] {today}")
print(f"  Всего строк:              {total_count:,}")
print(f"  Структурных ошибок:       {corrupt_count:,}  ({corrupt_pct:.2f}%)")
print(f"  NULL order_id (тип?):     {null_id_count:,}")
print(f"  Корректных:               {clean_count:,}")

cache() здесь принципиален: без него каждый вызов count() заново читал бы файлы с S3. Три count() без кэша = три полных сканирования. С cache() - одно чтение с S3, остальные работают из памяти executor-ов.

Шаг 3: Запись в Bronze Layer

# Добавляем технические метаданные перед записью - стандарт data lineage
bronze_df = (
    raw_df
    .filter(col("_corrupt_record").isNull())  # только корректные строки
    .drop("_corrupt_record")                  # убираем техническую колонку из основного слоя
    .withColumn("_ingest_ts",    current_timestamp())   # время загрузки
    .withColumn("_source_path",  lit(source_path))      # путь к исходному файлу
    .withColumn("_ingest_date",  lit(today))            # дата загрузки (для партиционирования)
)

# Записываем в Iceberg-таблицу Bronze-слоя с партиционированием по дате
(
    bronze_df
    .write
    .format("iceberg")
    .mode("append")
    .partitionBy("_ingest_date")
    .saveAsTable("bronze.orders")
)

raw_df.unpersist()  # освобождаем кэш после записи
print(f"Записано в bronze.orders: {clean_count:,} строк")

Технические поля _ingest_ts, _source_path, _ingest_date - стандарт Bronze-слоя. Они обеспечивают data lineage: в любой момент можно найти исходный файл, из которого пришла конкретная запись, и понять, когда она была загружена.

Шаг 4: Мониторинг и алерты

import json
import urllib.request

SLACK_WEBHOOK_URL = "https://hooks.slack.com/services/XXXXX"
ALERT_THRESHOLD_PCT = 1.0  # 1% брака - стандартный порог

def send_slack_alert(message: str) -> None:
    """Отправить сообщение в Slack через Incoming Webhook."""
    payload = json.dumps({"text": message}).encode("utf-8")
    req = urllib.request.Request(
        SLACK_WEBHOOK_URL,
        data=payload,
        headers={"Content-Type": "application/json"}
    )
    urllib.request.urlopen(req)

# Анализ содержимого DLQ - только если есть структурные ошибки
if corrupt_count > 0:
    dlq_df = spark.read.format("json").load(dlq_path)

    print("\n[DLQ Analysis] Топ-5 причин структурных ошибок:")
    (
        dlq_df
        .groupBy("reason")
        .agg(count("*").alias("count"))
        .orderBy(col("count").desc())
        .show(5, truncate=80)
    )

    print("\n[DLQ Analysis] Топ-5 проблемных файлов:")
    (
        dlq_df
        .groupBy("path")
        .agg(count("*").alias("bad_records"))
        .orderBy(col("bad_records").desc())
        .show(5, truncate=120)
    )

# Алерт при превышении порога
if corrupt_pct > ALERT_THRESHOLD_PCT:
    alert_msg = (
        f":warning: *Ingestion Alert* - {today}\n"
        f"Источник: `orders` (S3: `{source_path}`)\n"
        f"Структурный брак: *{corrupt_count:,}* из {total_count:,} строк "
        f"(*{corrupt_pct:.1f}%* > порог {ALERT_THRESHOLD_PCT}%)\n"
        f"DLQ: `{dlq_path}`\n"
        f"Требуется проверка схемы или кодировки источника!"
    )
    send_slack_alert(alert_msg)
    print(f"\n[ALERT] Отправлен Slack-алерт: брак {corrupt_pct:.1f}%")

Шаг 5: Полный pipeline как функция

from dataclasses import dataclass

@dataclass
class IngestionMetrics:
    date: str
    total: int
    clean: int
    corrupt: int
    null_ids: int
    corrupt_pct: float
    alert_triggered: bool


def run_orders_ingestion(
    spark,
    source_path: str,
    dlq_path: str,
    target_table: str = "bronze.orders",
    alert_threshold_pct: float = 1.0,
) -> IngestionMetrics:
    """
    Production-grade ingestion заказов из JSON.
    Возвращает метрики запуска для записи в систему мониторинга.
    """
    raw_df = (
        spark.read
        .format("json")
        .option("mode", "PERMISSIVE")
        .option("badRecordsPath", dlq_path)
        .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSSZ")
        .option("nullValue", "N/A")
        .schema(ORDER_SCHEMA)
        .load(source_path)
    ).cache()

    total   = raw_df.count()
    corrupt = raw_df.filter(col("_corrupt_record").isNotNull()).count()
    null_id = raw_df.filter(col("order_id").isNull() & col("_corrupt_record").isNull()).count()
    clean   = total - corrupt
    bad_pct = corrupt / total * 100 if total > 0 else 0.0

    if clean > 0:
        (
            raw_df
            .filter(col("_corrupt_record").isNull())
            .drop("_corrupt_record")
            .withColumn("_ingest_ts",   current_timestamp())
            .withColumn("_source_path", lit(source_path))
            .withColumn("_ingest_date", lit(source_path.split("/")[-2]))
            .write
            .format("iceberg")
            .mode("append")
            .saveAsTable(target_table)
        )

    raw_df.unpersist()

    metrics = IngestionMetrics(
        date=today,
        total=total,
        clean=clean,
        corrupt=corrupt,
        null_ids=null_id,
        corrupt_pct=round(bad_pct, 4),
        alert_triggered=bad_pct > alert_threshold_pct,
    )

    if metrics.alert_triggered:
        send_slack_alert(format_ingestion_alert(metrics, source_path))

    return metrics


# Использование в Airflow DAG:
# metrics = run_orders_ingestion(spark, source_path, dlq_path)
# push_metrics_to_prometheus(metrics)

Функция инкапсулирует полный цикл: чтение → изоляция брака → запись в Bronze → сбор метрик → алертинг. Возвращаемый IngestionMetrics можно записать в Prometheus, ClickHouse или любую другую систему мониторинга для исторического анализа трендов качества данных.

Диагностика ingestion проблем

Когда ingestion даёт неожиданные результаты или падает, используйте следующий диагностический чеклист.

Шаг 1: Убедиться, что Spark видит файлы

# Список файлов в директории источника
dbutils.fs.ls(source_path)  # Databricks
# или через subprocess:
import subprocess
result = subprocess.run(["hadoop", "fs", "-ls", source_path], capture_output=True, text=True)
print(result.stdout)

Шаг 2: Посмотреть на сырые данные

# Прочитать первые строки как текст - без парсинга, «как есть»
spark.read.text(source_path).show(10, truncate=False)

Это позволяет увидеть реальный формат данных, кодировку, разделители - без интерпретации Spark.

Шаг 3: Анализ _corrupt_record

# Смотреть на реальные битые строки
df.filter(col("_corrupt_record").isNotNull()) \
  .select("_corrupt_record") \
  .show(10, truncate=False)

Содержимое _corrupt_record - исходная строка из файла. По ней сразу видно, что именно пришло от поставщика и где расхождение со схемой.

Шаг 4: Проверить план выполнения

# Физический план: убедиться, что колонки читаются правильно
df.select("order_id", "total").explain(mode="formatted")
# В плане ищите: FileScan json [...] PushedFilters: [...] ReadSchema: ...

Шаг 5: Spark UI - вкладка Stages

В Spark UI для каждого Stage видно количество Input Records и Output Records. Если Output Records значительно меньше Input Records - возможно, что-то фильтруется раньше ожидаемого. Также в Spark UI видны ошибки executor-ов в столбце Task List.

Best Practices

Никогда не используйте inferSchema=true в production. Это двойное сканирование, нестабильная схема и нарушение воспроизводимости. inferSchema - инструмент для EDA и прототипирования в Jupyter-ноутбуке.

Всегда объявляйте явную схему через StructType. Это контракт с источником данных. Если схема источника изменилась - вы узнаете немедленно (через _corrupt_record или FAILFAST), а не спустя неделю при сравнении агрегатов.

Выбирайте режим чтения осознанно. Bronze-слой: PERMISSIVE + _corrupt_record + badRecordsPath - максимальная видимость, ни одна строка не теряется. Финансы и справочники: FAILFAST - немедленное обнаружение проблем. DROPMALFORMED - избегать; исключение только при наличии внешнего reconciliation и мониторинга потерь.

Включайте _corrupt_record в схему явно. Без этого Spark не будет заполнять колонку, и информация о структурных ошибках молча теряется.

Всегда изолируйте брак через badRecordsPath. Это позволяет анализировать причины ошибок, видеть тренды качества данных по файлам источника и повторно обрабатывать данные после исправления схемы.

Мониторинг процента брака обязателен. Порог в 1% - разумный старт для большинства источников. Для финансовых данных - 0.001%. Записывайте метрики исторически: рост процента брака - ранний сигнал деградации источника.

Добавляйте технические метаданные при записи в Bronze. _ingest_ts, _source_path, _ingest_date - минимальный набор для data lineage. Без них восстановить происхождение данных через месяц практически невозможно.

CSV и JSON - форматы для ingestion, не для хранения. Конвертируйте в Parquet или Iceberg как можно раньше в Bronze-слое. Одно IO на конвертацию экономит в 10–100× на каждом последующем аналитическом запросе.