DE. Путь файла по слоям

от автора

Меня зовут Дмитрий, я Data Engineer. В этой статье хочу на небольшом примере показать, что происходит с данными между исходным файлом и готовой BI-витриной и почему даже успешно завершившийся пайплайн не гарантирует правильный результат. Погнали!

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

На практике обычный CSV успевает создать достаточно проблем ещё до BI. В одной строке сумма записана как 1299.90, в другой как "1 299,90". У заказа нет даты, один order_id повторяется, а вчерашний файл источник присылает ещё раз. Итоговая сумма, даже если и появится, то может оказаться неверной.

Так вот, теперь к примеру — проведем один orders.csv по цепочке:

CSV -> RAW -> STG -> CORE -> MARTS -> BI

Для всего этого дела я использую несколько контейнеров Docker:

  • MinIO

  • Postgres

  • Airflow

В исходном файле 11 строк (тут я ограничился для простоты и визуальной легкости).

До CORE доходят 7, ещё 4 строки попадают в rejects с понятной причиной. Валовая сумма принятых заказов равна 4720.30, а выручка после применения бизнес-правил составляет 2200.30.

В этой статье покажу, где именно меняются данные и какие проверки контролируют сумму.

Общая схема CSV -> RAW -> STG -> CORE -> MARTS -> BI.

Общая схема CSV -> RAW -> STG -> CORE -> MARTS -> BI.

Названия слоев вашего хранилища

Сразу зафиксирую, что не в каждой DWH есть именно RAW, STG, CORE и MARTS. На одном проекте RAW лежит в S3, на другом это таблицы в Postgres. STG и CORE иногда объединяют, а преобразования могут выполняться до загрузки в хранилище. То есть у всех по-разному, но смысл плюс минус один.

Здесь я использую одну из частых схем:

  • RAW хранит исходный файл и метаданные загрузки;

  • STG приводит данные к технически корректному виду;

  • CORE добавляет бизнес-правила;

  • MARTS собирает данные под конкретного потребителя;

  • BI читает готовую витрину.

Тут просто важно понимать, что произошло с полем на каждом переходе и где это можно проверить.

Что лежит в исходном CSV

Для примера я взял небольшой файл и такие строки:

order_id,customer_id,order_date,status,amount1001,C001,2026-07-18,paid,1299.901002,C002,2026-07-18,paid,"1 299,90"1003,C003,,paid,750.001004,C004,2026-07-19,created,500.001005,C005,2026-07-19,refunded,800.001006,C006,2026-07-20,paid,not_a_number1007,C007,2026-07-20,paid,100.001007,C007,2026-07-20,paid,100.001008,C008,2026-07-20,cancelled,420.001009,C009,20.07.2026,paid,"300,50"1010,C010,2026-07-21,paid,-50.00

Проблемы здесь специально собраны в одном месте:

  • два формата даты;

  • два формата суммы;

  • пустая дата;

  • текст вместо числа;

  • повторный order_id;

  • отрицательная сумма;

  • статусы, которые по-разному влияют на выручку.

В файле на несколько миллионов строк подход «открыть и посмотреть» уже не работает 🙂

orders.csv и проблемные строки

orders.csv и проблемные строки

RAW: сохраняем то, что получили

Первый шаг DAG вычисляет SHA-256 содержимого файла. Имя файла для защиты от дублей не очень подходит, поскольку источник может прислать одинаковые данные под двумя именами или, наоборот, заменить содержимое файла с прежним именем.

file_bytes = source_file.read_bytes()file_sha256 = hashlib.sha256(file_bytes).hexdigest()

После этого код проверяет реестр загрузок:

select object_keyfrom raw.habr_file_registrywhere file_sha256 = %s;

Если такого хеша ещё нет, исходный CSV без изменений отправляется в MinIO.
В ключ объекта я тоже добавляю SHA-256:

raw/habr/orders/sha256=<hash>/orders.csv

В Postgres остаётся техническая запись:

file_name, file_sha256, object_key, source_row_count,valid_row_count, rejected_count, status, loaded_at

После первого запуска реестр выглядит так:

orders.csv | 11 | 7 | 4 | completed

