diff --git a/Makefile b/Makefile index 7d50653..2ddabcb 100644 --- a/Makefile +++ b/Makefile @@ -42,11 +42,11 @@ bookings-clone-demodb: git -C bookings/demodb checkout $(DEMODB_COMMIT); \ fi # Патчим generate/continue: при jobs=1 запускаем process_queue синхронно, без dblink - if ! grep -q "CALL process_queue(end_date);" bookings/demodb/engine.sql; then \ + if ! grep -q "Job 1 (local): ok" bookings/demodb/engine.sql; then \ patch -d bookings/demodb -p1 --forward < bookings/patches/engine_jobs1_sync.patch || true; \ fi - # Делаем установку идемпотентной: DROP DATABASE IF EXISTS demo - if ! grep -q "DROP DATABASE IF EXISTS demo;" bookings/demodb/install.sql; then \ + # Делаем установку идемпотентной и принудительной: DROP DATABASE IF EXISTS demo WITH (FORCE) + if ! grep -q "DROP DATABASE IF EXISTS demo WITH (FORCE);" bookings/demodb/install.sql; then \ patch -d bookings/demodb -p1 --forward < bookings/patches/install_drop_if_exists.patch || true; \ fi diff --git a/README.md b/README.md index a18d57b..97057f0 100644 --- a/README.md +++ b/README.md @@ -216,12 +216,17 @@ docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POS - Рекомендуемый учебный сценарий для DAG `bookings_to_gp_stage`: запускать DAG по одному дню вперёд, выбирая в форме Trigger логическую дату `Execution Date (ds)`, совпадающую с тем днём, который вы хотите загрузить (например, `2017-01-01`, затем `2017-01-02` и т.д.). - Быстрее всего: `make bookings-generate-day` — читает GUC и сам вызывает `continue`. -- Вручную из psql/DBeaver: +- Вручную из psql/DBeaver (подзапросы в аргументах CALL не работают, поэтому через DO-блок): ```sql - CALL continue( - (SELECT date_trunc('day', max(book_date)) + interval '1 day' FROM bookings.bookings) - ); - -- или с параллельностью: CALL continue((SELECT ...), 4); + DO $$ + DECLARE + v_next_day timestamptz; + BEGIN + SELECT date_trunc('day', max(book_date)) + interval '1 day' + INTO v_next_day + FROM bookings.bookings; + CALL continue(v_next_day); -- или CALL continue(v_next_day, 4) для параллельности + END $$; ``` - Не вызывайте `CALL generate(...)` поверх существующих данных: она делает TRUNCATE и создаёт демобазу заново. diff --git a/bookings/README.md b/bookings/README.md index 34a79b5..cce2c2f 100644 --- a/bookings/README.md +++ b/bookings/README.md @@ -9,3 +9,33 @@ Основные команды см. в корневом `Makefile` (`bookings-init`, `bookings-generate-day`, `bookings-psql`) и в `README.md` проекта. +## Источник и версия +- Репозиторий демобазы: `postgrespro/demodb`. +- Закреплённый коммит: `d68de192850237719f09b47688d5f3fc94653ca6` (см. `DEMODB_COMMIT` в корневом `Makefile`). + +## Что мы патчим в demodb +- `install.sql`: `DROP DATABASE IF EXISTS demo WITH (FORCE)` — установка не падает, даже если демобазу держат активные сессии (например, из Airflow). +- `engine.sql`: два изменения в `engine_jobs1_sync.patch`: + - `busy()` игнорирует свой `pid`, чтобы не считать собственное подключение занятым; + - `continue()` при `jobs=1` вызывает `process_queue` синхронно (без `dblink`), иначе генерация обрывается при выходе из `psql` и данных не появляется. +- Патчи применяются автоматически в `make bookings-init`. Если что-то пошло не так, их можно накатить вручную: + ``` + patch -d bookings/demodb -p1 --forward < bookings/patches/install_drop_if_exists.patch + patch -d bookings/demodb -p1 --forward < bookings/patches/engine_jobs1_sync.patch + ``` + +## Быстрая проверка после init/обновления +- `make bookings-init` должен завершиться без ошибок; в `bookings.bookings` ожидаем >0 строк (примерно 15k). +- `make bookings-generate-day` добавляет следующий день после `max(book_date)`. +- Ручной вызов генерации из psql/DBeaver — только через DO-блок (подзапрос в аргументах `CALL` не работает): + ```sql + DO $$ + DECLARE + v_next_day timestamptz; + BEGIN + SELECT date_trunc('day', max(book_date)) + interval '1 day' + INTO v_next_day + FROM bookings.bookings; + CALL continue(v_next_day); -- или CALL continue(v_next_day, 4) + END $$; + ``` diff --git a/bookings/patches/engine_jobs1_sync.patch b/bookings/patches/engine_jobs1_sync.patch index 6c09c6c..d2faefb 100644 --- a/bookings/patches/engine_jobs1_sync.patch +++ b/bookings/patches/engine_jobs1_sync.patch @@ -1,6 +1,6 @@ --- a/engine.sql +++ b/engine.sql -@@ -230,7 +230,8 @@ +@@ SELECT count(*) > 0 FROM pg_stat_activity WHERE application_name = 'Airlines processor' @@ -8,4 +8,36 @@ + AND state != 'idle' + AND pid <> pg_backend_pid(); END; - + +@@ + TRUNCATE TABLE gen.stat_jobs; + + -- disconnect all previously opened connections + PERFORM dblink_disconnect(unnest(dblink_get_connections())); + +- -- start parallel jobs +- FOR i IN 1 .. jobs LOOP +- connname := 'job' || i; +- PERFORM dblink_connect(connname, current_setting('gen.connstr')); +- res := CASE dblink_send_query(connname, format('CALL process_queue(%L)',end_date)) +- WHEN 1 THEN 'ok' ELSE 'FAILED' +- END CASE; +- CALL log_message(0, format('Job %s (connname=%s): %s', i, connname, res)); +- RAISE NOTICE 'Starting job %: %', i, res; +- END LOOP; ++ IF jobs = 1 THEN ++ -- jobs=1: синхронно, чтобы не убивать dblink-сессию при выходе из psql ++ CALL process_queue(end_date); ++ CALL log_message(0, 'Job 1 (local): ok'); ++ ELSE ++ -- start parallel jobs ++ FOR i IN 1 .. jobs LOOP ++ connname := 'job' || i; ++ PERFORM dblink_connect(connname, current_setting('gen.connstr')); ++ res := CASE dblink_send_query(connname, format('CALL process_queue(%L)',end_date)) ++ WHEN 1 THEN 'ok' ELSE 'FAILED' ++ END CASE; ++ CALL log_message(0, format('Job %s (connname=%s): %s', i, connname, res)); ++ RAISE NOTICE 'Starting job %: %', i, res; ++ END LOOP; ++ END IF; diff --git a/bookings/patches/install_drop_if_exists.patch b/bookings/patches/install_drop_if_exists.patch index 8a75dbb..0582bfd 100644 --- a/bookings/patches/install_drop_if_exists.patch +++ b/bookings/patches/install_drop_if_exists.patch @@ -5,7 +5,7 @@ */ -DROP DATABASE demo; -+DROP DATABASE IF EXISTS demo; ++DROP DATABASE IF EXISTS demo WITH (FORCE); CREATE DATABASE demo; \c demo CREATE EXTENSION btree_gist; diff --git a/sql/stg/bookings_dq.sql b/sql/stg/bookings_dq.sql index 17f3e5c..62fbb73 100644 --- a/sql/stg/bookings_dq.sql +++ b/sql/stg/bookings_dq.sql @@ -23,6 +23,12 @@ BEGIN FROM stg.bookings_ext WHERE book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике bookings_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (make bookings-init / make bookings-generate-day или таск generate_bookings_day).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + END IF; + -- Считаем строки, реально вставленные в stg.bookings в этом батче SELECT COUNT(*) INTO v_stg_count