Data Lakehouse в облаке (часть 2): Слой Silver и Iceberg

О чем: Управление бюджетом, запуск динамических Data Proc кластеров через Airflow и YAML-конфиги.

В первой части мы зафиксировали поток сырых JSON-данных на слое Bronze в Yandex Object Storage. Теперь перед нами стоит более сложная задача: превратить этот массив слабоструктурированных файлов в чистый, валидированный и транзакционный слой Silver.

Именно на этом этапе закладывается надежность всей дальнейшей аналитики. Мы разберем, как построить этот процесс эффективно, не перегружая систему тяжелыми внешними фреймворками, и как формат Apache Iceberg решает проблемы, которые раньше казались неразрешимыми в Data Lake.

Слой Silver: Очистка, валидация и кастомный карантин (DLQ)

Основное правило слоя Silver — данные здесь должны быть очищены, приведены к строгой схеме и готовы к использованию внутри компании. Но что делать, если часть записей из источника приходит поврежденной или не соответствует бизнес-требованиям?

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

Паттерн S3 Quarantine (Dead Letter Queue)

Процесс обработки строится по следующему алгоритму: 1. Чтение и разбор схемы: PySpark считывает инкремент из Bronze-слоя и накладывает целевую структуру на JSON.

2. Валидация: С помощью набора кастомных функций (проверки на NOT NULL, соответствие регулярным выражениям, корректность диапазонов дат) каждая строка размечается флагом is_valid.

3. Разделение потоков: - Строки с is_valid = true отправляются дальше по конвейеру. - «Битые» записи (is_valid = false) вместе с метаданными об ошибке изолируются и складываются в отдельный бакет — S3 Quarantine (DLQ). Это позволяет аналитикам и бэкенд-разработчикам исследовать аномалии, не останавливая основной пайплайн.

4. Контроль критического порога (Circuit Breaker): В код джобы заложена проверка: если процент поврежденных данных в текущем батче превышает критический порог (например, более 5%), PySpark прерывает выполнение с ошибкой, а Apache Airflow мгновенно отправляет алерт инженерам. Это защищает систему от массовых сбоев на стороне источников.

Оптимизация PySpark: Как мы победили Out-of-Memory (OOM)

Работа с глубоко вложенными и тяжелыми JSON-структурами в Spark таит в себе скрытые угрозы для оперативной памяти кластера. Парсинг таких данных часто приводит к избыточному раздуванию графа вычислений (Execution Plan) и падению воркеров по OOM.

Чтобы стабилизировать пайплайн, мы применили стратегию контролируемого кэширования промежуточных состояний. На этапе, когда JSON уже распарсен, но валидация и разделение на чистый поток и карантин еще не начались, мы вызываем метод .persist(StorageLevel.MEMORY_AND_DISK).

Почему Apache Iceberg, а не просто Parquet-файлы?

Исторически слой Silver организовывали в виде обычных файлов Parquet. Но в реальном продакшене это быстро упирается в ограничения файловых систем. Apache Iceberg поверх Yandex Object Storage полностью меняет правила игры благодаря трем фичам:

1. Полноценный ACID: Больше нет риска прочитать «наполовину записанные» данные, если Spark-джоба упала посреди транзакции. Читатели видят только успешные коммиты.

2. Schema Evolution: Если бэкенд добавил новое поле в исходный JSON или изменил тип данных, Iceberg позволяет обновить схему таблицы на лету без необходимости полной перезаписи терабайтов исторических данных.

3. Hidden Partitioning: Инженерам больше не нужно вручную подставлять фильтры по партициям в каждый запрос. Iceberg сам знает, как физически лежат данные, и оптимизирует чтение на основе контента, что критически важно для последующей эффективной работы dbt и ClickHouse.

В следующей части мы поднимемся на уровень выше и разберем облачный FinOps: как экономить бюджет Yandex Cloud, управляя жизненным циклом вычислительных кластеров с помощью Apache Airflow и динамических конфигураций.

Data Lakehouse в облаке (часть 2): Слой Silver и Iceberg | Сетка — социальная сеть от hh.ru