Зачем

У банков и ритейлеров, с которыми я работала, почти всегда одна и та же боль: данные о клиенте размазаны по CRM, core‑banking и платёжным системам, и никто не может ответить на простой вопрос — а что мы вообще знаем об этом клиенте прямо сейчас. Чтобы попрактиковаться в решении этой задачи и показать подход публично, я собрала pet‑проект: end‑to‑end пайплайн, который строит единую витрину Customer 360 из синтетических данных.

Сразу оговорюсь: все данные генерируются локально скриптом и полностью синтетические. Кода или данных работодателя в проекте нет — это самостоятельная реализация архитектуры, которую я использую в повседневной работе DWH/data‑аналитика.

Архитектура

Пайплайн собран слоями, как я бы хотела видеть production DWH: raw → staging → intermediate → mart.

  • raw — сырые данные как есть у источника, без изменений и трансформаций, просто загруженные в PostgreSQL

  • staging (stg_customers, stg_accounts, stg_transactions) — типизация, нейминг, базовая валидация

  • intermediate (int_customer_identity) — детерминированное разрешение идентичности клиента

  • mart (customer_360) — одна строка на клиента с поведенческими KPI

Ключевое архитектурное решение — разрешение идентичности вынесено в отдельный слой, а не размазано по mart‑модели. Это значит, что логику мэтчинга клиентов можно менять и улучшать, не трогая бизнес‑метрики выше по потоку. На практике это то решение, которое экономит больше всего нервов при рефакторинге DWH.

Вот как выглядит слой identity resolution в dbt:

select
    customer_id,
    full_name,
    email,
    city,
    created_at,
    md5(lower(trim(email))) as customer_match_key
from {{ ref('stg_customers') }}

Здесь customer_match_key — детерминированный ключ на основе нормализованного email. В реальном enterprise‑контуре сюда обычно добавляется fuzzy‑мэтчинг (телефон, ФИО, устройство), но начинать стоит именно с детерминированного и объяснимого ключа — его легко дебажить, и на нём проще строить более сложную логику позже.

Финальная mart‑модель агрегирует метрики по продуктам и транзакциям и приклеивает их к identity‑слою:

with account_metrics as (
    select
        customer_id,
        count(*) as product_count,
        count(*) filter (where status = 'ACTIVE') as active_product_count,
        sum(balance) as total_balance
    from {{ ref('stg_accounts') }}
    group by customer_id
),
transaction_metrics as (
    select
        a.customer_id,
        count(*) as transaction_count,
        sum(t.amount) filter (where t.direction = 'CREDIT') as total_inflow,
        sum(t.amount) filter (where t.direction = 'DEBIT') as total_outflow,
        max(t.transaction_ts) as last_activity_at
    from {{ ref('stg_transactions') }} t
    join {{ ref('stg_accounts') }} a using (account_id)
    group by a.customer_id
)
select
    c.customer_id,
    c.customer_match_key,
    c.full_name,
    c.email,
    c.city,
    c.created_at,
    current_date - c.created_at as customer_tenure_days,
    coalesce(a.product_count, 0) as product_count,
    coalesce(a.active_product_count, 0) as active_product_count,
    coalesce(a.total_balance, 0) as total_balance,
    coalesce(t.transaction_count, 0) as transaction_count,
    coalesce(t.total_inflow, 0) as total_inflow,
    coalesce(t.total_outflow, 0) as total_outflow,
    t.last_activity_at,
    case
        when coalesce(t.transaction_count, 0) >= 15 then 'HIGH_ACTIVITY'
        when coalesce(t.transaction_count, 0) >= 5 then 'MEDIUM_ACTIVITY'
        else 'LOW_ACTIVITY'
    end as activity_segment
from {{ ref('int_customer_identity') }} c
left join account_metrics a using (customer_id)
left join transaction_metrics t using (customer_id)

Обратите внимание: агрегации по продуктам и транзакциям считаются в отдельных CTE и приклеиваются через left join — это позволяет добавлять новые источники метрик, не переписывая существующую логику.

Качество данных — не факультатив

В dbt‑проекте тесты лежат рядом с моделями, а не где‑то отдельно:

# models/marts/schema.yml
version: 2
models:
  - name: customer_360
    description: One trusted analytical row per retail-banking customer.
    columns:
      - name: customer_id
        tests: [unique, not_null]
      - name: customer_match_key
        tests: [unique, not_null]
      - name: activity_segment
        tests:
          - accepted_values:
              values: ['HIGH_ACTIVITY', 'MEDIUM_ACTIVITY', 'LOW_ACTIVITY']

Дополнительно пайплайн проверяет: валидность типов счетов и транзакций, корректность связи account‑to‑customer, отсутствие дублей клиента в финальной витрине и сверку количества клиентов между источником и mart‑слоем. Важный момент — dbt build запускает трансформации и тесты в одном шаге, поэтому сломанная модель физически не может тихо проехать как готовая: пайплайн просто упадёт на этом шаге.

Оркестрация

Airflow‑DAG простой и линейный — генерация синтетических источников, загрузка в raw, dbt build, финальный quality gate:

from datetime import datetime

from airflow import DAG
from airflow.operators.bash import BashOperator

with DAG(
    dag_id="customer_360_daily",
    start_date=datetime(2026, 1, 1),
    schedule="0 6 * * *",
    catchup=False,
    tags=["portfolio", "dwh", "customer-360"],
) as dag:
    generate_sources = BashOperator(
        task_id="generate_synthetic_sources",
        bash_command="cd /opt/project && python -m src.generate_data",
    )
    load_raw = BashOperator(
        task_id="load_raw_layer",
        bash_command="cd /opt/project && python -m src.load_raw",
    )
    dbt_build = BashOperator(
        task_id="build_and_test_models",
        bash_command="cd /opt/project/dbt/customer360 && dbt build --profiles-dir .",
    )
    quality_gate = BashOperator(
        task_id="customer_360_quality_gate",
        bash_command="cd /opt/project && python -m src.quality_gate",
    )

    generate_sources >> load_raw >> dbt_build >> quality_gate

DAG идемпотентен: сырые таблицы перезаписываются для конкретного синтетического снэпшота, а dbt‑модели пересобираются по объявленным зависимостям. Это значит, что перезапуск пайплайна на том же снэпшоте всегда даёт одинаковый результат — что сильно упрощает отладку и демонстрацию.

Вся инфраструктура поднимается через Docker Compose одной командой, юнит‑тесты на pytest можно гонять и без поднятия инфраструктуры, а CI на GitHub Actions прогоняет их на каждый пуш.

При чём тут моя основная работа

Эта архитектура — не абстрактное упражнение. Сейчас в БКС я участвую в миграции корпоративного DWH с MS SQL Server на Greenplum и строю похожую клиентоцентричную модель данных на слоях ODS/DDS/DMA. Реальная реализация там сейчас близка к полному запуску. Продовый код и данные я, разумеется, показать не могу — этот pet‑проект как раз и существует для того, чтобы отработать и продемонстрировать ту же архитектуру на нейтральных синтетических данных.

Что можно докрутить

  • инкрементальные dbt‑модели и snapshot‑история изменений;

  • fuzzy‑мэтчинг идентичности с confidence score вместо чистого детерминированного ключа;

  • anomaly detection для нетипичного поведения клиента;

  • публикация data lineage и метрик observability;

  • CI‑пайплайн с поднятием PostgreSQL как сервис‑контейнера прямо в тестах.

Буду рада фидбэку и вопросам в комментариях — особенно от тех, кто тоже строил Customer 360 или похожие клиентоцентричные витрины в проде.