BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс по Apache Airflow и NiFi » Ассеты в оркестраторах данных: концепция и реализация в Airflow

Ассеты в оркестраторах данных: концепция и реализация в Airflow

Что такое ассеты?

В контексте систем оркестрации данных (Airflow, Dagster, Prefect) ассет — это абстракция, представляющая собой материальный объект данных: таблицу в базе данных, файл в облачном хранилище, модель в ML-реестре, агрегированную витрину в DWH и т.д.

У каждого ассета есть уникальный идентификатор, который, как правило, оформлен в стиле URI:

  • postgresql://db/schema/table
  • s3://bucket/path/to/file.parquet
  • mlflow://model_registry/model_name/version

 

Ассеты позволяют описывать не просто задачи (tasks) и их порядок, а реальные данные и зависимости между ними. Таким образом мы получаем граф данных, а не только граф задач.

 

Ключевые особенности ассетов:

  • Фокус не на коде, а на данных.
  • Возможность отслеживания и документирования зависимости между данными.
  • Прозрачность: легко понять, какие процессы влияют на конкретный набор данных.
  • Возможность построения lineage-графов (Data Lineage).

 

 

Почему эта концепция важна?

Рассмотрим аналогию из описания Prefect: в игре SimCity вы управляете одним и тем же городом, но можете переключаться между разными “линзами” — транспорт, электричество, водоснабжение.

С ассетами то же самое: вы можете смотреть на свою систему через призму зависимостей данных, истории обновлений или производных объектов.

Преимущества:

  • Упрощение сопровождения: легче понять, что и когда сломалось.
  • Минимизация переработок: если upstream-ассет не изменился, downstream-ассет можно не пересчитывать.
  • Возможность частичного пересчёта цепочек (incremental builds).

 

Как Airflow работает с ассетами

История появления ассетов в Airflow

В Airflow 2.4 появился объект Dataset, позволяющий определять зависимость между задачами на основе данных, а не только графа задач. Это было революционным шагом для Airflow, который до этого жил исключительно в парадигме DAG → Tasks.

В Airflow 3.0 Dataset был переименован в Asset, что подчёркивает фокус на данных. Появился декоратор @asset для объявления функций, которые создают ассеты.

 

Пример использования в Airflow

from airflow.decorators import asset
from airflow.datasets import Asset
sales_table = Asset("postgresql://analytics/sales")
report_file = Asset("s3://reports/monthly_report.csv")
@asset(sales_table)
def extract_sales():
    # Логика извлечения данных из внешнего источника
    ...
@asset(report_file, inputs=[sales_table])
def generate_report():
    # Генерация отчёта на основе sales_table
    ...

 

Что происходит:

  • Мы явно указываем, какие ассеты производит каждая функция.
  • Airflow строит граф зависимостей ассетов, который существует параллельно с DAG.
  • Можно запускать DAG, когда обновляется определённый ассет.

 

Ограничения ассетов в Airflow

В отличие от Prefect и Dagster:

  • Нет динамической материализации: нельзя “на лету” добавить новый ассет в процессе выполнения.
  • Нет императивного создания ассета: ассеты создаются строго в момент выполнения задачи с @asset.
  • Слабее интеграция с lineage — нужен дополнительный инструмент (например, OpenLineage).

 

Как строится граф зависимостей в Airflow

Airflow отслеживает “события материализации” — факт создания или обновления ассета.

Механизм:

  1. Task, помеченная @asset, в конце работы “помечает” указанный ассет как обновлённый.
  2. Все DAG-и, которые зависят от этого ассета, могут быть автоматически запущены.
  3. Это реализует data-driven scheduling — запуск по готовности данных, а не только по расписанию.

 

Практическое применение ассетов в Airflow

  • Автоматический запуск цепочек по готовности данных (например, после загрузки новой партии данных в S3).
  • Построение lineage-графов для аудита.
  • Интеграция с DQ-системами: можно “замораживать” downstream-задачи, если upstream-ассет не прошёл проверку качества.
  • Инкрементальные пайплайны — запуск только тех частей, которые реально затронуты изменениями.

 

Технические риски и ограничения