Теперь можно ответить хотя бы на базовые вопросы: какой файл пришёл, когда его загрузили, сколько в нём было строк и где лежит оригинал.

Объект orders.csv в MinIO Console, бакет raw, путь

Объект orders.csv в MinIO Console, бакет raw, путь

STG: приводим данные в порядок

Следующая задача читает файл уже из RAW-бакета. Дальнейшая обработка опирается на сохранённый оригинал.

С датами всё относительно просто. В примере разрешены два формата:

def parse_date(value: str):    for pattern in ("%Y-%m-%d", "%d.%m.%Y"):        try:            return datetime.strptime(value.strip(), pattern).date()        except ValueError:            continue    raise ValueError("invalid_order_date")

С суммами немного интереснее. Пробел может быть разделителем тысяч, а запятая десятичным разделителем. После нормализации значение переводится в Decimal, а не в float:

def parse_amount(value: str) -> Decimal:    normalized = value.strip().replace(" ", "")    if "," in normalized:        normalized = normalized.replace(",", ".")    amount = Decimal(normalized).quantize(Decimal("0.01"))    if amount <= 0:        raise ValueError("non_positive_amount")    return amount

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

Ошибочную строку я не удаляю. Она записывается в stg.habr_orders_rejects вместе с номером строки, исходным JSON и причиной отказа.

По нашему файлу получилось четыре rejects:

4  | invalid_order_date         | order_id=10037  | invalid_amount             | order_id=10069  | duplicate_order_id_in_file | order_id=100712 | non_positive_amount        | order_id=1010

Получили такой баланс:

11 строк источника = 7 принятых + 4 отклонённых

Все некорректные строки теперь видны с причинами брака.

Отдельно про дубль. DISTINCT убрал бы вторую одинаковую строку, но причина появления дубля потерялась бы. Здесь повторный order_id остаётся в rejects. Вообще его можно показать владельцу источника и решить, какую запись считать правильной, то есть оговаривается такое обычно отдельно.

Результат запроса к stg.habr_orders_rejects

Результат запроса к stg.habr_orders_rejects

После STG у нас есть технически корректные строки: дата стала датой, сумма числом, обязательные поля заполнены, дубли отделены.

Однако, сумма заказа ещё не равна выручке. В примере действуют четыре правила:

case    when order_status = 'paid' then amount    when order_status = 'refunded' then -amount    else 0end as net_revenue

Получается так:

paid      -> сумма входит в выручкуrefunded  -> сумма вычитаетсяcreated   -> заказ ещё не оплачен, выручка равна нулюcancelled -> отменённый заказ не входит в выручку

Поэтому в CORE я храню и исходную сумму заказа, и рассчитанную выручку:

order_id | status    | gross_amount | net_revenue1001     | paid      | 1299.90      | 1299.901004     | created   | 500.00       | 0.001005     | refunded  | 800.00       | -800.001008     | cancelled | 420.00       | 0.00

Тут приходит понимание, почему сумма из CSV и сумма на дашборде могут не совпадать. Это не обязательно ошибка. Иногда между ними находится бизнес-правило, которое нужно явно назвать и проверить. Чаще всего тут помогают аналитики с правильной бизнес-логикой (спасибо вам большое!).

Таблица core.habr_orders: gross_amount, net_revenue и business_rule.

Таблица core.habr_orders: gross_amount, net_revenue и business_rule.

MARTS: выбираем гранулярность и считаем дальше

Витрина отвечает на конкретный вопрос: какая выручка была по дням. Её гранулярность можно сформулировать одной фразой:

Одна строка равна одному календарному дню заказа.

После этого агрегация читается нормально:

select    order_date,    count(*) as orders_count,    count(*) filter (where order_status = 'paid') as paid_orders,    count(*) filter (where order_status = 'refunded') as refunded_orders,    sum(gross_amount) as gross_amount,    sum(net_revenue) as net_revenuefrom core.habr_ordersgroup by order_date;

Результат:

order_date  | orders | gross_amount | net_revenue2026-07-18  | 2      | 2599.80      | 2599.802026-07-19  | 2      | 1300.00      | -800.002026-07-20  | 3      | 820.50       | 400.50

Итого:

gross_amount = 4720.30net_revenue  = 2200.30

