Фикс DAG загрузки данных - построение дельты
This commit is contained in:
@@ -342,7 +342,7 @@ load_bookings_to_stg = PostgresOperator(
|
|||||||
task_id="load_bookings_to_stg",
|
task_id="load_bookings_to_stg",
|
||||||
postgres_conn_id="greenplum_conn",
|
postgres_conn_id="greenplum_conn",
|
||||||
sql="stg/bookings_load.sql",
|
sql="stg/bookings_load.sql",
|
||||||
params={"batch_id": "{{ ds_nodash }}"},
|
params={"batch_id": "{{ run_id }}"},
|
||||||
)
|
)
|
||||||
```
|
```
|
||||||
|
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ from __future__ import annotations
|
|||||||
Каждый запуск DAG работает как «шаг по времени вперёд»:
|
Каждый запуск DAG работает как «шаг по времени вперёд»:
|
||||||
- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date);
|
- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date);
|
||||||
- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей;
|
- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей;
|
||||||
- логическая дата запуска (`ds`) используется как удобная метка запуска (через `ds_nodash` в `batch_id`, в логах и DQ).
|
- `run_id` используется как метка запуска (в `batch_id`, в логах и DQ).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|||||||
@@ -5,7 +5,7 @@
|
|||||||
|
|
||||||
DO $$
|
DO $$
|
||||||
DECLARE
|
DECLARE
|
||||||
v_batch_id text := '{{ ds_nodash }}'::text;
|
v_batch_id text := '{{ run_id }}'::text;
|
||||||
v_prev_ts timestamp;
|
v_prev_ts timestamp;
|
||||||
v_src_count bigint;
|
v_src_count bigint;
|
||||||
v_stg_count bigint;
|
v_stg_count bigint;
|
||||||
|
|||||||
@@ -12,19 +12,25 @@ INSERT INTO stg.bookings (
|
|||||||
batch_id
|
batch_id
|
||||||
)
|
)
|
||||||
SELECT
|
SELECT
|
||||||
book_ref::text,
|
ext.book_ref::text,
|
||||||
book_date::text,
|
ext.book_date::text,
|
||||||
total_amount::text,
|
ext.total_amount::text,
|
||||||
book_date::timestamp,
|
ext.book_date::timestamp,
|
||||||
now(),
|
now(),
|
||||||
'{{ ds_nodash }}'::text
|
'{{ run_id }}'::text
|
||||||
FROM stg.bookings_ext
|
FROM stg.bookings_ext AS ext
|
||||||
WHERE book_date > COALESCE(
|
WHERE ext.book_date > COALESCE(
|
||||||
(
|
(
|
||||||
SELECT max(src_created_at_ts)
|
SELECT max(src_created_at_ts)
|
||||||
FROM stg.bookings
|
FROM stg.bookings
|
||||||
WHERE batch_id <> '{{ ds_nodash }}'::text
|
WHERE batch_id <> '{{ run_id }}'::text
|
||||||
OR batch_id IS NULL
|
OR batch_id IS NULL
|
||||||
),
|
),
|
||||||
TIMESTAMP '1900-01-01 00:00:00'
|
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
|
||||||
|
);
|
||||||
|
|||||||
Reference in New Issue
Block a user