Troubleshooting — Data Engineering для джунов
База знаний типичных ошибок курса Data Engineering для джунов.
Категория
Симптомы
- После повторного запуска failed task в Airflow в целевой таблице появляются строки-двойники с одинаковыми business-ключами. SELECT COUNT(*) показывает в 2-3 раза больше записей, чем должно быть.
Причина
Пайплайн использует INSERT вместо MERGE/UPSERT Нет дедупликации по ingestion-ключу или event_id Retry-логика не очищает партицию перед повторной загрузкой Источник присылает события с at-least-once гарантией
Решение
- Заменить INSERT на MERGE по natural key или (event_id, event_timestamp)
- Перед загрузкой удалять целевую партицию: DELETE FROM target WHERE dt = '{ds}'
- Добавить уникальный constraint на ingestion-ключ и обрабатывать конфликт через ON CONFLICT
- Сделать пайплайн идемпотентным: переписывать партицию целиком вместо инкрементального INSERT
Симптомы
- В Airflow UI один task в середине DAG в статусе failed, все downstream tasks в upstream_failed или skipped. DAG run в статусе failed, но данные не обновились за день.
Причина
Default trigger_rule всех downstream задач — all_success Нет настройки retries для падающей задачи Sensor poll-ит несуществующий ресурс и ловит timeout Зависимость от внешнего сервиса без circuit breaker
Решение
- Настроить retries=3 и retry_delay для нестабильной задачи
- Использовать trigger_rule=all_done для задач, которые должны выполниться даже при сбое предыдущих (cleanup, alerting)
- Заменить poke-sensor на reschedule-mode, чтобы освобождать worker slot
- Добавить SLA и алерты в Slack/PagerDuty, чтобы дежурный мог быстро ручно перезапустить
Симптомы
- Перезапустили DAG за прошлый месяц, в итоге часть данных в целевых витринах исчезла. Сравнение row_count до и после backfill показывает уменьшение.
Причина
DAG помечает партиции как обработанные, но не очищает downstream таблицы Источник за это время потерял данные (TTL, удалили старые записи) Backfill использует SELECT с фильтром, который не учитывает late-arriving events Изменилась схема источника, новые поля не подцепились
Решение
- Перед backfill сделать снапшот через CREATE TABLE AS SELECT для возможности отката
- Использовать time travel в Snowflake / Delta Lake для проверки старого состояния
- Backfill-ить с запасом по диапазону дат (захватить ±1 день для late events)
- Хранить ingestion-метку в источнике (high-water mark) и backfill-ить по ней, а не только по event_date
Симптомы
- Аналитик жалуется, что данные за вчера в дашборде неполные. При проверке источника видно, что часть событий пришла в систему с задержкой 6-24 часа.
Причина
Streaming-pipeline закрывает window слишком рано по event_time Watermark выставлен агрессивно — отбрасывает события старше N минут Batch-загрузка фиксирует партицию строго по календарным суткам без grace period Источник присылает события с неправильным timestamp (created_at vs ingested_at)
Решение
- Сдвинуть watermark на разумный SLA источника (1-6 часов)
- В batch — пересчитывать предыдущий день после окончания grace period (например, в 6 утра пересчитать вчера)
- Использовать ingestion_time как partition key, event_time как column — backfill становится проще
- Для streaming — отправлять late events в side output для отдельной обработки
Симптомы
- Утром Airflow DAG падает с ошибкой column not found или type mismatch. В источнике (Postgres-таблица, API-ответ) изменилась схема — добавлено или переименовано поле.
Причина
Команда источника изменила схему без уведомления Ingestion-скрипт жёстко описывает список колонок без graceful handling Нет contract testing между producer и consumer Schema-on-write DWH отвергает новые столбцы
Решение
- Использовать SELECT * с loose schema в staging-слое, валидация — в трансформациях
- Внедрить Fivetran/Airbyte с auto-schema-evolution для управляемого добавления колонок
- Зафиксировать контракт через JSON Schema / dbt source freshness + tests
- В Snowflake — использовать VARIANT для приёма semi-structured данных без жёсткой схемы
Симптомы
- Spark-job идёт 4 часа, на Spark UI видно что 99% executors завершились за 5 минут, а один тащит свою партицию ещё 3.5 часа. CPU utilization кластера около 5%.
Причина
Data skew: один ключ в groupBy/join имеет диспропорционально много значений (NULL, default user_id, hot key) Партиционирование по полю с низкой cardinality Join с большой таблицей без broadcast Сериализация Python-функций через UDF без vectorization
Решение
- Использовать salting: добавить случайный суффикс к ключу для распределения hot key по нескольким партициям
- Включить Adaptive Query Execution (spark.sql.adaptive.enabled=true) — Spark сам разрулит skew
- Заменить join с маленькой dimension на broadcast join
- Перейти с PySpark UDF на pandas_udf / Spark SQL functions
Симптомы
- Запросы поверх S3-партиций становятся всё медленнее. Spark тратит больше времени на listing файлов, чем на их чтение. В партиции — десятки тысяч мелких Parquet-файлов по 1-10 KB.
Причина
Streaming-ingestion пишет по одному файлу на батч/окно Spark coalesce/repartition не настроен на нужное число партиций Нет compaction-джоба, который периодически объединяет мелкие файлы Каждый retry создаёт новый файл вместо перезаписи
Решение
- Запустить OPTIMIZE на Delta-таблице или rewrite_data_files в Iceberg для компактизации
- В Spark перед записью делать repartition(N), целясь в файлы по 128-512 MB
- Настроить cron-job для еженедельной компактизации старых партиций
- Использовать table format (Iceberg/Delta) вместо raw Parquet — они дают встроенные инструменты компактизации
Симптомы
- Аналитик считает SUM(revenue) — получает в 3 раза больше, чем в источнике. При группировке по date выясняется, что одна транзакция записана несколько раз с разными dimension-комбинациями.
Причина
Fact построен через join нескольких источников без явного определения grain Денормализация раздула таблицу: одна продажа умножилась на N позиций Двойной join по dimension-таблицам без уникального ключа Изменение granularity между версиями модели
Решение
- Зафиксировать grain в комментариях/документации модели dbt: «one row = one transaction line»
- Добавить uniqueness-тест в dbt: unique по (transaction_id, line_number)
- Разделить fact на несколько таблиц по grain: fact_sales_header (на транзакцию) и fact_sales_line (на позицию)
- Перед UNION проверять row_count: SELECT COUNT(*) на каждом шаге через DBT-тесты
Симптомы
- В dim_customer Type 2 у клиента видно только одну строку с is_current=true, хотя за последний год было 5 изменений. История потеряна.
Причина
Пайплайн SCD2 каждый раз делает full reload вместо инкремента Логика SCD2 не закрывает старую версию (valid_to остаётся NULL у нескольких строк) Hash-сравнение полей пропускает изменения в значимых атрибутах При retry merge перезаписал valid_to у уже закрытых версий
Решение
- Использовать готовый snapshot из dbt — он реализует SCD2 правильно
- Тесты на уникальность по (natural_key, valid_to) и проверку «valid_from < valid_to»
- Хешировать только бизнес-значимые поля, исключая audit-поля типа updated_at
- Восстановить историю из source CDC-логов или audit-таблицы источника
Симптомы
- Spark-job падает с java.lang.OutOfMemoryError: Java heap space или Container killed by YARN for exceeding memory limits. Падает чаще всего на shuffle-стадии или collect.
Причина
Слишком много данных в одной партиции после shuffle Использование collect() на больших датасетах — всё стягивается в driver Cache/persist слишком больших DataFrame без unpersist Spark.executor.memory недостаточно для объёма данных, executor.memoryOverhead занижен
Решение
- Увеличить spark.executor.memory и spark.executor.memoryOverhead (10-20% от executor memory)
- Сделать repartition перед shuffle-операциями, чтобы выровнять размер партиций
- Заменить collect() на write() или take(N) для семплов
- Включить off-heap memory: spark.memory.offHeap.enabled=true, spark.memory.offHeap.size=2g
Симптомы
- Месячный счёт за DWH вырос в 2-3 раза без увеличения объёма данных или числа пользователей. Финансы спрашивают, почему.
Причина
Кто-то запускает SELECT * на огромных таблицах без LIMIT Snowflake warehouse не auto-suspend-ится после запросов BI-инструмент дёргает DWH запросами каждые 30 секунд для realtime-дашборда Materialized views пересчитываются на каждое изменение источника
Решение
- Включить query_history-аудит, найти топ-10 запросов по cost (Snowflake) или bytes_billed (BigQuery)
- Настроить auto_suspend=60 на warehouse Snowflake
- В BigQuery — добавить partition filter requirement и cap на bytes_billed
- Кэшировать BI-запросы на стороне приложения (Looker persistent derived tables, Tableau extracts)
Симптомы
- Бизнес-команда обнаруживает в дашборде клиента с возрастом 150 лет, отрицательную выручку, дубликаты заказов. Доверие к данным падает.
Причина
Нет DQ-тестов на ключевые поля (range, nullability, uniqueness) Тесты только на staging-слое, не на витринах Алерты не настроены, ошибки обнаруживаются вручную Источник присылает мусорные значения (-1, 9999-12-31, empty string)
Решение
- Добавить dbt-тесты: not_null, unique, accepted_values, relationships
- Внедрить Great Expectations или Soda для более сложных проверок
- Настроить Monte Carlo / Soda для observability с алертами в Slack
- Очищать мусорные значения на staging-слое: replace(-1, NULL), CASE WHEN age > 120 THEN NULL
Симптомы
- Запрос с join fact_orders (1 млрд строк) на dim_customer (100 млн строк) в Spark/Snowflake идёт часами. EXPLAIN показывает full table scan на dim.
Причина
Dimension слишком большая для broadcast join Нет clustering / sort key по customer_id Spark делает shuffle hash join вместо sort-merge В Snowflake — нет clustering key, micro-partitions перепутаны
Решение
- Кластеризовать dim и fact по join-ключу: ALTER TABLE ... CLUSTER BY (customer_id)
- Денормализовать часть dim в fact (Big Table approach в BigQuery)
- Использовать bucketing в Hive/Spark по customer_id с одинаковым числом бакетов
- Pre-aggregate в materialized view или snapshot, если запросы предсказуемы
Симптомы
- В витрине fact_revenue колонка discount_amount, аналитик спрашивает «откуда она?». Никто не знает: модель собрана из 5 источников через 3 уровня трансформаций.
Причина
SQL-модели написаны через хардкод имён таблиц вместо ref()/source() Нет автоматизированного data catalog Документация модели пуста или устарела Несколько команд пишут трансформации в разных репозиториях
Решение
- Перевести трансформации в dbt с использованием ref() — автоматический lineage
- Внедрить OpenLineage + Marquez / DataHub для автоматического сбора метаданных
- Добавить description в YAML-конфиг dbt-моделей и колонок
- Регулярно ревьюить data catalog в командных митингах
Симптомы
- DAG-и, которые должны были запуститься в 02:00, фактически стартуют в 02:45. UI показывает большое scheduler heartbeat delay. Sensor-задачи висят в очереди.
Причина
Слишком много активных DAG-ов на одном scheduler-е Скрипты DAG-ов делают тяжёлые операции на парсинге (db-запросы на верхнем уровне) Недостаточно worker slot-ов в пуле Слишком много poke-sensor-ов забивают слоты
Решение
- Включить HA-scheduler (Airflow 2+): несколько scheduler-экземпляров
- Уменьшить parsing-frequency: AIRFLOW__CORE__DAG_FILE_PROCESSOR_TIMEOUT=120
- Вынести тяжёлую логику из top-level DAG-кода в task-функции
- Перевести sensor-ы в reschedule mode
Симптомы
- Lag в Kafka топике растёт от часа к часу: consumer не успевает за producer. Дашборд показывает 10 минут, потом 30 минут, потом 2 часа отставания.
Причина
Слишком мало consumer-инстансов относительно числа партиций Тяжёлая обработка на каждое сообщение (внешний API, медленная запись в DB) Producer резко увеличил throughput (например, после маркетинговой кампании) Партиций меньше, чем нужно для параллелизма
Решение
- Увеличить число partition в топике (но это влияет на ordering)
- Запустить больше consumer-инстансов в той же consumer group (до числа партиций)
- Batch-ить запись в DWH/DB вместо одиночных insert
- Вынести тяжёлую обработку в отдельный downstream-процесс через Kafka Streams
Симптомы
- dbt run --select my_model завершается успешно, но в таблице нет последних данных. Source имеет новые строки, в инкрементальной модели — нет.
Причина
is_incremental() условие фильтрует по неправильному полю unique_key не уникален, MERGE не вставляет дубли Watermark подсчитывается через MAX(updated_at), но в источнике updated_at не обновляется Партиция уже была обработана при первом запуске
Решение
- Проверить SQL под is_incremental(): WHERE event_date >= (SELECT MAX(event_date) FROM {{ this }})
- Использовать корректный unique_key или композитный ключ
- Перейти на ingestion-timestamp от источника вместо source-side updated_at
- Run --full-refresh для пересборки модели с нуля
Симптомы
- Настроили Debezium для CDC из Postgres в Kafka. UPDATE в источнике есть, но соответствующего сообщения в Kafka нет. Часть строк теряется.
Причина
wal_level в Postgres не logical, а replica REPLICA IDENTITY таблицы не FULL — UPDATE без изменения PK не пишется Replication slot накопил слишком много WAL и был очищен Sequence изменений превышает retention в Kafka
Решение
- ALTER SYSTEM SET wal_level = logical, перезапустить Postgres
- ALTER TABLE my_table REPLICA IDENTITY FULL
- Мониторить размер replication slot: SELECT pg_size_pretty(pg_wal_lsn_diff(...))
- Увеличить retention в Kafka и настроить алерты на consumer lag
Симптомы
- Spark-job падает с сообщением «Job aborted due to stage failure: Task failed 4 times in stage X». В логах executor-а — разные ошибки на разных попытках.
Причина
Нестабильное окружение: один executor падает из-за OOM, другой — из-за shuffle service Bad data в источнике: одна строка с невалидным JSON ломает парсинг Race condition при записи: несколько executors пишут в один файл Внешний API в UDF timeout-ит
Решение
- Открыть Spark UI -> Stages -> найти failed task -> посмотреть exception в exec logs
- Добавить try-catch в UDF и помечать bad rows вместо exception
- Увеличить spark.task.maxFailures с 4 до 8 для retry-able задач
- Если разные tasks падают с одной ошибкой — проблема в коде, не в инфре
Симптомы
- В Slack ежедневно вопросы «где данные о клиентах?», «есть ли таблица с заказами по регионам?». Один и тот же датасет создаётся параллельно в 3 командах.
Причина
Нет data catalog — таблицы прячутся в произвольных schemas Naming convention отсутствует или не соблюдается Документация моделей не существует или устарела Каждая команда строит свои витрины в изоляции
Решение
- Внедрить data catalog: DataHub, OpenMetadata, Amundsen
- Зафиксировать naming convention: prefix_team_entity_grain (fact_sales_order_daily)
- В dbt — обязать description для каждой модели и колонки через CI-проверку
- Назначить data steward-ов от каждой команды
Симптомы
- Финансовый дашборд показывает данные за позавчера, хотя пайплайн должен обновлять каждое утро. Бизнес не понимает, можно ли доверять цифрам.
Причина
Пайплайн упал ночью, но никто не отреагировал на alert Cache в BI-инструменте не инвалидируется Источник не обновился из-за проблем на их стороне DAG висит в queued из-за нехватки worker slot-ов
Решение
- Добавить freshness-метрики: dbt source freshness, table_last_modified_at в каталоге
- На дашборде показывать last_updated_at и подсветка если данные >24h old
- Настроить алерты на SLA через PagerDuty / Slack с дежурной ротацией
- Инвалидировать BI-кэш по событию завершения пайплайна, не по таймеру
Симптомы
- В дашборде маркетинга revenue за апрель — 1.2M, в дашборде финансов — 1.18M. Бизнес требует объяснить разницу, обе стороны уверены в своих числах.
Причина
Разные источники истины: маркетинг тянет из CRM, финансы — из ERP Разная семантика метрики: gross revenue vs net revenue с возвратами Разные временные срезы (timezone, дата заказа vs дата отгрузки) Один из дашбордов использует stale snapshot
Решение
- Создать single source of truth: certified marts с подписью owner
- Согласовать определения метрик в semantic layer (Cube, dbt metrics, LookML)
- Документировать timezone и event-time семантику в каждой витрине
- При расхождениях — двигать дашборды на certified marts
Симптомы
- Мигрировали data lake с CSV на Parquet, рассчитывая на ускорение. На практике — некоторые запросы стали медленнее.
Причина
Parquet файлы слишком мелкие (по нескольку KB) — небольшие batch-загрузки Запросы с SELECT * — не использует преимущества колоночного формата Predicate pushdown не работает из-за отсутствия min/max статистик Compression codec неоптимален (gzip vs snappy vs zstd)
Решение
- Настроить compaction: целиться в файлы 128-512 MB
- Использовать SELECT only_needed_columns вместо *
- Переключиться с gzip на snappy или zstd: лучше баланс CPU/IO
- Проверить статистики через parquet-tools: должны быть min/max по колонкам
Симптомы
- Запрос на партиционированную таблицу сканирует все партиции вместо одной. SELECT * FROM events WHERE date = '2026-05-17' — обещали быстро, по факту минуты.
Причина
Тип в WHERE не совпадает с типом partition column (STRING vs DATE) Используется функция на partition column: WHERE DATE(event_ts) = '2026-05-17' Partition value — string, а в фильтре — date через implicit cast Snowflake clustering key неверно подобран
Решение
- Явно кастить в нужный тип: WHERE date = CAST('2026-05-17' AS DATE)
- Избегать функций на partition column в WHERE
- Проверить EXPLAIN: должно показывать partition pruning
- В BigQuery — использовать _PARTITIONTIME или DATE-partition правильно
Симптомы
- На локальной машине python pipeline.py отрабатывает за 5 минут. В Airflow тот же код падает с timeout / OOM / connection refused.
Причина
Локально — маленький семпл данных, в проде — миллионы строк Зашиты локальные пути или localhost connection strings В проде нет нужных Python-зависимостей Airflow worker не имеет доступа к нужным секретам/сетям
Решение
- Тестировать на репрезентативной выборке данных
- Всё конфигурировать через env variables и Airflow Connections
- Зафиксировать зависимости в requirements.txt, использовать virtualenv operator
- Проверить network policy / firewall между worker и target системами
Симптомы
- Запустили backfill за 2 года на Snowflake, через 6 часов финдир спрашивает «что вы натворили». Месячный лимит warehouse-а сожжён за сутки.
Причина
Backfill запускается сразу за весь период вместо разбиения по партициям Используется warehouse слишком большого размера Параллельные runs всех партиций одновременно Нет ограничения по cost / runtime в job-е
Решение
- Разбить backfill на партиции по дням и запускать последовательно
- Использовать отдельный backfill-warehouse с меньшим размером
- Установить statement_timeout / resource monitor с hard cap
- Запускать backfill в выходные на off-peak ценах
Симптомы
- SQL-запрос с CTE или window function на 100M строк висит часами или падает по timeout. EXPLAIN показывает огромный shuffle.
Причина
Window function без PARTITION BY — окно по всей таблице Recursive CTE с неконтролируемой глубиной Самoself-join по неиндексированному полю Большие промежуточные результаты в materialized CTE
Решение
- Добавить PARTITION BY user_id / region в window function
- Ограничить рекурсию через MAXRECURSION или явный exit
- Заменить self-join на join с pre-aggregated CTE
- Pre-aggregate в incremental model вместо on-the-fly расчёта
Симптомы
- Старые consumer-ы Kafka-топика начинают падать после деплоя нового producer-а. Spark не может прочитать старые Parquet-файлы после изменения схемы.
Причина
Удалили required-поле — backward-incompatible change Переименовали колонку без alias Изменили тип поля (int -> string) Нет schema registry, контракт нигде не зафиксирован
Решение
- Использовать Confluent Schema Registry с backward/forward compatibility-режимом
- Добавлять новые поля как optional с default, не удалять старые сразу
- Использовать column mapping в Delta Lake / Iceberg для безопасного rename
- Тесты compatibility в CI: pre-merge проверка new schema vs old data