Разница объясняется созданным, отменённым и возвращённым заказами. Если оставить в витрине только одну колонку amount, через месяц уже будет сложно вспомнить, что именно она означает.

В реальном проекте здесь часто появляется ещё одна проблема: JOIN с таблицей позиций или платежей размножает строки заказа. Поэтому перед SUM нужно проверить гранулярность обеих таблиц и кардинальность соединения (какой именно JOIN использовать).

График в Metabase

График в Metabase

Airflow: порядок шагов и место ошибки

DAG состоит из шести задач:

prepare_lab  -> store_file_in_raw  -> clean_to_stg  -> apply_business_rules_to_core  -> build_daily_mart  -> check_pipeline

Airflow просто запускает шаги в нужном порядке, передаёт метаданные между задачами и останавливает цепочку при ошибке.

Такой граф полезен ещё и для диагностики:

  • Если файл не попал в MinIO, смотрим store_file_in_raw.

  • Если четыре строки ушли в rejects, открываем clean_to_stg.

  • Если CORE собрался, а сумма витрины разошлась, проблема находится между apply_business_rules_to_core, build_daily_mart и проверками.

Успешный DAG

Успешный DAG

Проверки после успешного DAG

Зелёный DAG классически говорит, что скрипты технически выполнились. Но мы же знаем, что это не всегда равно тому, что все отбежало и теперь можно идти пить кофе.

Поэтому последняя задача запускает шесть проверок:

raw_row_balance                    11 = 11no_duplicate_order_ids_in_stg       0 = 0no_null_required_fields_in_stg      0 = 0stg_to_core_row_count               7 = 7core_to_mart_net_revenue      2200.30 = 2200.30sample_expected_net_revenue   2200.30 = 2200.30

Если хотя бы одна проверка не проходит, задача падает с перечнем нарушенных условий. В логах остаются actual, expected и короткое объяснение.

Для рабочего проекта последнюю проверку с жёстко заданной суммой я бы заменил сверкой с предыдущим слоем, историческим диапазоном или внешним контрольным источником. Здесь фиксированное значение удобно, поскольку датасет игрушечный, поэтому мы заранее знаем правильный ответ.

Лог задачи check_pipeline с шестью passed=True

Лог задачи check_pipeline с шестью passed=True

Что произойдёт, если источник пришлёт тот же файл второй раз

После успешной загрузки я запускаю DAG ещё раз, не меняя исходный orders.csv.

Задача prepare_lab проходит, а store_file_in_raw находит SHA-256 в реестре и получает состояние skipped. Остальные задачи тоже пропускаются. В RAW не появляется вторая копия, а данные не загружаются повторно.

Имя файла при этом не участвует в решении. Проверяется содержимое.

Такая простая защита сработала. В production ещё нужно определить политику повторной обработки: можно ли переигрывать старые загрузки, что делать с исправленным файлом, как хранить версии и кто имеет право запускать перезапуск. Для примера важно, что повторный запуск не ломает нам цифры и выручку.

Второй DAG-run: store_file_in_raw и последующие задачи отмечены как skipped.

Второй DAG-run: store_file_in_raw и последующие задачи отмечены как skipped.

Что проверить в своём пайплайне

В примере было всего 11 строк, но на больших данных возникают те же вопросы:

  • совпало ли количество строк в источнике с суммой принятых и отклонённых;

  • есть ли понятная причина для каждого reject;

  • на каком этапе изменились типы и значения;

  • какие бизнес-правила повлияли на итоговую сумму;

  • не создаёт ли повторная загрузка дубли.

Повторюсь, что названия слоёв и набор инструментов могут отличаться (скорее всего) — в одном проекте будет MinIO и Airflow, в другом S3 и собственный оркестратор.

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

Интересно, как такие проверки устроены в ваших проектах. Куда вы складываете неликвид, как сверяете количество строк и деньги между слоями и что происходит при повторной загрузке файла? Расскажите в комментариях. Если тема окажется полезной, в следующей статье разберу инкрементальную загрузку или запуск Spark jobs через Airflow.

Если статья вам понравилась, заходите в мой блог https://t.me/kuzmin_dmitry91.

ссылка на оригинал статьи https://habr.com/ru/articles/1062446/