Риски

  1. Сложность отладки — граф ассетов может быть неочевиден при большом количестве зависимостей.
  2. Ложные срабатывания — автоматический запуск DAG-а из-за изменения ассета, даже если изменение было незначительным.
  3. Отсутствие нативной CDC-логики — ассеты фиксируют факт обновления, но не передают дельту.
  4. Сложность при микросервисной архитектуре — ассеты в разных Airflow-инстансах не синхронизируются напрямую.

 

Ограничения Airflow-реализации

  • Нет встроенной визуализации lineage без плагинов.
  • Нет возможности группировать ассеты в “логические наборы” нативными средствами (только через DAG).
  • Отсутствует полноценная интеграция с внешними каталогами данных (Data Catalog).

 

Лучшие практики при работе с ассетами в Airflow

  1. Используйте осмысленные URI
    Не table1, а postgresql://prod.analytics.sales_by_region.
  2. Разделяйте слой ассетов и слой задач
    Tasks — это операции, ассеты — это данные. Смешивание может усложнить поддержку.
  3. Интегрируйте с Data Catalog / OpenLineage
    Чтобы визуализировать зависимости и автоматизировать документирование.
  4. Внедряйте тесты качества данных
    Чтобы downstream-ассеты не обновлялись при ошибочных данных.

 

Сравнение подходов: Airflow, Dagster, Prefect

Характеристика

Airflow

Dagster

Prefect

Фокус

DAG-и и ассеты как дополнение

Ассеты в центре архитектуры

Гибкая смесь задач и ассетов

Динамическая материализация

×

✓

✓

Data-driven запуск

✓

✓

✓

Интеграция с lineage

Ограничено

Глубокая

Автоматическая

Кривая обучения

Средняя

Выше средней

Низкая

 

Ассеты в Airflow — это шаг в сторону data-centric парадигмы оркестрации, где важнее понимать, какие данные производятся и используются, чем просто управлять порядком задач.

Хотя реализация Airflow пока уступает Prefect и Dagster в гибкости, она уже даёт:

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

 

Для команд, которые уже используют Airflow и хотят начать мыслить категориями данных, ассеты — это хороший старт.

 

Высокоуровневая схема

+------------------------------ Airflow Control Plane ---------------------------+
|                                                                                |
|   +-----------+     +-----------+     +-----------+     +------------------+   |
|   |  Web UI   |<--->| Scheduler |<--->| Triggerer |<--->|    Executor      |   |
|   +-----------+     +-----------+     +-----------+     +---------+--------+   |
|                                                             ^       |          |
|                                                             |       v          |
|                                                   +--------+-----+  +-------+  |
|                                                   |  Worker Pool |  | Sensors| |
 
|                                                   +--------------+  +-------+  |
|                                                                                |
|   +--------------+     +-----------------+      +---------------------------+  |
|   |  Metastore   |<--->|   Lineage/DQ    |<---->|   Logging/Observability  |   |
|   | (DB: DAGs,   |     | (OpenLineage,   |      |  (S3/GCS, Loki, ELK)     |   |
|   |  runs, assets)|    |  GreatExpect.)  |      +---------------------------+  |
+--------------------------------------------------------------------------------+
 
                                   |
                                   | emits/consumes "Asset Materialization" events
                                   v
 
+-------------------------------------- Data Plane --------------------------------------+
|                                                                                        |
|   +------------------+      +-------------------+      +---------------------------+   |
|   |  Sources         | ---> |  Tasks/@asset     | ---> |  Sinks / Target Assets    |   |
|   | (DB, APIs, Kafka)|      |  (Python, SQL)    |      |  (DWH tables, S3 files,   |   |
|   +------------------+      +-------------------+      |   ML registry models)     |   |
|                                                        +---------------------------+   |
|                                                                                        |
+----------------------------------------------------------------------------------------+

 

Идея: задачи (Tasks) с декоратором @asset производят/обновляют ассеты (таблицы/файлы/модели). Факт обновления фиксируется как событие материализации ассета. Планировщик (Scheduler) использует эти события, чтобы запускать зависимые DAG’и/задачи по готовности данных (data-driven scheduling). Метастор хранит реестр DAG’ов, ран-метаданные и карту ассетов. DQ/Lineage-сервисы получают события, строят lineage и валидируют данные. Логи и метрики идут в систему наблюдаемости.

 

Детальный поток событий (sequence view)

(1) Upstream Task (@asset) стартует
    |
    |  читает исходные данные
    v
