Перейти к содержанию
Learning Platform
Глоссарий
Troubleshooting

Troubleshooting — Data Engineering для джунов

База знаний типичных ошибок курса Data Engineering для джунов.

Категория

Показано 28 из 28 ошибок

Симптомы

  • После повторного запуска failed task в Airflow в целевой таблице появляются строки-двойники с одинаковыми business-ключами. SELECT COUNT(*) показывает в 2-3 раза больше записей, чем должно быть.

Причина

Пайплайн использует INSERT вместо MERGE/UPSERT Нет дедупликации по ingestion-ключу или event_id Retry-логика не очищает партицию перед повторной загрузкой Источник присылает события с at-least-once гарантией

Решение

  1. Заменить INSERT на MERGE по natural key или (event_id, event_timestamp)
  2. Перед загрузкой удалять целевую партицию: DELETE FROM target WHERE dt = '{ds}'
  3. Добавить уникальный constraint на ingestion-ключ и обрабатывать конфликт через ON CONFLICT
  4. Сделать пайплайн идемпотентным: переписывать партицию целиком вместо инкрементального 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

Решение

  1. Настроить retries=3 и retry_delay для нестабильной задачи
  2. Использовать trigger_rule=all_done для задач, которые должны выполниться даже при сбое предыдущих (cleanup, alerting)
  3. Заменить poke-sensor на reschedule-mode, чтобы освобождать worker slot
  4. Добавить SLA и алерты в Slack/PagerDuty, чтобы дежурный мог быстро ручно перезапустить

Симптомы

  • Перезапустили DAG за прошлый месяц, в итоге часть данных в целевых витринах исчезла. Сравнение row_count до и после backfill показывает уменьшение.

Причина

DAG помечает партиции как обработанные, но не очищает downstream таблицы Источник за это время потерял данные (TTL, удалили старые записи) Backfill использует SELECT с фильтром, который не учитывает late-arriving events Изменилась схема источника, новые поля не подцепились

Решение

  1. Перед backfill сделать снапшот через CREATE TABLE AS SELECT для возможности отката
  2. Использовать time travel в Snowflake / Delta Lake для проверки старого состояния
  3. Backfill-ить с запасом по диапазону дат (захватить ±1 день для late events)
  4. Хранить 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)

Решение

  1. Сдвинуть watermark на разумный SLA источника (1-6 часов)
  2. В batch — пересчитывать предыдущий день после окончания grace period (например, в 6 утра пересчитать вчера)
  3. Использовать ingestion_time как partition key, event_time как column — backfill становится проще
  4. Для 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 отвергает новые столбцы

Решение

  1. Использовать SELECT * с loose schema в staging-слое, валидация — в трансформациях
  2. Внедрить Fivetran/Airbyte с auto-schema-evolution для управляемого добавления колонок
  3. Зафиксировать контракт через JSON Schema / dbt source freshness + tests
  4. В 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

Решение

  1. Использовать salting: добавить случайный суффикс к ключу для распределения hot key по нескольким партициям
  2. Включить Adaptive Query Execution (spark.sql.adaptive.enabled=true) — Spark сам разрулит skew
  3. Заменить join с маленькой dimension на broadcast join
  4. Перейти с PySpark UDF на pandas_udf / Spark SQL functions

Симптомы

  • Запросы поверх S3-партиций становятся всё медленнее. Spark тратит больше времени на listing файлов, чем на их чтение. В партиции — десятки тысяч мелких Parquet-файлов по 1-10 KB.

Причина

Streaming-ingestion пишет по одному файлу на батч/окно Spark coalesce/repartition не настроен на нужное число партиций Нет compaction-джоба, который периодически объединяет мелкие файлы Каждый retry создаёт новый файл вместо перезаписи

Решение

  1. Запустить OPTIMIZE на Delta-таблице или rewrite_data_files в Iceberg для компактизации
  2. В Spark перед записью делать repartition(N), целясь в файлы по 128-512 MB
  3. Настроить cron-job для еженедельной компактизации старых партиций
  4. Использовать table format (Iceberg/Delta) вместо raw Parquet — они дают встроенные инструменты компактизации

