Меня зовут Дмитрий, я 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,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 и проблемные строки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, путь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 у нас есть технически корректные строки: дата стала датой, сумма числом, обязательные поля заполнены, дубли отделены.
Однако, сумма заказа ещё не равна выручке. В примере действуют четыре правила:
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.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 использовать).
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 = 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Что произойдёт, если источник пришлёт тот же файл второй раз
После успешной загрузки я запускаю 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.
ссылка на оригинал статьи https://habr.com/ru/articles/1062446/