+-----------------------------------+
| Преобразование/загрузка данных    |
| (Python/SQL/библиотеки, Spark etc)|
+-----------------------------------+
    |
    |  записывает целевой объект
    v
Target Asset (URI: postgresql://prod/analytics/sales_by_region)
    |
    |  --> emit AssetMaterialized {uri, run_id, ts, version, checksum, metadata}
    v
+------------------- Airflow Scheduler -------------------+
|  - регистрирует событие ассета                         |
|  - вычисляет затронутые зависимости                    |
|  - помещает в очередь запусков нужные DAG runs         |
+--------------------------------------------------------+
    |
    v
Dependent DAG(s) (dataset/asset-triggered) -> Executor -> Workers
    |
    |  для каждого downstream шага: проверка предусловий
    |  (в т.ч. DQ / SLAs / freshness / partition filters)
    v
Downstream Task(s) (@asset or @task) -> производят след. ассеты
    |
    +--> повтор цикла (каскадная материализация)

 

Модель данных и адресация ассетов (ключи в стиле URI)

Asset Key (URI)       : <scheme>://<authority>/<path>[?<query>][#<fragment>]
Примеры:
  - postgresql://prod/analytics/sales_by_region
  - s3://datalake/bronze/sales/2025/08/13/part-0001.parquet
  - mlflow://model-registry/churn_model/versions/42
 
Metadata (минимум):
  - logical_name, uri, producing_task_id, producing_dag_id
  - materialized_at (ts), run_id, data_version/hash/checksum
  - optional: partition_keys, freshness_slo, schema_fingerprint

 

Определение зависимостей (data-driven scheduling)

+--------------------------+
|  DAG_B (consumer)        |
|  Triggers on:            |
|    Asset A = s3://...    |
|    Asset B = postgresql://...  (AND condition) 
+--------------------------+

 

Условие запуска:

  • обе материализации (A и B) новее последнего успешного run DAG_B
  • опционально: совпали partition keys / data_version
  • не нарушены DQ/SLAs

 

Ж​изненный цикл ассета

[REGISTER] -> [MATERIALIZE] -> [VALIDATE] -> [PUBLISH] -> [SERVE] -> [RETIRE]
                 |                |            |
                 |                |            +--> открытие потребителям (BI/ML)
                 |                +--> DQ checks (Great Expectations, custom)
                 +--> событие AssetMaterialized (Airflow datasets/assets)

 

Смысл: после записи ассета происходит материализация и фиксация события. Далее проходят проверки качества/свежести, затем ассет публикуется и потребляется системами BI/ML. Когда схема или SLA устаревают — ассет снимается с публикации (retire) и замещается новой версией.

 

Встраивание DQ и Lineage

Upstream @asset ----> emits materialization
        |                 |
        |                 +--> OpenLineage exporter ---> Lineage Backend (Marquez/Collibra/Atlas)
        |
        +--> DQ Hook ---> Great Expectations (context)
                            |
                            +--> Validation Result (pass/fail, metrics)
                                   |
                                   +--> back to Scheduler (block/unblock downstream)

 

Сенсоры и внешние триггеры

External System (e.g., Data Lake Loader) 
    |
    +--> writes s3://landing/inbox/file.csv
            |
            +--> FileSensor / AssetSensor detects
                       |
                       +--> emits AssetMaterialized (for landing asset)
                               |
                               +--> Scheduler triggers ingestion DAG

 

Пример (сквозной, с URI и зависимостью)

Assет 1 (bronze):
  s3://datalake/bronze/sales/2025/08/13/*.json
  produced by: dag=ingest_sales, task=@asset bronze_sales
 
Assет 2 (silver):
  s3://datalake/silver/sales/2025/08/13/*.parquet
  produced by: dag=transform_sales, task=@asset silver_sales
  depends_on: bronze_sales
 
Ассет 3 (gold, DWH):
  postgresql://prod/analytics/sales_by_region
  produced by: dag=mart_sales, task=@asset mart_sales_by_region
  depends_on: silver_sales
 
BI-витрина (Looker/Power BI):
  connects to: postgresql://prod/analytics/sales_by_region
  freshness SLO: <= 2h since last bronze ingest

 

Поток:

  1. ingest_sales кладет сырые JSON в bronze и эмитит материализацию.
  2. transform_sales срабатывает по готовности bronze, пишет silver Parquet и эмитит материализацию.
  3. mart_sales срабатывает по silver, обновляет витрину в Postgres и эмитит материализацию.
  4. BI проверяет свежесть gold ассета.

 

Обратная заправка (backfill) и частичная материализация

Backfill window: 2025-08-01..2025-08-13 (by day partitions)
 
for d in days:
  materialize(brownz/silver/gold for partition=d)
  emit events per partition
  trigger downstream only for d

 

Итог: пересчитываются только нужные партиции и каскады ниже по графу

 

Ошибки, риски, защиты

Риск: Оркестрация сработает от "шумной" материализации (пустые обновления).

Защита: сравнивать schema_fingerprint / checksum; игнорировать no-op.

 

Риск: Несогласованные версии данных между ассетами.

Защита: передача data_version/partition_keys в событии; strict matching при триггере.

 

Риск: Медленные сенсоры / дупликаты событий.

Защита: Reschedule-сенсоры, идемпотентные мьютексы (task instance lock), дедупликация по (uri, version).

 

Риск: Невидимые ошибки качества (битые партиции).

Защита: обязательные DQ-чекпоинты до эмиссии "валидной" материализации (или блокировка downstream при fail).

 

Риск: Разрыв lineage при сторонних загрузках.

Защита: стандартизованные эмиттеры событий, policy "no write without materialize event".

 

Мини-карта реализации в Airflow (понятиями)

DAG code:
  - @asset(outputs=[Asset("s3://.../silver/...")], inputs=[Asset("s3://.../bronze/...")])
  - Dataset/Asset dependencies declared per task
  - DatasetTriggeredDagRun (Airflow 2.4+) / Assets (Airflow 3.x)
  - Sensors (FileSensor, ExternalTask/AssetSensor)
  - DQ hooks (Great Expectations) as pre/post-task checks
  - OpenLineage provider to emit lineage on task success
 
Control:
  - Scheduler listens for dataset/asset updates (metastore)
  - Triggerer wakes deferred tasks (sensors/async ops)
  - Executor schedules on Workers (Celery/K8s/Local)
 
State:
  - Metastore keeps DAG runs, task instances, asset events
  - Logs to object storage; metrics to Prometheus/StatsD

 

Эта текстовая схема отражает, как ассеты адресуются URI, как события материализации двигают вычисления «по данным», и где Airflow встраивает DQ, сенсоры и lineage, чтобы получить предсказуемый и воспроизводимый конвейер

 

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Настройка пайплайна с использованием Airflow и PostgreSQL
Следующая статья →
Разработка и внедрение производственных ML-пайплайнов на Apache Airflow

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • ЭГИС - международная фармацевтическая компания, основанная в 1907 году в Венгрии. Компания имеет представительства более чем в 60 странах мира, в том числе в России. Компания ЭГИС является одним из ведущих производителей дженерических лекарственных средств в Центральной и Восточной Европе. Её деятельность охватывает все звенья производственно-сбытовой фармацевтической цепочки.

  • В 2003 году Мерсико и пятью микрокредитными агентствами Мерсико было принято историческое решение о консолидации активов по всей территории Кыргызстана в целях образования национального финансового института по развитию сообществ - Компаньона. В октябре 2004 года Компаньон был зарегистрирован Национальным банком Кыргызской Республики.

  • ООО «Ай Пи Ти Групп» (IPT Group) — многопрофильный консалтинговый холдинг, специализирующийся на юридическом и финансовом сопровождении бизнеса. IPT Group занимает высокие позиции в профессиональных рейтингах, входит в ТОП-30 лучших юридических компаний России по версии «Право.ru-300», Global Law Experts и др.

  • «Лента» – первая по величине сеть гипермаркетов и четвертая среди крупнейших розничных сетей страны. Компания была основана в 1993 г. в Санкт-Петербурге.

    «Лента» управляет 249 гипермаркетами в 88 городах России и 131 супермаркетом в Москве, Санкт-Петербурге, Сибири, Уральском и Центральном регионах с общей торговой площадью около 1 494 тыс. кв. м. Средняя торговая площадь одного гипермаркета «Лента» составляет около 5 500 кв.м, средняя площадь супермаркета – 800 кв.м. Компания оперирует двенадцатью распределительными центрами. Штат компании – около 50, 5 тыс. человек.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.