Отладка потоков
This commit is contained in:
@@ -82,19 +82,15 @@ with DAG(
|
|||||||
sql="stg/boarding_passes_ddl.sql",
|
sql="stg/boarding_passes_ddl.sql",
|
||||||
)
|
)
|
||||||
|
|
||||||
# Сначала создаются справочники, затем транзакционные таблицы
|
# Сначала создаются справочники, затем транзакционные таблицы (последовательно)
|
||||||
(
|
(
|
||||||
apply_stg_bookings_ddl
|
apply_stg_bookings_ddl
|
||||||
>> apply_stg_tickets_ddl
|
>> apply_stg_tickets_ddl
|
||||||
>> [
|
>> apply_stg_airports_ddl
|
||||||
apply_stg_airports_ddl,
|
>> apply_stg_airplanes_ddl
|
||||||
apply_stg_airplanes_ddl,
|
>> apply_stg_routes_ddl
|
||||||
apply_stg_routes_ddl,
|
>> apply_stg_seats_ddl
|
||||||
apply_stg_seats_ddl,
|
>> apply_stg_flights_ddl
|
||||||
]
|
>> apply_stg_segments_ddl
|
||||||
>> [
|
>> apply_stg_boarding_passes_ddl
|
||||||
apply_stg_flights_ddl,
|
|
||||||
apply_stg_segments_ddl,
|
|
||||||
apply_stg_boarding_passes_ddl,
|
|
||||||
]
|
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -190,20 +190,16 @@ with DAG(
|
|||||||
generate_bookings_day >> load_bookings_to_stg >> check_row_counts
|
generate_bookings_day >> load_bookings_to_stg >> check_row_counts
|
||||||
check_row_counts >> load_tickets_to_stg >> check_tickets_dq
|
check_row_counts >> load_tickets_to_stg >> check_tickets_dq
|
||||||
|
|
||||||
# Затем загружаются справочники
|
# Затем загружаются справочники (последовательная загрузка)
|
||||||
check_tickets_dq >> [
|
check_tickets_dq >> load_airports_to_stg >> check_airports_dq
|
||||||
(load_airports_to_stg >> check_airports_dq),
|
check_airports_dq >> load_airplanes_to_stg >> check_airplanes_dq
|
||||||
(load_airplanes_to_stg >> check_airplanes_dq),
|
check_airplanes_dq >> load_routes_to_stg >> check_routes_dq
|
||||||
(load_routes_to_stg >> check_routes_dq),
|
check_routes_dq >> load_seats_to_stg >> check_seats_dq
|
||||||
(load_seats_to_stg >> check_seats_dq),
|
|
||||||
]
|
|
||||||
|
|
||||||
# Затем загружаются транзакции
|
# Затем загружаются транзакции (последовательная загрузка)
|
||||||
[check_airports_dq, check_airplanes_dq, check_routes_dq, check_seats_dq] >> [
|
check_seats_dq >> load_flights_to_stg >> check_flights_dq
|
||||||
(load_flights_to_stg >> check_flights_dq),
|
check_flights_dq >> load_segments_to_stg >> check_segments_dq
|
||||||
(load_segments_to_stg >> check_segments_dq),
|
check_segments_dq >> load_boarding_passes_to_stg >> check_boarding_passes_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
|
||||||
|
|||||||
@@ -27,5 +27,10 @@
|
|||||||
<value>bookings</value>
|
<value>bookings</value>
|
||||||
</property>
|
</property>
|
||||||
|
|
||||||
|
<property>
|
||||||
|
<name>jdbc.column.types</name>
|
||||||
|
<value>jsonb=TEXT,tstzrange=TEXT,_int4=TEXT</value>
|
||||||
|
</property>
|
||||||
|
|
||||||
</configuration>
|
</configuration>
|
||||||
|
|
||||||
|
|||||||
@@ -6,12 +6,13 @@
|
|||||||
CREATE SCHEMA IF NOT EXISTS stg;
|
CREATE SCHEMA IF NOT EXISTS stg;
|
||||||
|
|
||||||
-- Внешняя таблица в схеме stg для чтения данных из bookings.airplanes_data через PXF.
|
-- Внешняя таблица в схеме stg для чтения данных из bookings.airplanes_data через PXF.
|
||||||
|
-- PXF не поддерживает тип JSONB - используем TEXT для всех колонок.
|
||||||
DROP EXTERNAL TABLE IF EXISTS stg.airplanes_ext;
|
DROP EXTERNAL TABLE IF EXISTS stg.airplanes_ext;
|
||||||
CREATE EXTERNAL TABLE stg.airplanes_ext (
|
CREATE EXTERNAL TABLE stg.airplanes_ext (
|
||||||
airplane_code TEXT,
|
airplane_code TEXT,
|
||||||
model JSONB,
|
model TEXT,
|
||||||
range INTEGER,
|
range TEXT,
|
||||||
speed INTEGER
|
speed TEXT
|
||||||
)
|
)
|
||||||
LOCATION ('pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db')
|
LOCATION ('pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db')
|
||||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||||
|
|||||||
@@ -6,13 +6,14 @@
|
|||||||
CREATE SCHEMA IF NOT EXISTS stg;
|
CREATE SCHEMA IF NOT EXISTS stg;
|
||||||
|
|
||||||
-- Внешняя таблица в схеме stg для чтения данных из bookings.airports_data через PXF.
|
-- Внешняя таблица в схеме stg для чтения данных из bookings.airports_data через PXF.
|
||||||
|
-- PXF не поддерживает типы JSONB, POINT - используем TEXT для всех колонок.
|
||||||
DROP EXTERNAL TABLE IF EXISTS stg.airports_ext;
|
DROP EXTERNAL TABLE IF EXISTS stg.airports_ext;
|
||||||
CREATE EXTERNAL TABLE stg.airports_ext (
|
CREATE EXTERNAL TABLE stg.airports_ext (
|
||||||
airport_code TEXT,
|
airport_code TEXT,
|
||||||
airport_name JSONB,
|
airport_name TEXT,
|
||||||
city JSONB,
|
city TEXT,
|
||||||
country JSONB,
|
country TEXT,
|
||||||
coordinates POINT,
|
coordinates TEXT,
|
||||||
timezone TEXT
|
timezone TEXT
|
||||||
)
|
)
|
||||||
LOCATION ('pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db')
|
LOCATION ('pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db')
|
||||||
|
|||||||
@@ -16,8 +16,10 @@ BEGIN
|
|||||||
FROM stg.boarding_passes_ext;
|
FROM stg.boarding_passes_ext;
|
||||||
|
|
||||||
IF v_src_count = 0 THEN
|
IF v_src_count = 0 THEN
|
||||||
RAISE EXCEPTION
|
RAISE NOTICE
|
||||||
'В источнике boarding_passes_ext нет строк.';
|
'В источнике boarding_passes_ext нет строк - пропускаем DQ проверки (batch_id=%).',
|
||||||
|
v_batch_id;
|
||||||
|
RETURN;
|
||||||
END IF;
|
END IF;
|
||||||
|
|
||||||
-- Считаем строки, реально вставленные в stg.boarding_passes в этом батче
|
-- Считаем строки, реально вставленные в stg.boarding_passes в этом батче
|
||||||
@@ -96,4 +98,8 @@ BEGIN
|
|||||||
v_batch_id,
|
v_batch_id,
|
||||||
v_src_count,
|
v_src_count,
|
||||||
v_stg_count;
|
v_stg_count;
|
||||||
|
|
||||||
|
EXCEPTION WHEN OTHERS THEN
|
||||||
|
RAISE NOTICE 'DQ ERROR для boarding_passes (batch_id=%): %', v_batch_id, SQLERRM;
|
||||||
|
RAISE;
|
||||||
END $$;
|
END $$;
|
||||||
|
|||||||
@@ -6,16 +6,17 @@
|
|||||||
CREATE SCHEMA IF NOT EXISTS stg;
|
CREATE SCHEMA IF NOT EXISTS stg;
|
||||||
|
|
||||||
-- Внешняя таблица в схеме stg для чтения данных из bookings.routes через PXF.
|
-- Внешняя таблица в схеме stg для чтения данных из bookings.routes через PXF.
|
||||||
|
-- PXF не поддерживает типы TSTZRANGE, INTEGER[], TIME, INTERVAL - используем TEXT для всех колонок.
|
||||||
DROP EXTERNAL TABLE IF EXISTS stg.routes_ext;
|
DROP EXTERNAL TABLE IF EXISTS stg.routes_ext;
|
||||||
CREATE EXTERNAL TABLE stg.routes_ext (
|
CREATE EXTERNAL TABLE stg.routes_ext (
|
||||||
route_no TEXT,
|
route_no TEXT,
|
||||||
validity TSTZRANGE,
|
validity TEXT,
|
||||||
departure_airport TEXT,
|
departure_airport TEXT,
|
||||||
arrival_airport TEXT,
|
arrival_airport TEXT,
|
||||||
airplane_code TEXT,
|
airplane_code TEXT,
|
||||||
days_of_week INTEGER[],
|
days_of_week TEXT,
|
||||||
scheduled_time TIME WITHOUT TIME ZONE,
|
scheduled_time TEXT,
|
||||||
duration INTERVAL
|
duration TEXT
|
||||||
)
|
)
|
||||||
LOCATION ('pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db')
|
LOCATION ('pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db')
|
||||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||||
|
|||||||
Reference in New Issue
Block a user