Уход от запуска DAG за конкретную дату.
Теперь новый запуск генерирует и переливает новый, следующий день
This commit is contained in:
@@ -84,11 +84,11 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
- Этот DAG не загружает данные, только подготавливает структуру.
|
||||
- Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT).
|
||||
|
||||
### 4.2. DAG для ежедневной загрузки
|
||||
### 4.2. DAG для пошаговой загрузки
|
||||
|
||||
- `dag_id`: `bookings_to_gp_stage`.
|
||||
- Основные параметры:
|
||||
- `load_date` (по умолчанию `{{ ds }}`) — учебный день, за который генерим и грузим данные;
|
||||
- `batch_id` (по умолчанию `{{ ds_nodash }}`) — метка батча, которая попадает в `stg.bookings.batch_id`;
|
||||
- подключения:
|
||||
- `bookings_db_conn_id` — Airflow connection к `bookings-db` (в коде DAG — `BOOKINGS_CONN_ID = "bookings_db"`);
|
||||
- `greenplum_conn_id` — Airflow connection к Greenplum (`GREENPLUM_CONN_ID = "greenplum_conn"`).
|
||||
@@ -97,22 +97,22 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
|
||||
1. `generate_bookings_day`
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет скрипт `/sql/bookings/generate_day_if_missing.sql`;
|
||||
- скрипт сам идемпотентно проверяет, есть ли данные за `load_date` в `bookings.bookings`:
|
||||
- если день уже есть — выдаёт NOTICE и ничего не делает;
|
||||
- если нет — вызывает генератор демобазы (логика как в `bookings/generate_next_day.sql`).
|
||||
- выполняет скрипт `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт смотрит на `max(book_date)` и:
|
||||
- если база пуста — берёт стартовую дату из конфигурации (`bookings.start_date`) и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` (логика как в `bookings/generate_next_day.sql`).
|
||||
2. `load_bookings_to_stg`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_load.sql`;
|
||||
- внутри SQL считается `max(src_created_at_ts)` по «старым» батчам и по нему строится окно инкремента:
|
||||
- первая загрузка (full) — берём все строки из `stg.bookings_ext`;
|
||||
- последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума и не позже конца учебного дня;
|
||||
- последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума (верхняя граница по дате не задаётся явно);
|
||||
- при вставке заполняются тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`.
|
||||
3. `check_row_counts`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_dq.sql`;
|
||||
- скрипт заново считает окно инкремента по тем же правилам, что и загрузка, и сравнивает:
|
||||
- количество строк в `stg.bookings_ext` за окно,
|
||||
- количество строк в `stg.bookings_ext` с `book_date` позже «старого» максимума,
|
||||
- количество строк в `stg.bookings` для текущего `batch_id`;
|
||||
- при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
4. `finish_summary`
|
||||
|
||||
@@ -32,8 +32,9 @@ _Внутренний файл, чтобы не забыть договорён
|
||||
- `generate_bookings_day`:
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет SQL `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт сам проверяет, есть ли строки за `load_date` (`book_date::date = load_date`);
|
||||
- если день уже сгенерирован — ничего не делает (идемпотентность), только пишет NOTICE.
|
||||
- скрипт смотрит на `max(book_date)` в `bookings.bookings`:
|
||||
- если база пустая — берёт стартовую дату из GUC и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` и пишет NOTICE с интервалом генерации.
|
||||
- `load_bookings_to_stg`:
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет SQL `/sql/stg/bookings_load.sql`;
|
||||
@@ -54,15 +55,15 @@ _Внутренний файл, чтобы не забыть договорён
|
||||
- `make bookings-init`
|
||||
- `make ddl-gp`
|
||||
2. Открыть Airflow UI (`http://localhost:8080`) и включить DAG `bookings_to_gp_stage`.
|
||||
3. Вызвать `Trigger` DAG с конкретной датой (например, `2017-01-01` → зависит от стартовой конфигурации демобазы).
|
||||
3. Вызвать `Trigger` DAG (дату логического запуска можно оставить по умолчанию — она используется только как метка `batch_id`).
|
||||
4. Посмотреть:
|
||||
- в `bookings-db` ― что появился день с бронированиями;
|
||||
- в Greenplum (`make gp-psql`) — данные в `stg.bookings`:
|
||||
- `SELECT * FROM stg.bookings LIMIT 10;`
|
||||
- `SELECT src_created_at_ts, load_dttm, batch_id FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;`
|
||||
5. Перезапустить DAG для следующей даты и увидеть, что:
|
||||
- генерация в `bookings.bookings` идёт по одному дню вперёд;
|
||||
- в `stg.bookings` появляются только новые записи (delta).
|
||||
5. Перезапустить DAG ещё несколько раз и увидеть, что:
|
||||
- генерация в `bookings.bookings` идёт по одному дню вперёд от текущего `max(book_date)`;
|
||||
- в `stg.bookings` появляются только новые записи (delta), помеченные разными `batch_id`.
|
||||
|
||||
## 5. Примечания «на потом»
|
||||
|
||||
|
||||
Reference in New Issue
Block a user