Write–Audit–Publish (WAP): как спроектировать «чистые» пайплайны данных
Паттерн Write-Audit-Publish (WAP) — проверенная практика, которая «развязывает» запись данных и их публикацию, добавляя обязательный рубеж контроля качества между ними. В современном стеке данных WAP позволяет снижать операционные риски, предотвращать протечку «грязи» в продуктив и добиться трассируемости изменений. Коротко и очень понятно паттерн описан в блоге Dagster, а впервые в индустрии его популяризировала команда Netflix в 2017 году (доклад Michelle Ufford «Whoops, the Numbers are wrong! Scaling Data Quality @ Netflix»).
Что такое WAP простыми словами
Write — записываем результаты этапа в изолированную зону (staging/ветка/схема). Эти данные ещё не видят потребители (BI, ML, API).
Audit — запускаем проверки качества: схемы, уникальность, not null, ссылочная целостность, допустимые значения, а также бизнес-правила и статистические тесты.
Publish — только если все проверки пройдены, «переключаем указатель» на новую версию или быстро «перекладываем» проверенные данные в прод. Иначе — откатываемся/блокируем публикацию и уведомляем.
Ключевая идея — публикация невозможна без успешного аудита. Именно это отделение делает пайплайн устойчивым к ошибкам и безопасным для потребителей.
Архитектурные варианты WAP
Ветвление на уровне хранилища/табличного формата
Современные «табличные» форматы для data lake (Apache Iceberg, Delta Lake, Hudi) умеют версии снимков (snapshots), а Iceberg — ещё и ветки/теги. Обычно делают так:
- Включаем WAP-режим и создаём ветку для проверки:
- -- (Spark SQL с Iceberg)
- ALTER TABLE db.table SET TBLPROPERTIES ('write.wap.enabled'='true');
- ALTER TABLE db.table CREATE BRANCH audit_branch RETAIN 7 DAYS;
- Пишем данные в ветку и выполняем проверки.
- Если всё чисто — fast-forward основного указателя main на голову ветки:
- CALL catalog.system.fast_forward('db.table', 'main', 'audit_branch');
- Ветка затем автоматически протухает по политике retention.
Где хранить «черновик»:
- Iceberg branches (нативно, через Spark/Trino/Flink).
- lakeFS (S3/HDFS-надстройка с Git-подобными ветками).
- Project Nessie (транзакционный каталог с Git-семантикой поверх lakehouse).
- Bauplan демонстрирует WAP на Nessie «из коробки» (Python-подход).
Плюсы: атомарность публикации, быстрые откаты, единый источник правды. Минусы: нужна дисциплина по «уборке» веток, понимание механики snapshot-жизненного цикла и доп. стоимость хранения.
Blue/Green в СУБД/SQL-warehouse (dbt)
Если вы в Snowflake/BigQuery/Databricks, удобна «синяя/зелёная» публикация:
- Build моделей в «синюю» схему/базу,
- Test (dbt tests/пакеты),
- Swap/Expose — переключить представления/алиасы/clone на «зелёную» сторону. Для dbt есть официальный рецепт на базе dbt clone.
Префиксы/«теневые» бакеты в объектном хранилище
Запись — в s3://raw/_staging/run_2025_08_22/…, аудит — по этим ключам, публикация — атомарным move/rename (или сменой указателя таблицы/вью). Работает, но атомарности на файловом уровне нет без формата с транзакционным метаданным журналом (Iceberg/Delta/Hudi).
Streaming-WAP
В потоках (Kafka/Flink/Spark Structured Streaming) применяют «теневые» топики/таблицы: пишем события в topic_ingest_wap, валидируем агрегатами/правилами, затем «перекладываем» в topic_ingest_main или делаем fast-forward на ветку Iceberg для таблиц sink’а. Подход и кейсы описывает Dagster.
Где брать тесты для Audit
- dbt tests: unique, not_null, accepted_values, relationships + пользовательские и пакеты типа dbt-expectations. Это быстро даёт базовую «сетку безопасности».
- Great Expectations (GX): декларативные «expectations», интеграции с оркестраторами. (Проект ребрендился в GX; остаётся де-факто стандартом в OSS.)
- AWS Glue Data Quality: удобно сочетать с Iceberg-ветками на AWS (пример от AWS).
- Databricks Delta Live Tables (DLT): декларативные правила качества (expectations) в самом пайплайне.
Связка с оркестраторами и CI/CD
- Dagster, Prefect, Airflow: оформляйте гейты: publish зависит от audit. Набор задач: write_raw → dq_checks → publish/swap. Dagster отдельно описывает WAP-workflow и стриминги.
- CI: прогон unit/интеграционных тестов SQL/pytest, линтеров схем и генерации миграций.
- CD: разрешает публикацию (промотирование ветки/схемы) только после зелёного CI и зелёного Audit в окружении.
Два «приземлённых» сценария
A. Iceberg + Nessie + GX (Spark/Trino)
- Создаём ветку audit_branch, включаем WAP: см. пример выше. 2) Записываем партицию в ветку, 3) Прогоняем GX-чекпоинт (ключевые ожидания: PK-уникальность, not null, допустимые категории, «свежесть»), 4) fast_forward main → audit_branch. В случае провала — удаляем ветку, шлём алерт, авто-ретрай на write/transform.
Когда уместно: lakehouse, требования к атомарной публикации нескольких таблиц, быстрые откаты, дешёвое S3-хранилище.
B. Snowflake/BigQuery/Databricks + dbt (Blue/Green)
- dbt build --vars target_schema:blue
- dbt test (базовые + dbt-expectations)
- dbt clone/SWAP view/матричных таблиц в прод-схеме.
- Маркеры развертывания + метрики качества (см. ниже).
Когда уместно: SQL-warehouse, минимальные изменения в архитектуре, ставка на dbt-экосистему.
Эксплуатация: практические детали
Идемпотентность и детерминизм. Шаг Write должен быть повторяемым: один и тот же вход → один и тот же результат. Для потоков — «ровно-один-раз» семантика на уровне sink’а (Iceberg writer + транзакции).
Метрики WAP:
- Coverage тестов (какая доля критичных таблиц покрыта базовыми проверками dbt/GX).
- Pass rate и MTTR по отказам на Audit.
- «Время до публикации» (end-to-end latency) и «задержка до свежести» (freshness lag).
- Чистота веток/схем (retention, orphan-ветки).
Data Contracts. Формализуйте SLA/SLI на уровне таблиц и полей: типы, обязательность, диапазоны, правила совместимости миграций.
Общие шаблоны тестов (минимум):
- PK: unique + not_null
- FK: relationships
- Категории: accepted_values
- Бизнес-ограничения: «Цена ≥ 0», «Возраст 14–120», «Дата ≤ сегодня», «Количество ≥ 0» (singular-тесты в dbt или custom GX).
Риски и как их снизить
- Ложные тревоги и «шум» → приоритезируйте критические таблицы и проверки, постепенно расширяйте покрытие; не превращайте Audit в «всё и сразу». Помогают пакеты вроде dbt-expectations с целевым применением.
- Рост задержек → отделяйте тяжёлые статистические проверки в асинхронный слой, а публикацию блокируйте только базовыми инвариантами (PK/FK/not_null/accepted_values/freshness).
- Захламление веток/схем → политики retention и регулярный expire_snapshots/auto-cleanup.
- Неатомарная публикация нескольких таблиц → используйте форматы с транзакционной метаданной моделью и fast-forward главной ветки/каталога.
- Схема-дрейф между ветками → чёткие правила эволюции схемы и миграции «сначала в ветке, потом публикация».
Чек-лист внедрения WAP
Минимальный DoD для Publish:
- Пройдены базовые тесты (PK уникален, not null ключевых полей, FK связаны, категории валидны, свежесть не просрочена).
- Нет разрывов по количеству строк/дубликатам vs прошлый снапшот (контроль «диффа»).
- Схема совместима (backward-compatible) или сопровождается миграцией.
- Ветка/схема готова к fast-forward/Swap; прописан авто-rollback.
Выбор архитектуры:
- Lakehouse с Iceberg/лапшой из файлов → ветки Iceberg / lakeFS / Nessie.
- SQL-warehouse и dbt-практики → Blue/Green через dbt clone.
- AWS-центричный стек → добавьте Glue Data Quality к Iceberg-веткам.
- Databricks → DLT-expectations как декларативный Audit.
Примеры «кусочков» конфигураций
dbt (база)
models:
- name: fct_orders
columns:
- name: order_id
tests: [unique, not_null]
- name: status
tests:
- accepted_values:
values: ['new','paid','cancelled','returned']
- name: customer_id
tests:
- relationships:
to: ref('dim_customers')
field: id
(Базовые проверки dbt: uniqueness, non-null, accepted values, relationships.)
Iceberg (Spark SQL) — ветка для аудита
ALTER TABLE sales.orders SET TBLPROPERTIES('write.wap.enabled'='true');
ALTER TABLE sales.orders CREATE BRANCH audit_daily RETAIN 3 DAYS;
SET spark.wap.branch = audit_daily;
INSERT INTO sales.orders SELECT * FROM staging.orders_run_2025_08_22;
-- После внешних проверок:
CALL mycat.system.fast_forward('sales.orders', 'main', 'audit_daily');(Ветка + fast-forward — стандартные процедуры/DDL Iceberg.)
Когда WAP обязателен, а когда избыточен
Обязателен: регуляторная отчётность, критичные витрины (P&L, выручка, остатки), ML-фичесторы, мульти-табличные публикации с требованием атомарности.
Опционален/упрощайте: прототипы, нерегулярные одноразовые загрузки, «песочницы» аналитиков — здесь хватит тестов «на входе» без full-blown ветвления.
Как это делали «у больших»
Netflix строил платформу качества с идеей «сначала запись, потом аудит, потом публикация» и автоматизацией правил (Quinto), развивая тему аномалий и метрик качества. С тех пор паттерн стал стандартом практики в индустрии.
Вопрос–ответ (FAQ)
Q: Чем WAP отличается от просто «записал → проверил»?
A: Изоляция. При WAP «черновик» не видят потребители, публикация атомарна, а откат — быстрый (переключили указатель/ветку).
Q: Можно ли сделать WAP в потоках?
A: Да. Пишите в «теневой» топик/таблицу или ветку Iceberg sink’а, валидируйте, затем fast-forward/копируйте в основную.
Q: Как минимизировать «шум» от алертов?
A: Начните с базовой четвёрки dbt-тестов + freshness, покрывайте «золотые» таблицы, остальное добавляйте постепенно; для сложных правил используйте GX.
Q: Как быстро откатиться?
A: В lakehouse — fast-forward/rollback снапшота или переключение ветки main на предыдущую «голову». В warehouse — swap на предыдущую схему/вью (blue/green).
Q: Это не то же, что «data contracts»?
A: Контракты задают «какие поля/типы/правила допустимы», WAP — про процесс безопасной публикации. Их лучше комбинировать.
Q: Где посмотреть хорошие обзоры?
A: Dagster-гайд по WAP, руководство AWS по Iceberg-веткам + Glue Data Quality, материалы lakeFS/Nessie и записи Netflix.
WAP — не «мода», а зрелая инженерная практика, которая делает конвейеры данных предсказуемыми и доверенными: пишем изолированно, проверяем по правилам, публикуем атомарно. Вариантов реализации несколько — от Iceberg-веток и lakeFS/Nessie до dbt blue/green и DLT-expectations. Выбирайте по инфраструктуре и требованиям к атомарности/задержке, но рубеж Audit и недоступность черновика потребителям — должны оставаться неизменными.
Практическое руководство по внедрению WAP
1. Подготовительный этап
Прежде чем «поднимать» WAP, нужно ответить на три вопроса:
-
Где будем хранить черновики (Write)?
- В lakehouse — через ветки Iceberg/Delta/Hudi или lakeFS/Nessie.
- В SQL warehouse (Snowflake, BigQuery, Databricks) — через отдельные схемы и dbt clone.
- В объектном хранилище (S3, GCS, Azure) — через staging-префиксы и атомарный swap таблиц.
- Какими инструментами будем проверять (Audit)?
- dbt tests (unique, not_null, accepted_values, relationships).
- Great Expectations (GX).
- DLT expectations (для Databricks).
- AWS Glue Data Quality (для AWS-стека).
- Fast-forward ветки (main ← audit_branch).
- Swap схемы/вью (blue/green).
- CI/CD пайплайн с условием: публикация разрешена только после зелёного аудита.
- Как будем публиковать (Publish)?
2. Пошаговая инструкция (сквозной сценарий)
Шаг 1. Настрой Write
- Создаём staging-ветку (например, в Iceberg):
- ALTER TABLE sales.orders CREATE BRANCH audit_daily RETAIN 3 DAYS;
- Записываем новые данные только в эту ветку:
- SET spark.wap.branch = audit_daily;
- INSERT INTO sales.orders SELECT * FROM staging.orders_run_2025_08_22;
Шаг 2. Настрой Audit
-
Определяем обязательные проверки:
- PK уникален
- Not Null на ключевых полях
- FK (связь с измерениями)
- Бизнес-правила: цена ≥ 0, дата ≤ сегодня, возраст 14–120
- В dbt (schema.yml):
models:
- name: fct_orders
columns:
- name: order_id
tests: [unique, not_null]
- name: status
tests:
- accepted_values:
values: ['new','paid','cancelled','returned']
- name: customer_id
tests:
- relationships:
to: ref('dim_customers')
field: id
В Great Expectations:
expect_column_values_to_be_between("price", min_value=0)
expect_column_values_to_not_be_null("order_id")
Шаг 3. Организуй Publish
- Если все проверки зелёные → делаем fast-forward:
- CALL catalog.system.fast_forward('sales.orders', 'main', 'audit_daily');
- Если проверки упали → шлём алерт в Slack/Teams, ветку удаляем, процесс идёт на ретрай.
3. Чек-лист внедрения WAP
Чек-лист по Write
- Данные пишутся не в продуктив, а в staging/ветку.
- Write идемпотентен (одни и те же входные → одни и те же выходные).
- Есть стратегия очистки временных веток (retention).
Чек-лист по Audit
- Базовые инварианты: PK уникален, not_null, FK, accepted_values.
- Бизнес-правила: цены ≥ 0, даты ≤ сегодня и т.д.
- Проверяется количество строк против предыдущего снапшота (дифф).
- Проверяется свежесть (freshness).
- Тесты покрывают все таблицы 1-го класса важности (выручка, остатки, финансы).
Чек-лист по Publish
- Публикация происходит атомарно (swap/fast-forward).
- Если аудит упал, публикация невозможна.
- Есть автоматический rollback.
- В CI/CD прописан шаг: публикация доступна только при зелёном Audit.
4. Практические примеры
Пример 1. Iceberg + GX в e-commerce
- Write: загрузка заказов из Kafka в ветку orders_audit_2025_08_22.
- Audit: GX проверяет: уникальность order_id, цены ≥ 0, статус в списке, дата заказа ≤ сегодня.
- Publish: fast-forward main → ветка.
- Fail-case: если встречается заказ со статусом undefined, ветка блокируется, отправляется оповещение в Slack, заказчику не виден.
Пример 2. dbt Blue/Green в Snowflake (финансовая отчётность)
- Write: dbt build → схемa finance_blue.
- Audit: dbt tests на уникальность транзакций, баланс Дт=Кт, отсутствие NULL в обязательных полях.
- Publish: dbt clone заменяет алиас finance_current на finance_blue.
- Fail-case: если тест «баланс Дт=Кт» упал, схема finance_blue остаётся «в тени», пользователи BI видят старые данные.
5. FAQ для внедрения
Q: WAP увеличивает задержку?
A: Да, добавляется слой тестов. Решение — проверять только критические инварианты синхронно, а тяжёлые статистические тесты запускать асинхронно.
Q: Нужно ли покрывать тестами все таблицы?
A: Нет. Начинайте с критичных витрин (финансы, продажи, остатки). Остальное покрывайте постепенно.
Q: Какой минимальный набор тестов обязателен?
A: Уникальность ключей, отсутствие NULL в ключевых полях, допустимые значения категориальных полей, свежесть.
Q: Можно ли WAP в стриминге?
A: Да, пишем в «теневой» топик/таблицу, проверяем агрегаты (например, доля ошибок ≤ 1%), после чего коммитим в основную таблицу или переключаем ветку.




