Учебный проект по построению загрузочного контура аналитического хранилища данных.
Проект демонстрирует полный путь обработки данных: получение CSV-файла из S3, загрузку данных в staging-слой Vertica, формирование DWH-таблиц и выполнение аналитических SQL-запросов с использованием CTE.
- автоматизировать загрузку данных из S3 с помощью Apache Airflow;
- загрузить исходные данные в staging-слой Vertica;
- сформировать связи между пользователями и группами;
- сохранить историю событий пользователей в DWH;
- получить аналитические показатели по активности пользователей в группах;
- использовать SQL-конструкции
JOIN,CTE,GROUP BY, агрегатные функции и оконные/аналитические подходы там, где это необходимо.
┌───────────────┐
│ S3 │
│ group_log.csv│
└───────┬───────┘
│
▼
┌─────────────────────┐
│ Apache Airflow │
│ get_group_log_dag │
└──────────┬──────────┘
│
▼
/data/group_log.csv
│
▼
┌─────────────────────┐
│ Apache Airflow │
│ load_group_log_dag │
└──────────┬──────────┘
│ COPY
▼
┌─────────────────────┐
│ Vertica │
│ STAGING.group_log │
└──────────┬──────────┘
│
│ INSERT / JOIN
▼
┌─────────────────────┐
│ Vertica DWH │
│ h_users / h_groups │
│ l_user_group_activity│
│ s_auth_history │
└──────────┬──────────┘
│
▼
┌─────────────────────┐
│ Аналитические SQL │
│ CTE / JOIN / GROUP BY│
└─────────────────────┘
analytical-data-warehouse/
├── README.md
├── .gitignore
└── src/
├── dags/
│ ├── get_group_log_dag.py
│ └── load_group_log_in_stg_dag.py
└── sql/
├── ddl_group_log.sql
├── ddl_l_user_group_activity.sql
├── ddl_s_auth_history.sql
├── dml_l_user_group_activity.sql
├── dml_s_auth_history.sql
├── cte_user_group_log.sql
├── cte_user_group_message.sql
└── analytical_answer_with_cte.sql
- Python — реализация Airflow DAG;
- Apache Airflow — оркестрация загрузочного процесса;
- Amazon S3 — источник исходного CSV-файла;
- Vertica — аналитическая СУБД и DWH;
- SQL — DDL, DML и аналитические запросы;
- Git — версионирование проекта.
DAG get_group_log_dag.py подключается к S3 через S3Hook, получает файл group_log.csv из bucket sprint6 и сохраняет его локально по пути /data/group_log.csv.
DAG load_group_log_in_stg_dag.py использует VerticaHook и команду COPY для загрузки CSV-файла в таблицу VT26052617E774__STAGING.group_log.
SQL-скрипты создают и заполняют две основные сущности:
l_user_group_activity— связь пользователя с группой;s_auth_history— история событий пользователей.
Для формирования surrogate/hash key связи пользователя и группы используется HASH(user_id, group_id).
Запросы в src/sql/ рассчитывают количество пользователей, добавленных в группы, количество пользователей, создавших сообщения, и коэффициент конверсии:
conversion = users_with_messages / added_users
Для обработки отдельных этапов аналитики используются CTE, JOIN, COUNT(DISTINCT ...), COALESCE и NULLIF.
Проект рассчитан на учебное окружение с предоставленным контейнером Yandex Practicum.
Запуск контейнера:
docker run \
-d \
-p 3000:3000 \
-p 3002:3002 \
-p 15432:5432 \
--mount src=airflow_sp5,target=/opt/airflow \
--mount src=lesson_sp5,target=/lessons \
--mount src=db_sp5,target=/var/lib/postgresql/data \
--name=de-project-adb-server-local \
cr.yandex/crp1r8pht0n0gl25aug1/de-pg-cr-af:latestПосле запуска:
- Airflow:
http://localhost:3000/airflow - Vertica/PostgreSQL endpoint учебного окружения:
localhost:15432
Конкретные Airflow connections (
s3_connection,vertica_conn) и структура учебной базы должны быть настроены в окружении проекта.
ddl_group_log.sql— создание staging-таблицы.ddl_l_user_group_activity.sql— создание link-таблицы DWH.ddl_s_auth_history.sql— создание satellite/history-таблицы DWH.dml_l_user_group_activity.sql— загрузка связей пользователь–группа.dml_s_auth_history.sql— загрузка истории событий.cte_user_group_log.sqlиcte_user_group_message.sql— отдельные аналитические расчёты.analytical_answer_with_cte.sql— итоговый аналитический запрос.
Проект показывает практическое применение нескольких компонентов Data Engineering:
- ingestion из объектного хранилища;
- оркестрация через Airflow;
- staging-слой;
- DWH-моделирование с hub/link/satellite-подходом;
- DDL и DML для аналитического хранилища;
- построение аналитических запросов на SQL;
- расчёт бизнес-метрики на основе нескольких источников данных внутри DWH.
Проект выполнен в рамках учебного спринта «Аналитические базы данных» и оформлен как самостоятельный портфолио-проект с сохранением исходной логики решения.