Модуль 4.3. Хранилища и интеграции с DWH/BI
Темы: OLTP vs OLAP, витрины, CDC, семантический слой, SLA данных. Артефакты: спецификация обмена с DWH/BI. Практика: описать контракт выгрузок/загрузок.
Как системный аналитик (SA) вы связываете продуктовые фичи с данными: от «источника правды» в OLTP до отчётов и ML-моделей в DWH/BI. От качества контрактов данных зависят корректность метрик, своевременность отчётов, надёжность витрин и стоимость владения.
В результате модуля у вас будут:
- чёткое разделение OLTP/OLAP-обязанностей;
- шаблоны CDC/ETL и требования к их спецификациям;
- семантический слой (метрики/измерения) в виде формальной спецификации;
- SLA/SLO данных и мониторинг свежести/полноты/целостности;
- артефакт «Спецификация обмена с DWH/BI» + инструкция по контрактам выгрузок/загрузок.
Картина мира: OLTP vs OLAP (как сотруднику — к исполнению)
OLTP (операционная БД)
- Назначение: транзакции, низкая задержка, строгие инварианты.
- Модель: нормализованная, много UPDATE/INSERT, небольшие запросы.
- «Источник правды» для: заказов, платежей, статусов, справочников.
OLAP / DWH (аналитика)
- Назначение: агрегаты, тренды, отчёты, ML.
- Модель: звезда/снежинка, большие сканы/агрегации, иммутабельность партиций.
- Слои: Raw/Staging → Core (модель домена) → Marts (витрины/BI).
- Разделяйте «время события» (event_time) и «время загрузки» (ingest_time).
Правило: OLTP для продуктовых консистентных операций; DWH — для аналитики/отчётности/витрин. Не строим тяжёлые отчёты на OLTP.
Архитектура DWH и витрины
Слои данных
- RAW/STAGING: «как пришло» (CDC/файлы/события), минимальная обработка; добавьте техполя: extract_ts, op, source_lsn, filename.
- CORE: приведение типов, дедуп, склейка инкрементов, SCD-история, бизнес-ключи, ссылки на справочники.
- MARTS (GOLD): звезда (факты + измерения) под задачи BI/продукта.
Моделирование витрин (звезда)
- Факты: fact_order, fact_payment, меры (суммы, количества), ключи времени/клиента/продукта.
- Измерения: dim_customer (SCD2), dim_product, dim_date, dim_geo.
-
SCD типы:
Type 1 — перезапись ошибок; Type 2 — история с valid_from/valid_to/is_current; Type 3 — редкие альтернативы.
Производительность витрин
- Партиционирование по дате события (event_date) и/или по ключам высокой кардинальности.
- Материализованные агрегаты для «тяжёлых» отчётов (день/неделя/месяц).
- Индексы/кластеризация по полям фильтра/сортировки; избегайте SELECT *.
CDC/ETL: как мы «несём правду» в DWH
Варианты CDC
- Логическая репликация/WAL (Debezium/Connect) — предпочтительно: малое влияние на OLTP, порядок per-key, тип операции (c/u/d), tombstone.
- Триггеры/аудит-таблицы — быстрый старт, но нагрузка и риск ошибок.
- Outbox + CDC — для доменных событий (см. модуль 3.5).
Обязательные элементы спецификации CDC
- Таблицы/топики источника; ключи/порядок (partition key).
- Семантика операций: op=c|u|d, «устранение» логических удалений (tombstone).
- Дедуп: правило уникальности (pk + op + ts), окно задержек.
- Watermark: поле и логика «догрузок» (replay/backfill).
- Эволюция схем: registry, совместимость (backward/full).
Базовые SQL-шаблоны (PostgreSQL/ANSI-like)
Дедуп и last-write-wins в STG→CORE:
WITH ranked AS (
SELECT t.*,
row_number() OVER (
PARTITION BY business_key
ORDER BY event_ts DESC, source_lsn DESC
) AS rn
FROM stg_order t
)
SELECT * FROM ranked WHERE rn = 1;
SCD2 для измерения клиента:
MERGE INTO core.dim_customer d USING stg_customer s ON (d.customer_id = s.customer_id AND d.is_current = true) WHEN MATCHED AND d.hash <> s.hash THEN UPDATE SET valid_to = s.valid_from - interval '1 millisecond', is_current = false WHEN NOT MATCHED BY TARGET THEN INSERT (customer_id, name, email, segment, hash, valid_from, valid_to, is_current) VALUES (s.customer_id, s.name, s.email, s.segment, s.hash, s.valid_from, '9999-12-31', true);
Контракты выгрузок/загрузок: форматы, SLA, ошибки
Варианты интеграции
- Стрим: Kafka/NATS (Avro/Proto), CDC/события, near-real-time.
- Файлы: S3/HDFS/FTP, Parquet/CSV (gzip), партиции dt=YYYY-MM-DD.
- API-пулл: REST/gRPC для BI-инструмента/витрин (редко, дороже).
Именование и партиционирование файлов
- Путь: s3://datalake/staging/<system>/<table>/dt=YYYY-MM-DD/hour=HH/part-0000.snappy.parquet
- Манифест: manifest.json с списком файлов, record_count, schema_version, watermark.
Пример manifest.json:
{
"table": "order",
"date": "2025-08-19",
"schema_version": "1.3.0",
"files": [
{"path": "part-0000.snappy.parquet", "rows": 500000, "crc32": "AB12CD34"}
],
"watermark": "2025-08-19T23:59:59Z",
"producer": "oltp-repl",
"op_distribution": {"c": 340000, "u": 140000, "d": 20000}
}
Схема и типы
- Денежные — DECIMAL(18,2) + currency CHAR(3).
- Даты — TIMESTAMP WITH TIME ZONE (timestamptz), в данных — UTC.
- Статусы — коды, валидируемые справочниками.
- Политика NULL/обязательность полей — явно в контракте.
Ошибки и повторы
- Повторная поставка файла/партии — допустима (идемпотентность загрузчика по file_id/crc32).
- Частично обработанные батчи — «атомарность» на партицию/файл.
- DLQ: ошибки формирования — в карантин-каталог + тикет.
Семантический слой (Semantic Layer): метрики и словарь
Зачем
Единые определения метрик = отсутствие «двух разных выручек». Семантический слой — описательная модель Measures / Dimensions / Entities / Relationships, которую используют BI-инструменты, SQL/DBT и продуктовые команды.
Состав семслоя
- Entities: Order, Payment, Customer.
- Dimensions: дата/гео/канал/сегмент.
- Measures: gmv, captured_amount, refund_rate.
- Rules: как фильтровать статусы, как считать возвраты/частичные платежи.
- Time grains: day/week/month, timezone политики.
Пример YAML-описания (фрагмент):
entities:
- name: order
primary_key: order_id
time: created_at
dimensions:
- name: status
values: [CREATED, PAID, SHIPPED, DELIVERED, CANCELLED]
- name: currency
measures:
- name: gmv
expr: "SUM(total_amount)"
filters: ["status IN ('PAID','SHIPPED','DELIVERED')"]
- name: orders_cnt
expr: "COUNT_DISTINCT(order_id)"
relationships:
- name: order_payment
from: order.order_id
to: payment.order_id
type: many_to_one
Качество определения метрик
- Любая мера имеет: формулу, фильтры, time grain, timezone, исключения/исправления, владельца и SLO точности.
SLA/SLO данных и наблюдаемость
SLI/метрики
- Freshness: lag между event_time (или extract_ts) и доступностью в MART (p95 ≤ 15 мин в 08:00–23:00).
- Completeness: доля записей по бизнес-ключу (пакет/день) ≥ 99.9%.
- Integrity: инварианты (суммы/валюты/FK), нарушения = 0.
- Accuracy: выборочный контроль (сквозные сверки OLTP↔DWH).
- Latency витрин: время пересчёта.
Алерты и гейты
- freshness_p95_minutes > 15 (SEV-2), dq_invalid_rate > 0.1% (SEV-2).
- «Гейт публикации витрины»: при критических нарушениях — стоп релиза.
Безопасность и комплаенс
- PII/PCI: минимизация полей, токенизация, шифрование в покое/полёте, маскировка при выдаче.
- RLS (row-level security): доступ к витринам по ролям/региону.
- Аудит: кто/когда изменил схему/правила; lineage (OpenLineage/каталог данных).
Практические примеры (домен «Заказы–Платежи–Возвраты»)
Таблицы CORE
-- Факт платежей CREATE TABLE core.fact_payment ( payment_id uuid PRIMARY KEY, order_id uuid NOT NULL, event_time timestamptz NOT NULL, status text NOT NULL, -- AUTHORIZED/CAPTURED/FAILED amount numeric(18,2) NOT NULL, currency char(3) NOT NULL ); -- Измерение клиентов (SCD2) CREATE TABLE core.dim_customer ( surrogate_key bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, customer_id uuid NOT NULL, email text, segment text, valid_from timestamptz NOT NULL, valid_to timestamptz NOT NULL, is_current boolean NOT NULL );
Витрина MART «Ежедневная выручка»
CREATE MATERIALIZED VIEW mart.daily_revenue AS
SELECT
date_trunc('day', p.event_time AT TIME ZONE 'UTC')::date AS d,
p.currency,
SUM(CASE WHEN p.status = 'CAPTURED' THEN p.amount ELSE 0 END) AS captured_amount,
SUM(CASE WHEN r.status = 'COMPLETED' THEN r.amount ELSE 0 END) AS refunded_amount,
SUM(CASE WHEN p.status = 'CAPTURED' THEN p.amount ELSE 0 END)
- COALESCE(SUM(CASE WHEN r.status='COMPLETED' THEN r.amount END),0) AS net_amount
FROM core.fact_payment p
LEFT JOIN core.fact_refund r ON r.payment_id = p.payment_id
GROUP BY 1,2;
DQ-проверки (ядро)
-- Валюта платежа совпадает с валютой заказа
SELECT p.payment_id
FROM core.fact_payment p
JOIN core.fact_order o USING(order_id)
WHERE p.currency <> o.currency;
-- Сумма возвратов ≤ сумма списаний
WITH agg AS (
SELECT payment_id,
SUM(CASE WHEN status='CAPTURED' THEN amount ELSE 0 END) AS captured,
SUM(CASE WHEN status='COMPLETED' THEN amount ELSE 0 END) AS refunded
FROM core.fact_payment
LEFT JOIN core.fact_refund USING(payment_id)
GROUP BY 1
)
SELECT payment_id FROM agg WHERE refunded > captured;
Типовые риски и как их гасить
|
Риск |
Проявление |
Меры |
|---|---|---|
|
Сдвиг схемы |
Падение пайплайна при новом поле |
Schema Registry + совместимость, «ленивые» колонки, контракт-ревью |
|
Поздние события |
«прыгающие» метрики, некорректные окна |
Watermark + допересчёт партиций за X дней, idempotent upsert |
|
Часовые пояса |
Дни/Mtd расходятся |
Везде UTC в данных, перестройки в «локали» на витринах |
|
PII-утечки |
Персональные данные в витринах/логах |
Маски/токены, минимизация, DLP-сканы |
|
Дубли |
Двойные суммы/счётчики |
Дедуп по бизнес-ключам + окно, уникальные constraints в CORE |
|
Нагрузка на OLTP |
Тормоза в продукте |
CDC через WAL, чтение с реплик, окна инкрементов |
|
Несход отчётов |
«две выручки» |
Семантический слой, патч-таблицы, согласованные правила фильтров |
Вопрос–Ответ
Q: Что лучше для свежести: CDC или ночные дампы?
A: CDC. Ночные дампы — только для бэкапов/массовых backfill. CDC даёт near-real-time и меньше нагрузку на OLTP.
Q: Как считать «выручку» при частичных платежах/возвратах?
A: Через семантический слой: captured_amount минус refunded_amount, статусные фильтры (CAPTURED, COMPLETED), таймзона отчёта — договорная.
Q: Можно ли брать данные напрямую из OLTP для BI?
A: Только в исключениях. Стандарт — витрины MART, где применены DQ, SCD и правила семслоя.
Q: Что делать с логическими удалениями?
A: Хранить op=d (tombstone) в STG и отражать в CORE (флаг is_deleted или удаление из фактов в рамках бизнес-правил).
Q: SCD2 не «раздует» измерения?
A: Даёт рост, зато сохраняет историю. Смягчается партицией по valid_from, компрессией и агрегациями.
Артефакт: «Спецификация обмена с DWH/BI» (шаблон)
# Spec: DWH/BI Exchange — <Домен/Система> vX.Y.Z Owner: <Команда DWH/Бизнес> | Contact: <email/Slack> | SA: <ФИО> ## 1. Цель и границы - Назначение: <отчёты/метрики/ML> - Источник правды: <таблицы/топики OLTP/события> ## 2. Источники и инкременты - Тип: CDC (Debezium) | Файлы (S3/Parquet) | События (Kafka/Avro) - Объекты: <order, payment, refund> - Ключи: <order_id/payment_id> - Порядок: per-key, поле event_ts - Watermark: <event_ts | extract_ts>, окно <N часов/дней> ## 3. Схема и типы (по объектам) - Поля: имя | тип | nullable | домен | пример | описание - Денежные: DECIMAL(18,2) + currency CHAR(3) - Даты: timestamptz (UTC) - Статусы: ref_<...> ## 4. Контракты поставок/загрузок - Частота: <каждые 5 мин / раз в час / ежедневно 08:00> - Формат: Parquet/Avro/JSON | сжатие - Партиции: dt=YYYY-MM-DD[/hour=HH] - Манифест: manifest.json (schema_version, watermark, files) - Идемпотентность: file_id/crc32; upsert по ключу ## 5. DQ/метрики качества - Freshness p95 ≤ <...> (окно …) - Completeness ≥ <...> - Integrity: инварианты (списком) - Uniqueness: дубли по email/order_id — 0 - Отчёт/дашборд: <ссылка> ## 6. Семантический слой - Entities/Dimensions/Measures (YAML-ссылка) - Определения ключевых мер: GMV, NetRevenue, ARPU (формулы/фильтры) - Timezone/гранулярность ## 7. SLA/SLO и алерты - SLO: <таблица SLO> - Алерты: правила, каналы, дежурства ## 8. Безопасность/Доступ - Роли/скимы/вьюхи (RLS) - Маскирование PII - Политика ретенции ## 9. Процессы - Backfill/replay: шаги, лимиты, окна - Версионирование схем (SemVer), deprecation policy - Incident response (классы инцидентов, RACI) ## 10. Тесты/Контроль - Набор SQL-проверок (ссылки) - Contract tests (schema-compat / consumer-driven) - Performance (SLI построения витрин) ## 11. Изменения/Changelog - vX.Y.Z — …
Практика: «Описать контракт выгрузок/загрузок» (90–120 мин)
Вход: домен «Заказы–Платежи–Возвраты» (или ваш).
Задачи:
- Заполнить артефакт «Спецификация обмена с DWH/BI» для трёх объектов: order, payment, refund.
- Определить каналы: CDC (orders/payments), файлы (refunds) — указать партиции/формат/манифест.
- Описать эволюцию схем: registry, политика совместимости, окно депрекейта.
- Задать SLO свежести/полноты и алерты.
- Подготовить 5 DQ-SQL (полнота, целостность сумм, валюта, уникальность idempotency, свежесть).
- Сформировать YAML семслоя с 3 мерами и 5 измерениями.
- Прописать backfill на 90 дней: лимиты, порядок, проверки.
Критерии зачёта
- Контракт однозначен: формат/схемы/водяные знаки/ошибки/идемпотентность.
- DQ и SLO измеримы и привязаны к метрикам.
- Семантический слой покрывает ключевые метрики.
- Безопасность/PII учтены; роли доступа заданы.
- Backfill/replay описаны и безопасны к повтору.
Шпаргалка (коротко)
- OLTP — «живая» транзакционная правда; DWH — история и аналитика.
- CDC предпочтительнее дампов; outbox — для событий домена.
- Слои: STAGING→CORE→MARTS; SCD2 для истории, меры/измерения — в семслое.
- Контракты данных = формат, схема, партиции, watermark, идемпотентность, SLA/SLO, DQ.
- Свежесть/полнота/целостность — главные SLI; без мониторинга и алертов «данных нет».
- PII/PCI — минимизируйте, шифруйте, ограничивайте доступ (RLS).



