Files
airflow-greenplum/airflow/dags/bookings_to_gp_dm.py
T
ddadmin 7d18b3fe5c feat(dm): добавлена эталонная витрина dm.sales_report
- Зачем:
  - требуется эталонная витрина для обучения паттернам DM-слоя
  - демонстрация UPSERT-логики с IS DISTINCT FROM для идемпотентности
- Что:
  - DDL: heap-таблица dm.sales_report с 18 полями, DISTRIBUTED BY (flight_date)
  - Load: UPSERT (UPDATE + INSERT) с JOIN dim_calendar, dim_airports (x2), dim_tariffs, fact_flight_sales
  - DQ: PL/pgSQL DO $$ с проверками непустоты, уникальности, tickets_sold >= passengers_boarded, boarding_rate BETWEEN 0 AND 1
  - DAGs: bookings_dm_ddl (DDL), bookings_to_gp_dm (ETL + DQ с параллельными ветками)
  - Tests: smoke-тесты для обоих DAG
  - sql/ddl_gp.sql: добавлен \i dm/sales_report_ddl.sql
- Проверка:
  - make fmt && make lint — passed
  - make test — 15 passed, 11 skipped
  - make ddl-gp — DDL applied successfully
  - airflow dags test bookings_to_gp_dm 2026-02-28T13:00:00 — 4 tasks SUCCESS
  - 9243 rows loaded, _load_id подставлен корректно (Jinja2 templating works)
  - DQ checks passed
2026-02-28 22:13:11 +03:00

90 lines
3.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
"""
Учебный DAG: загрузка из DDS в DM (Greenplum) по домену bookings.
Ключевая идея:
- DM читает из DDS (Star Schema);
- для каждой витрины выполняем пару задач load -> dq;
- все витрины загружаются параллельно (не зависят друг от друга);
- sales_report использует UPSERT (heap-таблица с UPDATE).
На данном этапе реализована только эталонная витрина sales_report.
Остальные витрины будут добавлены в последующих этапах.
"""
from datetime import timedelta
from logging import getLogger
import pendulum
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow import DAG
GREENPLUM_CONN_ID = "greenplum_conn"
log = getLogger(__name__)
default_args = {
"owner": "airflow",
"retries": 1,
"retry_delay": timedelta(seconds=30),
}
def _finish_summary() -> None:
"""Логирует краткий итог выполнения DM-ветки."""
log.info("DAG bookings_to_gp_dm завершён. Подробности смотрите в логах задач.")
with DAG(
dag_id="bookings_to_gp_dm",
start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
schedule=None,
catchup=False,
max_active_runs=1,
template_searchpath="/sql",
default_args=default_args,
tags=["demo", "bookings", "greenplum", "dm"],
description="Учебный DAG: загрузка DDS -> DM (Data Mart) + DQ проверки",
) as dag:
# start-задача (для структуры, вдруг понадобятся pre-checks)
start_dm = PythonOperator(
task_id="start_dm",
python_callable=lambda: log.info("Начало загрузки DM-слоя..."),
)
# === Эталонная витрина: sales_report ===
load_dm_sales_report = PostgresOperator(
task_id="load_dm_sales_report",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="dm/sales_report_load.sql",
)
dq_dm_sales_report = PostgresOperator(
task_id="dq_dm_sales_report",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="dm/sales_report_dq.sql",
)
# === Заглушки для будущих витрин (будут реализованы в этапах 2-5) ===
# load_dm_route_performance = PostgresOperator(...)
# dq_dm_route_performance = PostgresOperator(...)
# load_dm_passenger_loyalty = PostgresOperator(...)
# dq_dm_passenger_loyalty = PostgresOperator(...)
# load_dm_airport_traffic = PostgresOperator(...)
# dq_dm_airport_traffic = PostgresOperator(...)
# load_dm_monthly_overview = PostgresOperator(...)
# dq_dm_monthly_overview = PostgresOperator(...)
# finish-задача
finish_dm_summary = PythonOperator(
task_id="finish_dm_summary",
python_callable=_finish_summary,
)
# Зависимости: параллельные ветки load -> dq
# В будущем: start_dm >> [load_sales >> dq_sales, load_route >> dq_route, ...] >> finish
start_dm >> load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary