From 196171e3d6560a84cf2b2a59df7f46c1d47dfa98 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 9 Dec 2025 10:43:43 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A1=D0=BA=D0=BB=D0=B5=D0=BB=D0=B5=D1=82=20?= =?UTF-8?q?=D1=83=D1=87=D0=B5=D0=B1=D0=BD=D0=BE=D0=B3=D0=BE=20=D0=B4=D0=B0?= =?UTF-8?q?=D0=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow/dags/bookings_to_gp_stage.py | 140 +++++++++++++++++++++++++++ docs/internal/bookings_stg_design.md | 8 ++ 2 files changed, 148 insertions(+) create mode 100644 airflow/dags/bookings_to_gp_stage.py diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py new file mode 100644 index 0000000..18e027b --- /dev/null +++ b/airflow/dags/bookings_to_gp_stage.py @@ -0,0 +1,140 @@ +from __future__ import annotations + +""" +Учебный DAG для менти: показывает, как устроен поток +от источника bookings-db до слоя stg в Greenplum. + +Важно: этот файл специально должен быть хорошо задокументирован — +docstring и комментарии помогают студенту понять, «зачем» каждая задача, +а не только «что именно она делает». +""" + +import logging +from datetime import datetime, timedelta + +from airflow import DAG +from airflow.operators.python import PythonOperator + +from helpers.greenplum import get_gp_conn + + +default_args = { + "owner": "airflow", + "retries": 1, + "retry_delay": timedelta(seconds=30), +} + + +def _generate_bookings_day(load_date: str) -> None: + """ + Генерирует данные за указанный день в bookings-db. + + Важно: функция должна быть идемпотентной: + если данные за load_date уже есть в исходной БД, + повторно пересобирать день не нужно. + """ + logging.info( + "Генерация учебного дня в bookings-db за дату %s (заглушка)", load_date + ) + # TODO: реализовать проверку наличия дня в bookings.bookings + # и генерацию нового дня при его отсутствии. + + +def _get_last_loaded_ts_from_gp() -> str | None: + """ + Возвращает максимальное значение src_created_at_ts из stg.bookings. + + Пока функция возвращает None как заглушку, что соответствует + режиму полной загрузки (full). + """ + logging.info( + "Чтение последнего загруженного src_created_at_ts из stg.bookings (заглушка)" + ) + # Пример будущей реализации: + # with get_gp_conn() as conn, conn.cursor() as cur: + # cur.execute("SELECT max(src_created_at_ts) FROM stg.bookings") + # row = cur.fetchone() + # return row[0] + return None + + +def _extract_and_load_increment_via_pxf( + last_loaded_ts: str | None, + load_date: str, +) -> None: + """ + Читает дельту из stg.bookings_ext и вставляет её в stg.bookings. + + Логика: + - если last_loaded_ts is None → первая загрузка (full), + берём все данные за load_date и ранее; + - иначе берём только записи, где src_created_at_ts > last_loaded_ts + и не позже конца учебного дня. + """ + logging.info( + "Загрузка инкремента через PXF: last_loaded_ts=%s, load_date=%s (заглушка)", + last_loaded_ts, + load_date, + ) + # TODO: реализовать INSERT INTO stg.bookings (...) SELECT ... FROM stg.bookings_ext + # с учётом инкрементального окна по src_created_at_ts. + + +def _check_row_counts(load_date: str) -> None: + """ + Проверяет, что количество строк из источника и в stg.bookings совпадает. + + Эта проверка должна помочь студенту увидеть пример простой DQ‑проверки + для инкрементальной загрузки. + """ + logging.info("Проверка количества строк за %s (заглушка)", load_date) + # TODO: реализовать сравнение количества строк, + # например через SELECT COUNT(*) в источнике и в stg.bookings. + + +def _finish_summary() -> None: + """Логирует краткий итог выполнения DAG за один запуск.""" + logging.info("DAG bookings_to_gp_stage завершён (пока только скелет).") + + +with DAG( + dag_id="bookings_to_gp_stage", + start_date=datetime(2024, 1, 1), + schedule=None, + catchup=False, + default_args=default_args, + tags=["demo", "bookings", "greenplum", "stg"], + description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)", +) as dag: + generate_bookings_day = PythonOperator( + task_id="generate_bookings_day", + python_callable=_generate_bookings_day, + op_kwargs={"load_date": "{{ ds }}"}, + ) + + get_last_loaded_ts = PythonOperator( + task_id="get_last_loaded_ts_from_gp", + python_callable=_get_last_loaded_ts_from_gp, + ) + + extract_and_load_increment = PythonOperator( + task_id="extract_and_load_increment_via_pxf", + python_callable=_extract_and_load_increment_via_pxf, + op_kwargs={ + "last_loaded_ts": "{{ ti.xcom_pull(task_ids='get_last_loaded_ts_from_gp') }}", + "load_date": "{{ ds }}", + }, + ) + + check_row_counts = PythonOperator( + task_id="check_row_counts", + python_callable=_check_row_counts, + op_kwargs={"load_date": "{{ ds }}"}, + ) + + finish_summary = PythonOperator( + task_id="finish_summary", + python_callable=_finish_summary, + ) + + generate_bookings_day >> get_last_loaded_ts >> extract_and_load_increment >> check_row_counts >> finish_summary diff --git a/docs/internal/bookings_stg_design.md b/docs/internal/bookings_stg_design.md index 21ae467..18474a2 100644 --- a/docs/internal/bookings_stg_design.md +++ b/docs/internal/bookings_stg_design.md @@ -124,3 +124,11 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre Дальнейшая модель DWH (слои ODS/DDS/DM, факт/измерения, SCD) должна быть спроектирована студентом по статье о моделировании данных, используя `stg.bookings` как входной слой. +## 6. Требования к читаемости и комментариям + +- DAG’и `bookings_stg_ddl` и `bookings_to_gp_stage` — это учебный материал для менти. +- В коде DAG’ов должны быть: + - понятные docstring на русском у всех функций и DAG; + - короткие комментарии рядом с нетривиальной логикой (особенно вокруг инкремента и идемпотентности); + - говорящие `task_id` и названия функций, отражающие их роль в процессе. +- Цель: чтобы по одному только коду DAG студент мог восстановить архитектуру процесса и сопоставить её с теорией из статьи про моделирование DWH.