From 7312cdfb68db51caa45775c0e72b87a76764e9ec Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 18 Jan 2026 13:47:09 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9E=D1=82=D0=BB=D0=B0=D0=B4=D0=BA=D0=B0=20?= =?UTF-8?q?=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow/dags/bookings_stg_ddl.py | 20 ++++++++------------ airflow/dags/bookings_to_gp_stage.py | 24 ++++++++++-------------- pxf/servers/bookings-db/jdbc-site.xml | 5 +++++ sql/stg/airplanes_ddl.sql | 7 ++++--- sql/stg/airports_ddl.sql | 9 +++++---- sql/stg/boarding_passes_dq.sql | 10 ++++++++-- sql/stg/routes_ddl.sql | 9 +++++---- 7 files changed, 45 insertions(+), 39 deletions(-) diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py index f205e79..874ec16 100644 --- a/airflow/dags/bookings_stg_ddl.py +++ b/airflow/dags/bookings_stg_ddl.py @@ -82,19 +82,15 @@ with DAG( sql="stg/boarding_passes_ddl.sql", ) - # Сначала создаются справочники, затем транзакционные таблицы + # Сначала создаются справочники, затем транзакционные таблицы (последовательно) ( apply_stg_bookings_ddl >> apply_stg_tickets_ddl - >> [ - apply_stg_airports_ddl, - apply_stg_airplanes_ddl, - apply_stg_routes_ddl, - apply_stg_seats_ddl, - ] - >> [ - apply_stg_flights_ddl, - apply_stg_segments_ddl, - apply_stg_boarding_passes_ddl, - ] + >> apply_stg_airports_ddl + >> apply_stg_airplanes_ddl + >> apply_stg_routes_ddl + >> apply_stg_seats_ddl + >> apply_stg_flights_ddl + >> apply_stg_segments_ddl + >> apply_stg_boarding_passes_ddl ) diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index d978de8..354f241 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -190,20 +190,16 @@ with DAG( generate_bookings_day >> load_bookings_to_stg >> check_row_counts check_row_counts >> load_tickets_to_stg >> check_tickets_dq - # Затем загружаются справочники - check_tickets_dq >> [ - (load_airports_to_stg >> check_airports_dq), - (load_airplanes_to_stg >> check_airplanes_dq), - (load_routes_to_stg >> check_routes_dq), - (load_seats_to_stg >> check_seats_dq), - ] + # Затем загружаются справочники (последовательная загрузка) + check_tickets_dq >> load_airports_to_stg >> check_airports_dq + check_airports_dq >> load_airplanes_to_stg >> check_airplanes_dq + check_airplanes_dq >> load_routes_to_stg >> check_routes_dq + check_routes_dq >> load_seats_to_stg >> check_seats_dq - # Затем загружаются транзакции - [check_airports_dq, check_airplanes_dq, check_routes_dq, check_seats_dq] >> [ - (load_flights_to_stg >> check_flights_dq), - (load_segments_to_stg >> check_segments_dq), - (load_boarding_passes_to_stg >> check_boarding_passes_dq), - ] + # Затем загружаются транзакции (последовательная загрузка) + check_seats_dq >> load_flights_to_stg >> check_flights_dq + check_flights_dq >> load_segments_to_stg >> check_segments_dq + check_segments_dq >> load_boarding_passes_to_stg >> check_boarding_passes_dq # В конце финальный лог - [check_flights_dq, check_segments_dq, check_boarding_passes_dq] >> finish_summary + check_boarding_passes_dq >> finish_summary diff --git a/pxf/servers/bookings-db/jdbc-site.xml b/pxf/servers/bookings-db/jdbc-site.xml index dca82e9..6172ce5 100644 --- a/pxf/servers/bookings-db/jdbc-site.xml +++ b/pxf/servers/bookings-db/jdbc-site.xml @@ -27,5 +27,10 @@ bookings + + jdbc.column.types + jsonb=TEXT,tstzrange=TEXT,_int4=TEXT + + diff --git a/sql/stg/airplanes_ddl.sql b/sql/stg/airplanes_ddl.sql index d29196d..c25439d 100644 --- a/sql/stg/airplanes_ddl.sql +++ b/sql/stg/airplanes_ddl.sql @@ -6,12 +6,13 @@ CREATE SCHEMA IF NOT EXISTS stg; -- Внешняя таблица в схеме stg для чтения данных из bookings.airplanes_data через PXF. +-- PXF не поддерживает тип JSONB - используем TEXT для всех колонок. DROP EXTERNAL TABLE IF EXISTS stg.airplanes_ext; CREATE EXTERNAL TABLE stg.airplanes_ext ( airplane_code TEXT, - model JSONB, - range INTEGER, - speed INTEGER + model TEXT, + range TEXT, + speed TEXT ) LOCATION ('pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db') FORMAT 'CUSTOM' (formatter='pxfwritable_import'); diff --git a/sql/stg/airports_ddl.sql b/sql/stg/airports_ddl.sql index 3f27dd2..f953258 100644 --- a/sql/stg/airports_ddl.sql +++ b/sql/stg/airports_ddl.sql @@ -6,13 +6,14 @@ CREATE SCHEMA IF NOT EXISTS stg; -- Внешняя таблица в схеме stg для чтения данных из bookings.airports_data через PXF. +-- PXF не поддерживает типы JSONB, POINT - используем TEXT для всех колонок. DROP EXTERNAL TABLE IF EXISTS stg.airports_ext; CREATE EXTERNAL TABLE stg.airports_ext ( airport_code TEXT, - airport_name JSONB, - city JSONB, - country JSONB, - coordinates POINT, + airport_name TEXT, + city TEXT, + country TEXT, + coordinates TEXT, timezone TEXT ) LOCATION ('pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db') diff --git a/sql/stg/boarding_passes_dq.sql b/sql/stg/boarding_passes_dq.sql index a19e710..0a2f123 100644 --- a/sql/stg/boarding_passes_dq.sql +++ b/sql/stg/boarding_passes_dq.sql @@ -16,8 +16,10 @@ BEGIN FROM stg.boarding_passes_ext; IF v_src_count = 0 THEN - RAISE EXCEPTION - 'В источнике boarding_passes_ext нет строк.'; + RAISE NOTICE + 'В источнике boarding_passes_ext нет строк - пропускаем DQ проверки (batch_id=%).', + v_batch_id; + RETURN; END IF; -- Считаем строки, реально вставленные в stg.boarding_passes в этом батче @@ -96,4 +98,8 @@ BEGIN v_batch_id, v_src_count, v_stg_count; + +EXCEPTION WHEN OTHERS THEN + RAISE NOTICE 'DQ ERROR для boarding_passes (batch_id=%): %', v_batch_id, SQLERRM; + RAISE; END $$; diff --git a/sql/stg/routes_ddl.sql b/sql/stg/routes_ddl.sql index c3c93c4..dca3f56 100644 --- a/sql/stg/routes_ddl.sql +++ b/sql/stg/routes_ddl.sql @@ -6,16 +6,17 @@ CREATE SCHEMA IF NOT EXISTS stg; -- Внешняя таблица в схеме stg для чтения данных из bookings.routes через PXF. +-- PXF не поддерживает типы TSTZRANGE, INTEGER[], TIME, INTERVAL - используем TEXT для всех колонок. DROP EXTERNAL TABLE IF EXISTS stg.routes_ext; CREATE EXTERNAL TABLE stg.routes_ext ( route_no TEXT, - validity TSTZRANGE, + validity TEXT, departure_airport TEXT, arrival_airport TEXT, airplane_code TEXT, - days_of_week INTEGER[], - scheduled_time TIME WITHOUT TIME ZONE, - duration INTERVAL + days_of_week TEXT, + scheduled_time TEXT, + duration TEXT ) LOCATION ('pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db') FORMAT 'CUSTOM' (formatter='pxfwritable_import');