Симптомы

  • Аналитик считает SUM(revenue) — получает в 3 раза больше, чем в источнике. При группировке по date выясняется, что одна транзакция записана несколько раз с разными dimension-комбинациями.

Причина

Fact построен через join нескольких источников без явного определения grain Денормализация раздула таблицу: одна продажа умножилась на N позиций Двойной join по dimension-таблицам без уникального ключа Изменение granularity между версиями модели

Решение

  1. Зафиксировать grain в комментариях/документации модели dbt: «one row = one transaction line»
  2. Добавить uniqueness-тест в dbt: unique по (transaction_id, line_number)
  3. Разделить fact на несколько таблиц по grain: fact_sales_header (на транзакцию) и fact_sales_line (на позицию)
  4. Перед 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 у уже закрытых версий

Решение

  1. Использовать готовый snapshot из dbt — он реализует SCD2 правильно
  2. Тесты на уникальность по (natural_key, valid_to) и проверку «valid_from < valid_to»
  3. Хешировать только бизнес-значимые поля, исключая audit-поля типа updated_at
  4. Восстановить историю из 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 занижен

Решение

  1. Увеличить spark.executor.memory и spark.executor.memoryOverhead (10-20% от executor memory)
  2. Сделать repartition перед shuffle-операциями, чтобы выровнять размер партиций
  3. Заменить collect() на write() или take(N) для семплов
  4. Включить 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 пересчитываются на каждое изменение источника

Решение

  1. Включить query_history-аудит, найти топ-10 запросов по cost (Snowflake) или bytes_billed (BigQuery)
  2. Настроить auto_suspend=60 на warehouse Snowflake
  3. В BigQuery — добавить partition filter requirement и cap на bytes_billed
  4. Кэшировать BI-запросы на стороне приложения (Looker persistent derived tables, Tableau extracts)

Симптомы

  • Бизнес-команда обнаруживает в дашборде клиента с возрастом 150 лет, отрицательную выручку, дубликаты заказов. Доверие к данным падает.

Причина

Нет DQ-тестов на ключевые поля (range, nullability, uniqueness) Тесты только на staging-слое, не на витринах Алерты не настроены, ошибки обнаруживаются вручную Источник присылает мусорные значения (-1, 9999-12-31, empty string)

Решение

  1. Добавить dbt-тесты: not_null, unique, accepted_values, relationships
  2. Внедрить Great Expectations или Soda для более сложных проверок
  3. Настроить Monte Carlo / Soda для observability с алертами в Slack
  4. Очищать мусорные значения на 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 перепутаны

Решение

  1. Кластеризовать dim и fact по join-ключу: ALTER TABLE ... CLUSTER BY (customer_id)
  2. Денормализовать часть dim в fact (Big Table approach в BigQuery)
  3. Использовать bucketing в Hive/Spark по customer_id с одинаковым числом бакетов
  4. Pre-aggregate в materialized view или snapshot, если запросы предсказуемы

Симптомы

  • В витрине fact_revenue колонка discount_amount, аналитик спрашивает «откуда она?». Никто не знает: модель собрана из 5 источников через 3 уровня трансформаций.

Причина

SQL-модели написаны через хардкод имён таблиц вместо ref()/source() Нет автоматизированного data catalog Документация модели пуста или устарела Несколько команд пишут трансформации в разных репозиториях

Решение

  1. Перевести трансформации в dbt с использованием ref() — автоматический lineage
  2. Внедрить OpenLineage + Marquez / DataHub для автоматического сбора метаданных
  3. Добавить description в YAML-конфиг dbt-моделей и колонок
  4. Регулярно ревьюить data catalog в командных митингах

Симптомы

  • DAG-и, которые должны были запуститься в 02:00, фактически стартуют в 02:45. UI показывает большое scheduler heartbeat delay. Sensor-задачи висят в очереди.

Причина

Слишком много активных DAG-ов на одном scheduler-е Скрипты DAG-ов делают тяжёлые операции на парсинге (db-запросы на верхнем уровне) Недостаточно worker slot-ов в пуле Слишком много poke-sensor-ов забивают слоты

