даг наконец-то работает
This commit is contained in:
@@ -36,6 +36,12 @@
|
||||
- Запустить вручную после первого DAG.
|
||||
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
|
||||
|
||||
- DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG):
|
||||
- предварительно выполнить один раз: `make bookings-init` (инициализация демо‑БД bookings) и `make ddl-gp` (создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum);
|
||||
- включить DAG `bookings_to_gp_stage` и запустить `Trigger DAG`;
|
||||
- убедиться, что все задачи (`generate_bookings_day`, `load_bookings_to_stg`, `check_row_counts`, `finish_summary`) завершились со статусом Success;
|
||||
- при желании проверить данные: в `bookings-db` появился новый день, а в Greenplum в `stg.bookings` — строки с актуальным `batch_id` (см. пример запросов в разделе 5).
|
||||
|
||||
- (опционально, для менторов/разработчиков) Smoke-тест DAG через Airflow CLI без UI:
|
||||
- `docker compose -f docker-compose.yml exec gp_airflow_web airflow dags test bookings_to_gp_stage 2024-01-01` — прогоняет `bookings_to_gp_stage` целиком в «off-line» режиме;
|
||||
- `docker compose -f docker-compose.yml exec gp_airflow_web airflow dags trigger bookings_to_gp_stage` — создаёт реальный запуск DAG (логи и статус можно смотреть либо через UI, либо командой `airflow tasks list`/`airflow tasks logs` внутри контейнера).
|
||||
|
||||
@@ -69,9 +69,6 @@ with DAG(
|
||||
task_id="load_bookings_to_stg",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_load.sql",
|
||||
params={
|
||||
"batch_id": "{{ ds_nodash }}",
|
||||
},
|
||||
)
|
||||
|
||||
# 3. Проверяем количество строк между источником и stg.bookings
|
||||
@@ -79,9 +76,6 @@ with DAG(
|
||||
task_id="check_row_counts",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_dq.sql",
|
||||
params={
|
||||
"batch_id": "{{ ds_nodash }}",
|
||||
},
|
||||
)
|
||||
|
||||
# 4. Финальный лог/сводка
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
|
||||
DO $$
|
||||
DECLARE
|
||||
v_batch_id text := {{ params.batch_id | tojson }}::text;
|
||||
v_batch_id text := '{{ ds_nodash }}'::text;
|
||||
v_prev_ts timestamp;
|
||||
v_src_count bigint;
|
||||
v_stg_count bigint;
|
||||
|
||||
@@ -17,13 +17,13 @@ SELECT
|
||||
total_amount::text,
|
||||
book_date::timestamp,
|
||||
now(),
|
||||
{{ params.batch_id | tojson }}::text
|
||||
'{{ ds_nodash }}'::text
|
||||
FROM stg.bookings_ext
|
||||
WHERE book_date > COALESCE(
|
||||
(
|
||||
SELECT max(src_created_at_ts)
|
||||
FROM stg.bookings
|
||||
WHERE batch_id <> {{ params.batch_id | tojson }}::text
|
||||
WHERE batch_id <> '{{ ds_nodash }}'::text
|
||||
OR batch_id IS NULL
|
||||
),
|
||||
TIMESTAMP '1900-01-01 00:00:00'
|
||||
|
||||
Reference in New Issue
Block a user