Как грамотно реализовать перевод данных из STG в ODS при помощи SCD2 в PostgreSQL?
Ссылка скопирована
Сейчас ситуация такая: есть задание: нужно сделать task в Airflow, чтобы данные из БД STG переносились в БД ODS с реализацией SCD2. Хочу сделать универсальную процедуру, которую мы будем вызывать в нашем dag. Проблема с написанием самой функции. Нужно учесть, что таблицы имеют разное количество столбцов и разные названия.
По вводным: как основа у меня есть процедура для одной таблицы:
CREATE TABLE товары_scd2 ( id SERIAL PRIMARY KEY, название TEXT, код_товара TEXT NOT NULL, группа_товаров_id INT NOT NULL, бренд_id INT NOT NULL, линейка_товара_id INT NOT NULL, производитель_id INT NOT NULL, дата_начала DATE NOT NULL, дата_окончания DATE, признак_актуальности BOOLEAN NOT NULL ); CREATE OR REPLACE PROCEDURE update_товары_scd2() LANGUAGE plpgsql AS $$ DECLARE updated_old_versions INT := 0; inserted_new_versions INT := 0; inserted_new_products INT := 0; BEGIN DROP TABLE IF EXISTS tmp_source_data; CREATE TEMP TABLE tmp_source_data AS SELECT t.Товар AS название, t.КодТовара AS код_товара, g.id AS группа_товаров_id, b.id AS бренд_id, l.id AS линейка_товара_id, p.id AS производитель_id FROM dblink( 'host=localhost port=5432 dbname=postgres user=postgres password=1', 'SELECT DISTINCT Товар, КодТовара, ГруппаТовара, Бренд, ЛинейкаТовара, Производитель FROM товары' ) AS t( Товар VARCHAR, КодТовара VARCHAR, ГруппаТовара VARCHAR, Бренд VARCHAR, ЛинейкаТовара VARCHAR, Производитель VARCHAR ) JOIN бренды b ON b.название = t.Бренд JOIN группы_товаров g ON g.группа = t.ГруппаТовара AND g.бренд_id = b.id JOIN линейки_товара l ON l.название = t.ЛинейкаТовара JOIN производители p ON p.имя_производителя = t.Производитель WHERE t.Товар IS NOT NULL AND t.КодТовара IS NOT NULL AND t.ГруппаТовара IS NOT NULL AND t.Бренд IS NOT NULL AND t.ЛинейкаТовара IS NOT NULL AND t.Производитель IS NOT NULL; -- Закрываем старые версии UPDATE товары_scd2 t SET дата_окончания = CURRENT_DATE - 1, признак_актуальности = FALSE WHERE t.признак_актуальности = TRUE AND EXISTS ( SELECT 1 FROM tmp_source_data s WHERE s.код_товара = t.код_товара AND s.группа_товаров_id = t.группа_товаров_id AND ( s.название IS DISTINCT FROM t.название OR s.бренд_id IS DISTINCT FROM t.бренд_id OR s.линейка_товара_id IS DISTINCT FROM t.линейка_товара_id OR s.производитель_id IS DISTINCT FROM t.производитель_id ) ); GET DIAGNOSTICS updated_old_versions = ROW_COUNT; RAISE NOTICE 'Закрыто устаревших версий: %', updated_old_versions; -- Вставка новых версий INSERT INTO товары_scd2 ( название, код_товара, группа_товаров_id, бренд_id, линейка_товара_id, производитель_id, дата_начала, дата_окончания, признак_актуальности ) SELECT s.название, s.код_товара, s.группа_товаров_id, s.бренд_id, s.линейка_товара_id, s.производитель_id, CURRENT_DATE, NULL, TRUE FROM tmp_source_data s JOIN товары_scd2 t ON t.код_товара = s.код_товара AND t.группа_товаров_id = s.группа_товаров_id WHERE t.признак_актуальности = FALSE AND NOT EXISTS ( SELECT 1 FROM товары_scd2 t2 WHERE t2.код_товара = s.код_товара AND t2.группа_товаров_id = s.группа_товаров_id AND t2.название = s.название AND t2.бренд_id = s.бренд_id AND t2.линейка_товара_id = s.линейка_товара_id AND t2.производитель_id = s.производитель_id AND t2.признак_актуальности = TRUE ); GET DIAGNOSTICS inserted_new_versions = ROW_COUNT; RAISE NOTICE 'Добавлено новых версий: %', inserted_new_versions; -- Вставка новых товаров INSERT INTO товары_scd2 ( название, код_товара, группа_товаров_id, бренд_id, линейка_товара_id, производитель_id, дата_начала, дата_окончания, признак_актуальности ) SELECT s.название, s.код_товара, s.группа_товаров_id, s.бренд_id, s.линейка_товара_id, s.производитель_id, CURRENT_DATE, NULL, TRUE FROM tmp_source_data s LEFT JOIN товары_scd2 t ON s.код_товара = t.код_товара AND s.группа_товаров_id = t.группа_товаров_id WHERE t.id IS NULL; GET DIAGNOSTICS inserted_new_products = ROW_COUNT; RAISE NOTICE 'Добавлено новых товаров: %', inserted_new_products; RAISE NOTICE 'SCD2: Обновление завершено.'; END; $$;
Нужно решить такую задачу?
Заказать помощь
Опишите проблему, и специалист поможет с настройкой, исправлением ошибки или доработкой сайта. Подберём понятный план работ без лишней переписки.
Лучший ответ
1
Другие ответы (0)
Пока нет других ответов. Будьте первым, кто поможет автору.
Ответить на вопроскомментарий
Вам также может быть интересно
VPN
Как правильно настроить vless для Android TV?
1 ответ
Pyrogram
Как правильно зарегистрировать юзер бота в Telegram?
1 ответ
печатные-платы
Как заставить запускаться программу M3.exe от компании Hanxing AOI в инспекционной машине на Windows 7 Pro?
1 ответ
woocommerce
Как отображать при фильтрации товары со значением search параметра visibility: HUSKY — Products Filter for WooCommerce?
1 ответ