Решение

  1. Включить HA-scheduler (Airflow 2+): несколько scheduler-экземпляров
  2. Уменьшить parsing-frequency: AIRFLOW__CORE__DAG_FILE_PROCESSOR_TIMEOUT=120
  3. Вынести тяжёлую логику из top-level DAG-кода в task-функции
  4. Перевести sensor-ы в reschedule mode

Симптомы

  • Lag в Kafka топике растёт от часа к часу: consumer не успевает за producer. Дашборд показывает 10 минут, потом 30 минут, потом 2 часа отставания.

Причина

Слишком мало consumer-инстансов относительно числа партиций Тяжёлая обработка на каждое сообщение (внешний API, медленная запись в DB) Producer резко увеличил throughput (например, после маркетинговой кампании) Партиций меньше, чем нужно для параллелизма

Решение

  1. Увеличить число partition в топике (но это влияет на ordering)
  2. Запустить больше consumer-инстансов в той же consumer group (до числа партиций)
  3. Batch-ить запись в DWH/DB вместо одиночных insert
  4. Вынести тяжёлую обработку в отдельный downstream-процесс через Kafka Streams

Симптомы

  • dbt run --select my_model завершается успешно, но в таблице нет последних данных. Source имеет новые строки, в инкрементальной модели — нет.

Причина

is_incremental() условие фильтрует по неправильному полю unique_key не уникален, MERGE не вставляет дубли Watermark подсчитывается через MAX(updated_at), но в источнике updated_at не обновляется Партиция уже была обработана при первом запуске

Решение

  1. Проверить SQL под is_incremental(): WHERE event_date >= (SELECT MAX(event_date) FROM {{ this }})
  2. Использовать корректный unique_key или композитный ключ
  3. Перейти на ingestion-timestamp от источника вместо source-side updated_at
  4. 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

Решение

  1. ALTER SYSTEM SET wal_level = logical, перезапустить Postgres
  2. ALTER TABLE my_table REPLICA IDENTITY FULL
  3. Мониторить размер replication slot: SELECT pg_size_pretty(pg_wal_lsn_diff(...))
  4. Увеличить 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-ит

Решение

  1. Открыть Spark UI -> Stages -> найти failed task -> посмотреть exception в exec logs
  2. Добавить try-catch в UDF и помечать bad rows вместо exception
  3. Увеличить spark.task.maxFailures с 4 до 8 для retry-able задач
  4. Если разные tasks падают с одной ошибкой — проблема в коде, не в инфре

Симптомы

  • В Slack ежедневно вопросы «где данные о клиентах?», «есть ли таблица с заказами по регионам?». Один и тот же датасет создаётся параллельно в 3 командах.

Причина

Нет data catalog — таблицы прячутся в произвольных schemas Naming convention отсутствует или не соблюдается Документация моделей не существует или устарела Каждая команда строит свои витрины в изоляции

Решение

  1. Внедрить data catalog: DataHub, OpenMetadata, Amundsen
  2. Зафиксировать naming convention: prefix_team_entity_grain (fact_sales_order_daily)
  3. В dbt — обязать description для каждой модели и колонки через CI-проверку
  4. Назначить data steward-ов от каждой команды

Симптомы

  • Финансовый дашборд показывает данные за позавчера, хотя пайплайн должен обновлять каждое утро. Бизнес не понимает, можно ли доверять цифрам.

Причина

Пайплайн упал ночью, но никто не отреагировал на alert Cache в BI-инструменте не инвалидируется Источник не обновился из-за проблем на их стороне DAG висит в queued из-за нехватки worker slot-ов

Решение

  1. Добавить freshness-метрики: dbt source freshness, table_last_modified_at в каталоге
  2. На дашборде показывать last_updated_at и подсветка если данные >24h old
  3. Настроить алерты на SLA через PagerDuty / Slack с дежурной ротацией
  4. Инвалидировать BI-кэш по событию завершения пайплайна, не по таймеру

Симптомы

  • В дашборде маркетинга revenue за апрель — 1.2M, в дашборде финансов — 1.18M. Бизнес требует объяснить разницу, обе стороны уверены в своих числах.

