Меня зовут Дмитрий, я 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.Названия слоев вашего хранилища
Сразу зафиксирую, что не в каждой DWH есть именно RAW, STG, CORE и MARTS. На одном проекте RAW лежит в S3, на другом это таблицы в Postgres. STG и CORE иногда объединяют, а преобразования могут выполняться до загрузки в хранилище. То есть у всех по-разному, но смысл плюс минус один.
Здесь я использую одну из частых схем:
RAW хранит исходный файл и метаданные загрузки;
STG приводит данные к технически корректному виду;
CORE добавляет бизнес-правила;
MARTS собирает данные под конкретного потребителя;
BI читает готовую витрину.
Тут просто важно понимать, что произошло с полем на каждом переходе и где это можно проверить.
Что лежит в исходном CSV
Для примера я взял небольшой файл и такие строки:
order_id,customer_id,order_date,status,amount 1001,C001,2026-07-18,paid,1299.90 1002,C002,2026-07-18,paid,"1 299,90" 1003,C003,,paid,750.00 1004,C004,2026-07-19,created,500.00 1005,C005,2026-07-19,refunded,800.00 1006,C006,2026-07-20,paid,not_a_number 1007,C007,2026-07-20,paid,100.00 1007,C007,2026-07-20,paid,100.00 1008,C008,2026-07-20,cancelled,420.00 1009,C009,20.07.2026,paid,"300,50" 1010,C010,2026-07-21,paid,-50.00
Проблемы здесь специально собраны в одном месте:
два формата даты;
два формата суммы;
пустая дата;
текст вместо числа;
повторный
order_id;отрицательная сумма;
статусы, которые по-разному влияют на выручку.
В файле на несколько миллионов строк подход «открыть и посмотреть» уже не работает :)

orders.csv и проблемные строкиRAW: сохраняем то, что получили
Первый шаг DAG вычисляет SHA-256 содержимого файла. Имя файла для защиты от дублей не очень подходит, поскольку источник может прислать одинаковые данные под двумя именами или, наоборот, заменить содержимое файла с прежним именем.
file_bytes = source_file.read_bytes() file_sha256 = hashlib.sha256(file_bytes).hexdigest()
После этого код проверяет реестр загрузок:
select object_key from raw.habr_file_registry where 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, путь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=1003 7 | invalid_amount | order_id=1006 9 | duplicate_order_id_in_file | order_id=1007 12 | non_positive_amount | order_id=1010
Получили такой баланс:
11 строк источника = 7 принятых + 4 отклонённых
Все некорректные строки теперь видны с причинами брака.
Отдельно про дубль. DISTINCT убрал бы вторую одинаковую строку, но причина появления дубля потерялась бы. Здесь повторный order_id остаётся в rejects. Вообще его можно показать владельцу источника и решить, какую запись считать правильной, то есть оговаривается такое обычно отдельно.

stg.habr_orders_rejectsПосле STG у нас есть технически корректные строки: дата стала датой, сумма числом, обязательные поля заполнены, дубли отделены.
Однако, сумма заказа ещё не равна выручке. В примере действуют четыре правила:
case when order_status = 'paid' then amount when order_status = 'refunded' then -amount else 0 end as net_revenue
Получается так:
paid -> сумма входит в выручку refunded -> сумма вычитается created -> заказ ещё не оплачен, выручка равна нулю cancelled -> отменённый заказ не входит в выручку
Поэтому в CORE я храню и исходную сумму заказа, и рассчитанную выручку:
order_id | status | gross_amount | net_revenue 1001 | paid | 1299.90 | 1299.90 1004 | created | 500.00 | 0.00 1005 | refunded | 800.00 | -800.00 1008 | cancelled | 420.00 | 0.00
Тут приходит понимание, почему сумма из CSV и сумма на дашборде могут не совпадать. Это не обязательно ошибка. Иногда между ними находится бизнес-правило, которое нужно явно назвать и проверить. Чаще всего тут помогают аналитики с правильной бизнес-логикой (спасибо вам большое!).

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_revenue from core.habr_orders group by order_date;
Результат:
order_date | orders | gross_amount | net_revenue 2026-07-18 | 2 | 2599.80 | 2599.80 2026-07-19 | 2 | 1300.00 | -800.00 2026-07-20 | 3 | 820.50 | 400.50
Итого:
gross_amount = 4720.30 net_revenue = 2200.30
Разница объясняется созданным, отменённым и возвращённым заказами. Если оставить в витрине только одну колонку amount, через месяц уже будет сложно вспомнить, что именно она означает.
В реальном проекте здесь часто появляется ещё одна проблема: JOIN с таблицей позиций или платежей размножает строки заказа. Поэтому перед SUM нужно проверить гранулярность обеих таблиц и кардинальность соединения (какой именно JOIN использовать).

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 классически говорит, что скрипты технически выполнились. Но мы же знаем, что это не всегда равно тому, что все отбежало и теперь можно идти пить кофе.
Поэтому последняя задача запускает шесть проверок:
raw_row_balance 11 = 11 no_duplicate_order_ids_in_stg 0 = 0 no_null_required_fields_in_stg 0 = 0 stg_to_core_row_count 7 = 7 core_to_mart_net_revenue 2200.30 = 2200.30 sample_expected_net_revenue 2200.30 = 2200.30
Если хотя бы одна проверка не проходит, задача падает с перечнем нарушенных условий. В логах остаются actual, expected и короткое объяснение.
Для рабочего проекта последнюю проверку с жёстко заданной суммой я бы заменил сверкой с предыдущим слоем, историческим диапазоном или внешним контрольным источником. Здесь фиксированное значение удобно, поскольку датасет игрушечный, поэтому мы заранее знаем правильный ответ.

check_pipeline с шестью passed=TrueЧто произойдёт, если источник пришлёт тот же файл второй раз
После успешной загрузки я запускаю DAG ещё раз, не меняя исходный orders.csv.
Задача prepare_lab проходит, а store_file_in_raw находит SHA-256 в реестре и получает состояние skipped. Остальные задачи тоже пропускаются. В RAW не появляется вторая копия, а данные не загружаются повторно.
Имя файла при этом не участвует в решении. Проверяется содержимое.
Такая простая защита сработала. В production ещё нужно определить политику повторной обработки: можно ли переигрывать старые загрузки, что делать с исправленным файлом, как хранить версии и кто имеет право запускать перезапуск. Для примера важно, что повторный запуск не ломает нам цифры и выручку.

store_file_in_raw и последующие задачи отмечены как skipped.Что проверить в своём пайплайне
В примере было всего 11 строк, но на больших данных возникают те же вопросы:
совпало ли количество строк в источнике с суммой принятых и отклонённых;
есть ли понятная причина для каждого reject;
на каком этапе изменились типы и значения;
какие бизнес-правила повлияли на итоговую сумму;
не создаёт ли повторная загрузка дубли.
Повторюсь, что названия слоёв и набор инструментов могут отличаться (скорее всего) - в одном проекте будет MinIO и Airflow, в другом S3 и собственный оркестратор.
Важно, чтобы хороший пайп позволял быстро ответить, откуда взялась каждая строка, куда пропали остальные и почему цифра в BI отличается от исходного файла.
Интересно, как такие проверки устроены в ваших проектах. Куда вы складываете неликвид, как сверяете количество строк и деньги между слоями и что происходит при повторной загрузке файла? Расскажите в комментариях. Если тема окажется полезной, в следующей статье разберу инкрементальную загрузку или запуск Spark jobs через Airflow.
Если статья вам понравилась, заходите в мой блог https://t.me/kuzmin_dmitry91.