Для переноса STG -> ODS с SCD2 в PostgreSQL лучше делать не “универсальную магическую процедуру на всё”, а общий шаблон, в который явно передаются: таблица-источник, таблица-приёмник, бизнес-ключ, список сравниваемых атрибутов и дата загрузки. SCD2 требует аккуратности: вы не обновляете старую строку данными, а закрываете её период действия и вставляете новую версию.
Минимальная схема ODS-таблицы обычно такая:
business_key attr_1 attr_2 valid_from valid_to is_current hash_diff load_dttm
Удобнее сравнивать не каждое поле отдельно, а hash_diff по значимым атрибутам:
md5(concat_ws('|', attr_1, attr_2, attr_3)) AS hash_diff
Логика загрузки:
Примерно так закрывается старая версия:
UPDATE ods_table ods SET valid_to = :load_dttm - interval '1 second', is_current = false FROM stg_prepared stg WHERE ods.business_key = stg.business_key AND ods.is_current = true AND ods.hash_diff stg.hash_diff;
После этого вставляются новые и изменившиеся:
INSERT INTO ods_table (business_key, attr_1, attr_2, valid_from, valid_to, is_current, hash_diff, load_dttm) SELECT stg.business_key, stg.attr_1, stg.attr_2, :load_dttm, null, true, stg.hash_diff, now() FROM stg_prepared stg LEFT JOIN ods_table ods ON ods.business_key = stg.business_key AND ods.is_current = true WHERE ods.business_key IS NULL OR ods.hash_diff stg.hash_diff;
В Airflow DAG лучше делать task на каждую сущность, но вызывать одну процедуру/SQL-шаблон с параметрами. Так вы сохраните универсальность, но не потеряете контроль над ключами и колонками. Обязательно добавьте уникальный индекс на текущую версию business_key, иначе при сбоях легко получить две актуальные строки.