From ff831b6bb285501808482df3040e7a9b3b7d11ec Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 9 Jan 2026 14:51:21 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A4=D0=B8=D0=BA=D1=81=20DAG=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=B3=D1=80=D1=83=D0=B7=D0=BA=D0=B8=20=D0=B4=D0=B0=D0=BD=D0=BD?= =?UTF-8?q?=D1=8B=D1=85=20-=20=D0=BF=D0=BE=D1=81=D1=82=D1=80=D0=BE=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D0=B5=20=D0=B4=D0=B5=D0=BB=D1=8C=D1=82=D1=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 2 +- airflow/dags/bookings_to_gp_stage.py | 2 +- sql/stg/bookings_dq.sql | 2 +- sql/stg/bookings_load.sql | 24 +++++++++++++++--------- 4 files changed, 18 insertions(+), 12 deletions(-) diff --git a/README.md b/README.md index 6e30c6f..1a22e46 100644 --- a/README.md +++ b/README.md @@ -342,7 +342,7 @@ load_bookings_to_stg = PostgresOperator( task_id="load_bookings_to_stg", postgres_conn_id="greenplum_conn", sql="stg/bookings_load.sql", - params={"batch_id": "{{ ds_nodash }}"}, + params={"batch_id": "{{ run_id }}"}, ) ``` diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 078affc..1112106 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -12,7 +12,7 @@ from __future__ import annotations Каждый запуск DAG работает как «шаг по времени вперёд»: - генератор в демо-БД bookings добавляет следующий учебный день после max(book_date); - загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей; -- логическая дата запуска (`ds`) используется как удобная метка запуска (через `ds_nodash` в `batch_id`, в логах и DQ). +- `run_id` используется как метка запуска (в `batch_id`, в логах и DQ). """ from datetime import datetime, timedelta diff --git a/sql/stg/bookings_dq.sql b/sql/stg/bookings_dq.sql index 62fbb73..7d8cc06 100644 --- a/sql/stg/bookings_dq.sql +++ b/sql/stg/bookings_dq.sql @@ -5,7 +5,7 @@ DO $$ DECLARE - v_batch_id text := '{{ ds_nodash }}'::text; + v_batch_id text := '{{ run_id }}'::text; v_prev_ts timestamp; v_src_count bigint; v_stg_count bigint; diff --git a/sql/stg/bookings_load.sql b/sql/stg/bookings_load.sql index fdd028f..55cd551 100644 --- a/sql/stg/bookings_load.sql +++ b/sql/stg/bookings_load.sql @@ -12,19 +12,25 @@ INSERT INTO stg.bookings ( batch_id ) SELECT - book_ref::text, - book_date::text, - total_amount::text, - book_date::timestamp, + ext.book_ref::text, + ext.book_date::text, + ext.total_amount::text, + ext.book_date::timestamp, now(), - '{{ ds_nodash }}'::text -FROM stg.bookings_ext -WHERE book_date > COALESCE( + '{{ run_id }}'::text +FROM stg.bookings_ext AS ext +WHERE ext.book_date > COALESCE( ( SELECT max(src_created_at_ts) FROM stg.bookings - WHERE batch_id <> '{{ ds_nodash }}'::text + WHERE batch_id <> '{{ run_id }}'::text OR batch_id IS NULL ), TIMESTAMP '1900-01-01 00:00:00' -); +) + AND NOT EXISTS ( + SELECT 1 + FROM stg.bookings AS b + WHERE b.batch_id = '{{ run_id }}'::text + AND b.book_ref = ext.book_ref::text + );