Data Engineering и ETL Data Engineering: карта трека, роли и архитектура данных
0%

Data Engineering: карта трека, роли и архитектура данных

Data Engineering: карта трека, роли и архитектура данных

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

1. Зачем вообще нужен отдельный слой инженерии данных

Начнём с ситуации, знакомой почти всем. Есть продуктовая база (PostgreSQL, MySQL — неважно). Аналитику нужно число «выручка за вчера по странам». Кажется, что решение занимает одну строку: подключиться к реплике и написать SELECT. Первые полгода так и работает. Потом ломается, и ломается всегда одинаково — по пяти независимым причинам.

Первая: OLTP-схема оптимизирована не под тот вопрос. Продуктовая база нормализована, чтобы быстро писать и точечно читать строку по ключу. Аналитический запрос читает 200 миллионов строк и агрегирует три колонки из тридцати, а строковое хранение заставляет поднять с диска все тридцать. Разница в объёме чтения — порядок и больше (https://courses.digitable.life/post/data-engineering/06-storage-and-formats/).

Вторая: аналитика мешает продукту. Тяжёлый запрос на реплике держит снапшот, раздувает WAL, конкурирует за буферный кэш. Классическая авария: «отчёт положил чекаут».

Третья: OLTP не хранит историю. В таблице users поле plan перезаписывается при апгрейде, и вопрос «сколько было пользователей на Pro первого марта» не имеет ответа — прошлое состояние физически стёрто. Историю фиксируют заранее, медленно меняющимися измерениями (https://courses.digitable.life/post/data-engineering/02-data-modeling/).

Четвёртая: одного источника не хватает. Метрика «выручка по каналу привлечения» живёт в трёх системах: заказы в продуктовой базе, платежи в биллинге, атрибуция в рекламных кабинетах — с разными ключами, часовыми поясами и валютами.

Пятая, самая дорогая: определения расходятся. Маркетинг считает активного пользователя как «зашёл за 30 дней», продукт — «совершил действие за 7 дней», финансы — «оплатил». Три цифры на трёх дашбордах, и совещание превращается в спор о данных.

Data engineering закрывает все пять проблем одним способом: строит отдельный, воспроизводимый, версионируемый конвейер от источников к решениям. У Райса и Хаусли в «Fundamentals of Data Engineering» (O’Reilly, 2022) это звучит так: дата-инженер получает сырые данные и производит из них качественный консистентный продукт, готовый к потреблению.

Аналогия, которая хорошо работает: водопровод. Вода в реке есть всегда, но пить её нельзя. Между рекой и краном стоят водозабор, отстойники, фильтры, лаборатория контроля, насосы и трубы под давлением. Никто из жителей не думает про фильтры — они думают про «открыл кран, потекла чистая вода, напор нормальный». Ровно это и есть контракт платформы данных: свежесть (напор), корректность (чистота), доступность (кран работает).

2. Карта трека

Порядок статей не случаен: сначала способ движения (https://courses.digitable.life/post/data-engineering/01-etl-vs-elt/), потом форма (https://courses.digitable.life/post/data-engineering/02-data-modeling/), потом два режима исполнения — пакетный (https://courses.digitable.life/post/data-engineering/03-batch-processing/) и потоковый (https://courses.digitable.life/post/data-engineering/04-streaming/), затем то, что связывает всё в работающую систему (https://courses.digitable.life/post/data-engineering/05-orchestration/), физический субстрат (https://courses.digitable.life/post/data-engineering/06-storage-and-formats/), гарантии (https://courses.digitable.life/post/data-engineering/07-data-quality-and-governance/) и, наконец, потребители (https://courses.digitable.life/post/data-engineering/08-analytics-and-ml-handoff/).

3. Слои платформы: почему их именно столько

Любая зрелая платформа приходит к одной и той же слоёной структуре. У Databricks она называется «медальонной» (bronze/silver/gold), в dbt-мире — staging / intermediate / marts, у Кимбалла — ODS / DWH / data marts. Названия разные, логика одна.

Слои платформы данных: raw, staging, core, marts и как меняются свойства данных

Ключевой принцип: каждый слой снимает ровно один вид сложности, и никогда не берёт на себя чужую работу.

Слой Что снимает Что запрещено Идемпотентность
RAW недоступность источника любые трансформации append-only, партиция иммутабельна
STAGING грязь формата: типы, дубли, имена join между источниками полная перезапись партиции
CORE расхождение бизнес-смысла «удобные» костыли под один дашборд merge по ключу + история
MARTS дороговизна запроса обращение к RAW напрямую пересчёт из CORE в любой момент

Почему нельзя схлопнуть слои и писать сразу из RAW в витрину? Тогда дедупликация, приведение типов и определение бизнес-сущности размазываются по десяткам запросов, и через год у вас двенадцать реализаций «активного пользователя» без способа узнать, какая верна. Слои — не бюрократия, а вынесение общего подвыражения: тот же рефакторинг, что и извлечение функции в коде. Второй принцип, который часто нарушают: RAW не переписывается никогда — именно нетронутый RAW позволяет пересчитать всё сверху, когда обнаружится ошибка в трансформации. Это аналог git-истории: пока сырые данные целы, любая ошибка обратима; как только вы «почистили мусор» в RAW, ошибка трансформации становится потерей данных.

4. Референсная архитектура

Оркестратор, качество и каталог — сквозные, а не этап. Начинающие команды рисуют их как ещё один прямоугольник в цепочке и потом удивляются, что качество проверяется «в конце», когда плохие данные уже разошлись по десяти витринам. Правильная модель: это три перпендикулярные плоскости, пересекающие все слои.

Batch и stream — не альтернатива, а два режима одной модели. Исторически их разделяли в лямбда-архитектуре (два кода, два результата, вечное расхождение). Современный подход — каппа-архитектура и унифицированные модели вроде Dataflow Model (Akidau et al., VLDB 2015), где пакет — это частный случай потока с бесконечно большим окном. Про это подробно в https://courses.digitable.life/post/data-engineering/04-streaming/.

5. Жизненный цикл партиции данных

Единица работы в пакетном мире — не «таблица», а партиция: например, все события за dt=2026-07-15. Понимание её жизненного цикла объясняет половину решений в треке.

Обратите внимание на переход Опубликована --> Загружается. Он означает, что уже опубликованная партиция может быть пересчитана, и это нормальная, ожидаемая операция, а не авария. Отсюда прямо следует главное техническое требование всего трека:

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

Почему опубликованные данные приходится пересчитывать — видно на картинке про два времени.

Время события против времени обработки, водяной знак и опоздавшие события

Событие произошло в 23:58 у пользователя в офлайне, а долетело в 09:15 следующего дня: по event time оно относится к вчерашней партиции, по processing time — к сегодняшней. Вы либо ждёте (теряете свежесть), либо публикуете и потом пересчитываете (теряете стабильность цифр), либо теряете событие. Третьего варианта нет — это фундаментальный компромисс, а не недоработка инструмента.

6. Идемпотентность на практике

Наивный вариант загрузки:

# ПЛОХО: не идемпотентно. Второй запуск удвоит строки.
def load_day(conn, rows: list[dict]) -> None:
    conn.executemany(
        "INSERT INTO orders (id, user_id, amount, dt) VALUES (?, ?, ?, ?)",
        [(r["id"], r["user_id"], r["amount"], r["dt"]) for r in rows],
    )

Retry после сетевой ошибки, ручной перезапуск оператором, повторная доставка из очереди — любой из этих сценариев даёт задвоение выручки. Причём заметят это через месяц, когда финансы не сойдутся.

Корректный вариант строится на перезаписи целой партиции — простейшая форма идемпотентности, работающая почти везде:

from datetime import date
import pyarrow as pa
import pyarrow.parquet as pq

def partition_path(table: str, dt: date) -> str:
    # Hive-совместимая раскладка: движки умеют отсекать партиции прямо по пути
    return f"s3://lake/raw/{table}/dt={dt.isoformat()}"

def write_partition_atomically(table: str, dt: date, data: pa.Table, fs) -> None:
    """Идемпотентная запись: пишем во временный префикс, затем атомарно
    подменяем целевой. Повторный запуск даёт ровно тот же результат."""
    target, tmp = partition_path(table, dt), partition_path(table, dt) + "__tmp"

    if fs.exists(tmp):                  # хвост от упавшего прошлого запуска
        fs.rm(tmp, recursive=True)

    pq.write_to_dataset(
        data, root_path=tmp, filesystem=fs,
        compression="zstd",             # ~как gzip по размеру, заметно быстрее на чтении
        row_group_size=1_000_000,       # крупные row group → эффективное отсечение по min/max
    )

    # Публикация. На объектных хранилищах без атомарного rename эту роль
    # выполняет коммит табличного формата (Iceberg / Delta): новая версия
    # метаданных становится видимой одним атомарным обновлением указателя.
    if fs.exists(target):
        fs.rm(target, recursive=True)
    fs.mv(tmp, target, recursive=True)

Ключевая идея: читатель никогда не видит полузаписанную партицию. Пока идёт запись, данные лежат под временным префиксом; переключение — одна атомарная операция. Ровно этот приём — сердце форматов Iceberg и Delta Lake, где «переключение указателя» реализовано как коммит новой версии метаданных (спецификация Iceberg).

Второй способ идемпотентности — merge по ключу с дедупликацией, нужен там, где партиция не перезаписывается целиком (медленно меняющиеся сущности, CDC-поток):

-- Слой staging: снимаем дубликаты, оставляя последнюю версию каждой записи.
-- Дубликаты в CDC-потоке — норма: at-least-once доставка это гарантирует.
create or replace table stg_orders as
with ranked as (
    select *,
           row_number() over (partition by order_id
                              order by updated_at desc, ingested_at desc) as rn
    from raw_orders
    where dt = date '2026-07-15'
)
select order_id,
       user_id,
       cast(amount_cents as bigint) / 100.0 as amount,
       cast(created_at as timestamp)        as created_at,
       updated_at
from ranked
where rn = 1;

-- Слой core: MERGE, а не INSERT. Повторный запуск не создаёт дублей.
merge into core_orders t
using stg_orders s
    on t.order_id = s.order_id
when matched and s.updated_at > t.updated_at then update set *
when not matched then insert *;

Обратите внимание на order by updated_at desc, ingested_at desc: второй ключ нужен, потому что updated_at у источника может совпадать до секунды у двух версий строки. Без tie-breaker результат запроса недетерминирован — при пересчёте получите другую строку, и цифры «поедут» без единой ошибки в коде. Это одна из самых коварных ошибок в ETL: сортировка без полного порядка.

7. Сложность и экономика: почему инкремент, а не full refresh

Пусть таблица накопила N строк, а за сутки приходит Δ новых, где обычно Δ ≪ N.

Стратегия Время Чтение с диска Когда оправдана
Full refresh O(N) за запуск, O(N·D) за D дней весь датасет каждый раз N мал (до десятков млн), логика часто меняется
Инкремент по партиции O(Δ) одна партиция основной режим для фактов
Merge по ключу O(Δ · log N) или O(Δ + N) при hash-join партиция + индекс/сортировка целевой измерения, CDC
Backfill за период O(Δ · K) для K дней K партиций после исправления логики

Практический вывод, который экономит реальные деньги: стоимость пайплайна определяется не сложностью SQL, а объёмом просканированных байт. В облачных хранилищах вы платите буквально за сканирование (BigQuery on-demand тарифицирует TB прочитанных данных). Поэтому три приёма дают почти весь выигрыш:

  1. Партиционирование по времени — движок читает только нужные каталоги (partition pruning).
  2. Колоночный формат — читаются только запрошенные колонки, обычно 3 из 40.
  3. Статистика по row group — min/max в метаданных Parquet позволяют пропустить блок целиком.

Совокупно это легко даёт сокращение чтения в 50–200 раз против «прочитать всю таблицу построчно» (https://courses.digitable.life/post/data-engineering/06-storage-and-formats/). Контр-нюанс, о котором забывают: у партиционирования есть цена. Слишком мелкие партиции порождают проблему мелких файлов — сотни тысяч объектов, каждый со своими метаданными, и планировщик тратит больше времени на листинг, чем на чтение. Практическое правило: целевой размер файла 128 МБ – 1 ГБ; если партиция стабильно меньше, укрупняйте гранулярность или запускайте компакцию.

8. Роли: кто за что отвечает

Самый частый организационный сбой в работе с данными — не техника, а неопределённая ответственность на границах. Разберём типовое распределение.

Data engineer. Отвечает за движение и надёжность: приём, оркестрацию, инфраструктуру хранения, SLA по свежести, идемпотентность, стоимость. Владеет RAW и STAGING. Метрика качества его работы — «данные на месте, вовремя, полные, за предсказуемые деньги».

Analytics engineer. Роль, оформившаяся около 2019 года вместе с dbt (см. The Analytics Engineer). Отвечает за смысл: моделирование CORE, единые определения метрик, тесты бизнес-логики, документацию. Пишет в основном SQL, но применяет к нему инженерные практики — версионирование, код-ревью, CI, тесты. Владеет CORE и MARTS.

Аналитик / data scientist. Отвечает за вопрос и вывод: гипотеза, интерпретация, модель. Потребитель, а не владелец слоёв. ML engineer отвечает за то, чтобы фичи в обучении и в проде считались одинаково; его главная боль — training/serving skew — лежит на стыке с платформой (https://courses.digitable.life/post/data-engineering/08-analytics-and-ml-handoff/).

Data steward / владелец домена. Отвечает за смысл и допустимость: что означает поле, кто имеет право видеть, сколько хранить. Не техническая роль, но без неё governance не существует (https://courses.digitable.life/post/data-engineering/07-data-quality-and-governance/).

Отдельно — производитель данных. Это продуктовая команда, чей код порождает события. Ключевой сдвиг мышления, который принесла идея data mesh (Zhamak Dehghani, martinfowler.com): схема события — это публичный API, а не деталь реализации. Переименовать поле в таблице orders — то же самое, что сломать обратную совместимость REST-эндпоинта. Пока эта мысль не принята организационно, дата-инженеры обречены на бесконечную починку пайплайнов после чужих релизов. Отсюда общая эвристика распределения: ответственность идёт за возможностью починить — если сломанное поле может починить только продуктовая команда, значит SLA на схему её, и никакая героическая валидация на входе платформы этого не заменит.

9. Обещания платформы: SLA, SLO и сервисные уровни данных

Платформа без явных обещаний не отличается от отсутствия платформы: пользователи всё равно проверяют цифры вручную. Три измерения обещаний:

  • Свежесть (freshness). «Партиция за сутки доступна к 06:00 в 99% случаев».
  • Полнота (completeness). «Не менее 99.9% событий источника доехали до RAW за 24 часа».
  • Корректность (correctness). «Ключ уникален, суммы неотрицательны, ссылочная целостность держится».

Формализуются они как тесты, исполняемые на каждой публикации:

# dbt: тесты как часть определения модели — контракт живёт рядом с кодом
version: 2
models:
  - name: core_orders
    description: "Заказы, одна строка на order_id. Источник истины по выручке."
    config:
      contract: {enforced: true}     # схема фиксируется: изменение типа сломает сборку
    columns:
      - name: order_id
        data_type: string
        constraints: [{type: not_null}, {type: primary_key}]
        tests: [unique, not_null]
      - name: amount
        tests: [{dbt_utils.accepted_range: {min_value: 0, inclusive: true}}]
      - name: user_id
        tests: [{relationships: {to: ref('core_users'), field: user_id}}]
    tests:
      # SLA свежести как исполняемая проверка: данные не старше 26 часов
      - dbt_utils.recency: {datepart: hour, field: created_at, interval: 26}

Важный принцип: тест, который никого не будит, — это не тест, а логирование. Каждой проверке нужен уровень: error останавливает публикацию, warn создаёт задачу. Если всё выставлено в warn, через два месяца в канале 300 непрочитанных предупреждений и ноль доверия. Второй принцип: проверка стоит до публикации, а не после. Схема «записали в витрину → проверили → откатили» почти всегда означает, что кто-то уже успел увидеть неправильную цифру и принять по ней решение. Правильный порядок — write-audit-publish: записать в невидимую версию, проверить, атомарно опубликовать.

10. Типичные ошибки

Грабли из практики, каждая из которых стоит команде недель.

1. Пайплайн без идемпотентности. Проявляется как задвоение после ретрая. Проверка: можете ли вы прямо сейчас запустить вчерашнюю задачу второй раз и не бояться? Если нет — это долг номер один.

2. Трансформации в RAW. «Мы же всё равно фильтруем ботов на входе». Через полгода определение бота меняется, а исторические данные восстановить нечем.

3. Бизнес-логика в BI-инструменте. Метрика посчитана формулой внутри дашборда: невидима для тестов, не переиспользуется, дублируется на следующем дашборде с опечаткой. Логика должна жить в моделируемом слое под версионным контролем.

4. Оркестрация по расписанию вместо зависимостей. «Загрузка в 02:00, трансформация в 03:00 — за час точно успеет». В день, когда источник отдаёт данные медленнее, трансформация считает пустоту и публикует нули. Правильно — явные зависимости и сенсоры готовности данных (https://courses.digitable.life/post/data-engineering/05-orchestration/).

5. SELECT * в staging. Каждое новое поле источника протекает во все слои, включая PII, о которой никто не знал. Явный список колонок — это ещё и граница ответственности за приватность. Рядом стоит одна большая задача вместо графа: скрипт на 900 строк падает на 800-й, и перезапускать приходится с нуля — гранулярность задачи должна совпадать с гранулярностью восстановления.

6. Мониторинг инфраструктуры вместо мониторинга данных. Все зелёные галочки: DAG отработал, CPU в норме, ошибок ноль. При этом источник вчера начал присылать NULL в поле страны, и вся география поехала. Наблюдаемость данных — про значения, а не про процессы. И, наконец, данные без владельца: таблица есть, ей пользуются, автор уволился — каталог и явное владение суть условие того, что через два года систему можно будет менять.

11. Как это выглядит в проде

Стартап, 1–2 инженера, десятки ГБ. Managed-СУБД (Postgres или ClickHouse), загрузка готовыми коннекторами (Airbyte/Fivetran), трансформации на dbt, оркестрация — cron, BI поверх. Никакого Spark, никакой Kafka. Главное решение здесь — не строить платформу раньше, чем появились три источника и повторяющиеся вопросы.

Средняя компания, 5–15 инженеров, десятки ТБ. Object storage + открытый табличный формат (Iceberg/Delta), Spark или облачный движок для тяжёлого батча, Kafka для событий, Airflow или Dagster, dbt для CORE/MARTS, каталог с линиджем, дежурство по инцидентам данных. Появляется явный метрический слой и процесс изменения схем.

Крупная компания, десятки команд, петабайты. Federated-модель: платформенная команда даёт самообслуживание (шаблоны пайплайнов, стандарты контрактов, квоты), доменные команды владеют своими данными как продуктами. Это и есть практическая суть data mesh — не технология, а модель ответственности. Обязательны контроль стоимости по командам, автоматический линидж, политики retention и PII. Что общего у всех трёх масштабов: сначала вопрос бизнеса, потом данные, потом инструмент. Обратный порядок («внедряем Kafka, найдём применение») — надёжный способ получить дорогую неработающую платформу.

12. Сквозные понятия, которые встретятся в каждой статье

  • Партиция — физически отделяемый кусок таблицы, обычно по дате; единица работы, восстановления и удаления. Идемпотентность — повторный запуск даёт тот же результат. Backfill — пересчёт исторического диапазона.
  • Event time / processing time — время события у источника и время его обработки. Водяной знак (watermark) — граница, после которой считаем окно закрытым.
  • Схема-эволюция — совместимые изменения структуры (добавление поля — совместимо, переименование — нет). Data contract — соглашение производителя и потребителя о схеме, семантике и SLA.
  • Линидж — граф происхождения: из каких колонок каких таблиц получена данная колонка. SCD — способ хранить историю изменений атрибутов сущности.
  • Exactly-once / at-least-once — семантики доставки; на практике exactly-once достигается как at-least-once + идемпотентный приёмник.

13. Как пользоваться треком

Если вы приходите с задачей, а не с целью «прочитать всё»:

  • «Забрать данные из пяти систем в одно место» → https://courses.digitable.life/post/data-engineering/01-etl-vs-elt/, затем https://courses.digitable.life/post/data-engineering/05-orchestration/.
  • «Цифры на дашбордах не сходятся» → https://courses.digitable.life/post/data-engineering/02-data-modeling/ и https://courses.digitable.life/post/data-engineering/07-data-quality-and-governance/.
  • «Пайплайн не укладывается в окно или стоит слишком дорого» → https://courses.digitable.life/post/data-engineering/03-batch-processing/ и https://courses.digitable.life/post/data-engineering/06-storage-and-formats/.
  • «Нужны данные в реальном времени» → https://courses.digitable.life/post/data-engineering/04-streaming/ (и честно проверьте, действительно ли нужны: минутная задержка обычно решает задачу за долю цены).
  • «Модель в проде работает хуже, чем на валидации» → https://courses.digitable.life/post/data-engineering/08-analytics-and-ml-handoff/.

Смежные треки портала: https://courses.digitable.life/post/architecture-patterns/00-overview/ даёт язык для обсуждения границ и trade-offs, https://courses.digitable.life/post/ddd/00-overview/ — для разговора о доменах и едином языке (прямо применимо к моделированию CORE), https://courses.digitable.life/post/machine-learning/00-overview/ — про потребителя ваших витрин.

14. Мини-итог

  • Инженерия данных существует, потому что OLTP-база не отвечает на аналитические вопросы: другой профиль чтения, нет истории, нет сведения источников, нет единых определений.
  • Платформа расслаивается на RAW / STAGING / CORE / MARTS: каждый слой снимает ровно один вид сложности, RAW неприкосновенен, оркестрация и качество — сквозные плоскости.
  • Единица работы — партиция; её жизненный цикл включает пересчёт как штатный переход, поэтому идемпотентность обязательна везде.
  • Стоимость определяется просканированными байтами: партиционирование, колоночный формат и статистика блоков дают почти весь выигрыш.
  • Роли различаются по владению слоями; самый частый организационный сбой — производитель данных не знает, что его схема является публичным API.
  • Обещания платформы измеримы по трём осям: свежесть, полнота, корректность — и должны быть исполняемыми тестами до публикации.

Источники

Что дальше

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

ETL и ELT: этапы, различия и практика

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов