Модуль 10.3. Системный аналитик DWH/BI/AI, Data-intensive/Analytics
Темы: витрины, метрики, lineage, DQ, CDC, API для выгрузок, роль в ML/AI-проектах. Артефакт: спецификация обмена DWH↔продукт. Практика: подготовить контракт DWH↔продукт (batch/stream).
В «данных-тяжёлых» продуктах системный аналитик (SA) отвечает за контракт между операционными сервисами и аналитической платформой: что, когда и с какой гарантией попадает в DWH/BI/AI, как это измеряется и как эволюционирует без поломок. Этот модуль — рабочая инструкция: от CDC и витрин до семантического слоя, от DQ-правил до API выгрузок и роли SA в ML/AI.
Роль SA в Data-инженерии и аналитике
Зоны ответственности:
- Контракты данных: схемы, версии, ключи дедупликации, SLA свежести и полноты.
- CDC/интеграции: выбор паттерна (log-based/triggers/timestamps), гарантии доставки, порядок.
- Семантический слой: словарь метрик/измерений, инварианты, соответствие бизнесу (UL).
- DQ и наблюдаемость: правила качества (validity/completeness/consistency/uniqueness/timeliness/accuracy), мониторинги, алерты.
- Lineage: прослеживаемость источников→трансформаций→витрин→отчётов/ML-фич.
- Безопасность данных: PII/финданные — маскирование, RLS/CLS, аудит, ретенция.
- ML/AI: формулировка задачи, выбор целевых метрик (контрфакты/бизнес-метрики), контракт фич/лейблов, синхрон онлайн/офлайн.
Архитектура данных (ориентир)
- Слои хранения: Raw (Bronze) → Cleansed (Silver) → Business (Gold) → Semantic (Metrics/BI).
- Модели: Star-schema (факты/измерения), Data Vault 2.0 (Hubs/Links/Sats) — выбираем по домену и скорости изменений.
- Потоки: batch (D+1) и near-real-time (стрим через шину).
- Медальоны + Mesh: доменные «data products» с собственными SLA и владельцами.
CDC (Change Data Capture): что выбрать и как описать
Паттерны
- Log-based (binlog/redo, Debezium и пр.): высоконадежно, все операции I/U/D с порядком и транзакциями.
- Timestamp-based: «где updated_at > watermark» — просто, но есть дырки при clock skew.
- Trigger-based: универсально, но нагружает OLTP.
- Snapshot + Incremental: начальный снимок → инкременталь.
Требования к CDC-контракту
- Ключи: businessKey, surrogateKey, dedupeKey = (pk, op_ts); идемпотентность.
- Порядок: op_ts, event_sequence, гарантия «не позже N минут».
- Tombstones для delete.
- Схемы: эволюция «additive-only» для downstream; MAJOR при breaking.
- SLA: freshness (например, ≤ 15 мин p95), completeness (≥ 99.5% за окно).
Витрины (модели потребления)
Факт/измерения
- Факты (Tx/Periodic/Snapshot): «платежи», «сеансы», «выдачи кредитов».
-
Измерения (SCD):
- Type-1: перезапись (актуальное).
- Type-2: историзация (valid_from/to, is_current).
- Type-3: короткая история (prev_value).
- Голы и анти-факты: «анти-заказы», сторно — аккуратно отражайте в моделях.
Семантический слой и метрики
- Единый словарь: metric_id, формула, агрегация по времени, фильтры включения/исключения, join-ключи.
- Валидации метрик: тест-наборы «ручных расчётов», инварианты (например, GMV ≥ NetRevenue).
- Версионирование: metrics.yml в Git; «метрика v2» при изменении определения.
Data Quality (DQ): правила, метрики, алерты
Классы DQ:
- Validity: формат, домены, справочники.
- Completeness: доля non-null, обязательные атрибуты.
- Uniqueness: ключи без дублей.
- Consistency: межтабличные связи, баланс.
- Timeliness: лаг между источником и витриной.
- Accuracy: сверка с внешним эталоном (PSP реестры, бухгалтерия).
Практика: для каждой таблицы — dq_checks.yml с порогами и действиями (fail/alert/quarantine). Линия жизни — «сырьё валидируем/кварентируем, бизнес-слой — не пропускаем брак».
Lineage и наблюдаемость данных
- Технический lineage: граф «источники → джобы (ETL/ELT) → таблицы → отчёты/ML».
- Бизнес-lineage: «метрика X = сумма из … с фильтром …».
- Сигналы наблюдаемости: лаг, размер партиций, кардинальность, частота джоб, ошибки парсера схем.
- Артефакты: карточки датасетов с владельцем, SLA, схемой, DQ-правилами, связями.
API для выгрузок/доступа к данным
Паттерны доступа
- SQL-endpoint / Federated queries (для внутренних сервисов/BI).
- REST/gRPC (агрегаты, отбор по ключам/диапазонам).
- GraphQL (выбор полей, вложенные выборки с лимитами).
- Асинхронные экспорт-задачи: POST /exports → job → pre-signed URL (CSV/Parquet).
Обязательные требования
- Пагинация/курсоры, rate-limits, idempotency (Idempotency-Key).
- Фильтры: время (from/to), бизнес-ключи, статусы.
- Форматы: CSV (схема/разделители), Parquet/ORC (колоночные), JSONL.
- Безопасность: ABAC/скоупы, RLS (строчная безопасность), маскирование PII, аудит.
Роль SA в ML/AI-проектах (без «магии»)
- Формулировка задачи: таргет/переменные, сценарий принятия решений, онлайн/офлайн.
- Метрики успеха: бизнес-метрики (uplift, конверсия), качество модели (AUC/PR-AUC), затраты/риск.
- Данные: контракт фич/лейблов (частота, лаги, источники, ключи join, SCD), синхронизм offline/online.
- Feature Store: дефиниции фич (время округления, агрегации, fill-политики), тесты «тренировка=прод».
- Релизы: shadow/AB, canary, мониторинг drift (data/model), алерты, откат.
- Комплаенс/этика: чувствительные признаки, аудит объяснимости, ретенция и право на удаление.
Артефакт: спецификация обмена DWH↔продукт (шаблон)
# data_contract.yml
data_product: "payments_analytics"
owner: "Data Platform Team"
sa: "ФИО"
version: "1.4.0"
sla:
freshness_p95: "15m"
completeness_day: ">=99.5%"
availability_month: "99.9%"
security:
classification: "PII:low" # none/low/medium/high
rls: ["tenant_id"]
pii_masking: ["email", "phone"]
source_streams:
- name: "payments_cdc"
type: "kafka"
topic: "db.payments.v1"
key: ["payment_id"]
ordering_key: ["payment_id"]
schema_ref: "schemas/payments.avsc"
cdc: {mode: "log-based", deletes: "tombstone"}
dedupe_key: ["payment_id","op_ts"]
- name: "refunds_cdc"
type: "kafka"
topic: "db.refunds.v1"
targets:
- table: "bronze.payments_raw"
- table: "silver.payments_clean"
- table: "gold.fact_payments"
partition_by: ["event_date"]
primary_key: ["payment_id"]
watermarks:
late_data: "24h"
drop_late_policy: "quarantine" # accept/quarantine/drop
metrics:
- id: "gmv"
layer: "semantic"
formula: "sum(amount) FILTER (where status in ('CAPTURED'))"
grain: "day, tenant_id, currency"
scd_policy: "type1"
- id: "conversion_to_payment"
formula: "payments_count / sessions_count"
dq_checks:
- table: "silver.payments_clean"
rules:
- id: "not_null_payment_id"
expr: "payment_id is not null"
threshold: "100%"
action: "fail"
- id: "valid_currency"
expr: "currency in ('RUB','USD','EUR')"
threshold: "99.9%"
action: "alert"
exports_api:
- endpoint: "/v1/exports/payments"
method: "POST"
request:
filters: ["date_from","date_to","tenant_id","status"]
format: ["CSV","PARQUET"]
columns: "optional"
response:
"202": "job accepted"
"download_url": "pre-signed, TTL=24h"
limits:
max_rows: 5_000_000
rps: 5
lineage:
inputs: ["db.payments", "db.orders"]
transformations: ["sql/gold_fact_payments.sql"]
outputs: ["bi.dashboard.payments", "ml.feature_store.payments_features"]
change_policy:
versioning: "semver"
non_breaking: ["add_column_nullable","add_metric"]
breaking: ["rename_drop_column","change_type"]
deprecation: "announce>=30d; support 90d"
Примеры и шаблоны
CDC таблица payments (SQL — Silver→Gold)
-- Идемпотентная загрузка в факт
MERGE INTO gold.fact_payments AS t
USING (
SELECT
payment_id,
invoice_id,
tenant_id,
amount::decimal(18,2),
currency,
status,
event_time::timestamp as event_time,
date_trunc('day', event_time) as event_date
FROM silver.payments_clean
QUALIFY ROW_NUMBER() OVER (PARTITION BY payment_id ORDER BY op_ts DESC) = 1
) s
ON t.payment_id = s.payment_id
WHEN MATCHED AND t.hash <> hash(s.*) THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...columns...) VALUES (...);
REST экспорт (OpenAPI фрагмент)
paths:
/v1/exports/payments:
post:
summary: Export payments dataset
parameters:
- in: header
name: Idempotency-Key
required: true
schema: {type: string, maxLength: 128}
requestBody:
required: true
content:
application/json:
schema:
type: object
required: [date_from, date_to, format]
properties:
date_from: {type: string, format: date}
date_to: {type: string, format: date}
tenant_id: {type: string}
status: {type: string}
format: {type: string, enum: [CSV, PARQUET]}
responses:
"202": {description: Accepted; returns jobId}
"200": {description: Ready; returns download_url}
"413": {description: Too large}
"429": {description: Rate limit}
Каталог метрик (фрагмент)
|
id |
Наименование |
Формула |
Гранулярность |
Фильтры |
Владелец |
Версия |
|---|---|---|---|---|---|---|
|
gmv |
GTV/оборот |
∑ amount where status='CAPTURED' |
день/магазин/валюта |
исключить тест-платежи |
Finance |
2.1 |
|
net_revenue |
Чистая выручка |
gmv − refunds − fees |
день/магазин |
— |
Finance |
1.4 |
Риски и анти-паттерны
|
Риск |
Симптом |
Что делать |
|---|---|---|
|
Сломанные выгрузки при «тихой» смене схемы |
ETL упал ночью |
Data contract + semver + CCB, тесты схем |
|
Дубликаты из CDC |
Скачет подсчёт, завышен GMV |
Идемпотентные MERGE, dedupeKey, ROW_NUMBER() |
|
Несвоевременные данные |
Дэшборды «вчерашние» |
SLA freshness, водяные знаки (watermarks), алерты лагов |
|
Конфликты метрик |
BI vs отчёт бухгалтера |
Единый словарь метрик, версии, бизнес-lineage |
|
Протечки PII |
Публичные выгрузки с e-mail |
Классификация, RLS/CLS, маскирование, аудит |
|
Разный offline/online для ML |
«Сдвиг» и деградация модели |
Feature store, одинаковые трансформации, тест «offline=online» |
|
Отсутствие историзации |
Пересчёты «поплыли» |
SCD2, даты действия тарифов/ставок |
|
Латентные аномалии DQ |
Ошибки у источника не видны |
DQ-правила на Bronze/Silver, quarantine-слой |
Практика (90–150 мин): спецификация обмена DWH↔продукт
Задание: подготовьте контракт для домена «Платежи и возвраты» (batch + stream).
Включите:
- data_contract.yml (по шаблону §9): источники (CDC топики), цели (таблицы), SLA, безопасность, правила DQ, export API.
- Схемы: Avro/JSON Schema событий payments_cdc, refunds_cdc (tombstone для delete).
- SQL-трансформации: Bronze→Silver (очистка/типизация), Silver→Gold (MERGE с дедупом).
- Каталог метрик: gmv, refund_rate, auth_to_capture_ratio с формулами и владельцами.
- Lineage-карточку: входы/выходы/скрипты/владельцы, ссылки на дашборды/фичи ML.
Критерии зачёта:
- Контракт полон: ключи, SLA, версия, DQ, безопасность.
- Схемы событий покрывают операции I/U/D и эволюцию.
- MERGE идемпотентен, есть дедуп и партиции.
- Метрики воспроизводимы и задокументированы.
- Lineage прозрачен (от источника до BI/ML).
Вопрос–Ответ
В: Что выбрать — Star или Data Vault?
О: Для стабильных бизнес-процессов и BI — Star. Для разношёрстных источников/частых изменений — Data Vault (история и источники «из коробки»). Часто делают Vault→Star.
В: Как считать SLA свежести?
О: freshness = now() − max(event_time) на Gold/семантическом уровне. Считайте перцентили (p50/p95), держите алерты и водяные знаки.
В: Допустимо ли изменять тип колонки в Gold?
О: Это breaking. Делайте additive (новая колонка + backfill), затем депрекейт старую.
В: Что делать с поздними событиями (late data)?
О: Политика: accept (перерасчёт партиции), quarantine (в отдельную таблицу), или drop — явно в контракте. Храните watermark.
В: Можно ли отдавать BI-выгрузки «как есть» наружу?
О: Только через авторизованный API, с RLS/CLS и маскированием PII; предпочтительнее — агрегаты, а не деталь.
В: Как синхронизировать фичи ML онлайн/офлайн?
О: Один feature-definition (код/SQL), единый store, check «training-serving skew», тесты совместимости и мониторинг drift.
Шпаргалка
- Data contract + semver + CCB — против «тихих» изменений.
- CDC с ключами и порядком; идемпотентные MERGE.
- DQ на каждом слое; quarantine для сомнительных данных.
- SCD2 там, где важна история.
- Семантический слой = единый словарь метрик.
- Lineage и наблюдаемость — обязательно.
- Экспорт API — асинхронный, безопасный, с лимитами.
- ML — контракт фич/лейблов и мониторинг drift.