Причина

Разные источники истины: маркетинг тянет из CRM, финансы — из ERP Разная семантика метрики: gross revenue vs net revenue с возвратами Разные временные срезы (timezone, дата заказа vs дата отгрузки) Один из дашбордов использует stale snapshot

Решение

  1. Создать single source of truth: certified marts с подписью owner
  2. Согласовать определения метрик в semantic layer (Cube, dbt metrics, LookML)
  3. Документировать timezone и event-time семантику в каждой витрине
  4. При расхождениях — двигать дашборды на certified marts

Симптомы

  • Мигрировали data lake с CSV на Parquet, рассчитывая на ускорение. На практике — некоторые запросы стали медленнее.

Причина

Parquet файлы слишком мелкие (по нескольку KB) — небольшие batch-загрузки Запросы с SELECT * — не использует преимущества колоночного формата Predicate pushdown не работает из-за отсутствия min/max статистик Compression codec неоптимален (gzip vs snappy vs zstd)

Решение

  1. Настроить compaction: целиться в файлы 128-512 MB
  2. Использовать SELECT only_needed_columns вместо *
  3. Переключиться с gzip на snappy или zstd: лучше баланс CPU/IO
  4. Проверить статистики через 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 неверно подобран

Решение

  1. Явно кастить в нужный тип: WHERE date = CAST('2026-05-17' AS DATE)
  2. Избегать функций на partition column в WHERE
  3. Проверить EXPLAIN: должно показывать partition pruning
  4. В BigQuery — использовать _PARTITIONTIME или DATE-partition правильно

Симптомы

  • На локальной машине python pipeline.py отрабатывает за 5 минут. В Airflow тот же код падает с timeout / OOM / connection refused.

Причина

Локально — маленький семпл данных, в проде — миллионы строк Зашиты локальные пути или localhost connection strings В проде нет нужных Python-зависимостей Airflow worker не имеет доступа к нужным секретам/сетям

Решение

  1. Тестировать на репрезентативной выборке данных
  2. Всё конфигурировать через env variables и Airflow Connections
  3. Зафиксировать зависимости в requirements.txt, использовать virtualenv operator
  4. Проверить network policy / firewall между worker и target системами

Симптомы

  • Запустили backfill за 2 года на Snowflake, через 6 часов финдир спрашивает «что вы натворили». Месячный лимит warehouse-а сожжён за сутки.

Причина

Backfill запускается сразу за весь период вместо разбиения по партициям Используется warehouse слишком большого размера Параллельные runs всех партиций одновременно Нет ограничения по cost / runtime в job-е

Решение

  1. Разбить backfill на партиции по дням и запускать последовательно
  2. Использовать отдельный backfill-warehouse с меньшим размером
  3. Установить statement_timeout / resource monitor с hard cap
  4. Запускать backfill в выходные на off-peak ценах

Симптомы

  • SQL-запрос с CTE или window function на 100M строк висит часами или падает по timeout. EXPLAIN показывает огромный shuffle.

Причина

Window function без PARTITION BY — окно по всей таблице Recursive CTE с неконтролируемой глубиной Самoself-join по неиндексированному полю Большие промежуточные результаты в materialized CTE

Решение

  1. Добавить PARTITION BY user_id / region в window function
  2. Ограничить рекурсию через MAXRECURSION или явный exit
  3. Заменить self-join на join с pre-aggregated CTE
  4. Pre-aggregate в incremental model вместо on-the-fly расчёта

Симптомы

  • Старые consumer-ы Kafka-топика начинают падать после деплоя нового producer-а. Spark не может прочитать старые Parquet-файлы после изменения схемы.

Причина

Удалили required-поле — backward-incompatible change Переименовали колонку без alias Изменили тип поля (int -> string) Нет schema registry, контракт нигде не зафиксирован

Решение

  1. Использовать Confluent Schema Registry с backward/forward compatibility-режимом
  2. Добавлять новые поля как optional с default, не удалять старые сразу
  3. Использовать column mapping в Delta Lake / Iceberg для безопасного rename
  4. Тесты compatibility в CI: pre-merge проверка new schema vs old data