Портфолио-версия проекта cloud-technologies: событийно-ориентированный конвейер обработки данных, который получает события о заказах через Apache Kafka, обогащает их данными из Redis, сохраняет исходные события в STG, загружает данные в слой DDS по модели Data Vault в PostgreSQL и формирует агрегаты CDM по связям пользователей, продуктов и категорий.
┌──────────────────────┐
│ События о заказах │
│ Kafka │
└──────────┬───────────┘
│
▼
┌──────────────────────┐
│ STG-сервис │
│ Flask + APScheduler │
│ Kafka + Redis + PG │
└───────┬────────┬─────┘
│ │
исходные события │ обогащённые события
│ ▼
│ Kafka STG topic
▼ │
┌──────────────────────┐
│ PostgreSQL / STG │
└──────────────────────┘
│
▼
┌──────────────────────┐
│ DDS-сервис │
│ Data Vault │
│ Hubs / Links / Sats │
└──────────┬───────────┘
│
│ Kafka DDS topic
▼
┌──────────────────────┐
│ CDM-сервис │
│ счётчики / витрины │
│ + PG │
└──────────────────────┘
- событийно-ориентированный ETL/ELT с использованием Apache Kafka;
- разделение системы на микросервисы в соответствии со слоями хранилища данных;
- загрузку данных в STG и сохранение исходных событий;
- обогащение данных из Redis;
- моделирование слоя DDS по принципам Data Vault: hubs, links и satellites;
- использование детерминированных ключей UUID5 и безопасной обработки конфликтов;
- инкрементальное построение агрегатов CDM с помощью
ON CONFLICT ... DO UPDATE; - работу с PostgreSQL через
psycopg; - контейнеризацию с использованием Docker Compose;
- структуру проекта, ориентированную на развёртывание через Helm;
- TLS/SASL-аутентификацию при подключении к управляемым Kafka и Redis;
- health-check endpoints и фоновую обработку данных по расписанию.
| Область | Технологии |
|---|---|
| Язык | Python 3.10 |
| Потоковая обработка | Apache Kafka, confluent-kafka |
| Обработка | Flask, APScheduler |
| Хранилище | PostgreSQL |
| Кэш / обогащение | Redis |
| Контейнеризация | Docker, Docker Compose |
| Развёртывание | Helm / структура Kubernetes |
| Безопасность | SASL/SCRAM + SSL, CA-сертификат |
| Моделирование данных | STG, Data Vault-style DDS, CDM |
cloud-technologies/
├── README.md
├── docs/
│ └── architecture.md
├── solution/
│ ├── docker-compose.yaml
│ ├── service_stg/
│ │ ├── dockerfile
│ │ ├── requirements.txt
│ │ └── src/
│ ├── service_dds/
│ │ ├── dockerfile
│ │ ├── requirements.txt
│ │ └── src/
│ └── service_cdm/
│ ├── dockerfile
│ ├── requirements.txt
│ └── src/
└── deploy/
└── helm/
service_stg получает событие о заказе из Kafka, сохраняет исходный payload в stg.order_events, получает справочные данные о пользователе и ресторане из Redis, нормализует структуру заказа и публикует обогащённое событие в следующий Kafka-топик.
service_dds получает обогащённые события и преобразует их в модель данных в стиле Data Vault:
- Hubs: пользователь, ресторан, заказ, продукт, категория;
- Links: заказ–пользователь, заказ–продукт, продукт–категория, продукт–ресторан;
- Satellites: стоимость заказа, статус заказа, названия продуктов, названия ресторанов, имена пользователей.
Бизнес-ключи преобразуются в детерминированные идентификаторы UUID5. При вставке используется обработка конфликтов, что делает повторную обработку событий более безопасной.
service_cdm получает сообщения, сформированные на основе DDS, и поддерживает аналитические счётчики для связей пользователь–продукт и пользователь–категория.
Агрегированные значения обновляются атомарно с использованием PostgreSQL ON CONFLICT ... DO UPDATE.
Все параметры подключения передаются через переменные окружения. Учётные данные не должны храниться непосредственно в Git-репозитории.
Основные группы переменных:
KAFKA_*— брокер, параметры аутентификации, consumer group и топики;REDIS_*— параметры подключения к Redis;PG_WAREHOUSE_*— параметры подключения к PostgreSQL.
Для реального развёртывания секреты рекомендуется передавать через менеджер секретов или Kubernetes Secrets, а не хранить в открытом виде в переменных окружения или конфигурационных файлах.
- Создайте необходимые переменные окружения на основе параметров вашей облачной инфраструктуры.
- Поместите CA-сертификат Yandex Cloud в расположение, которое ожидает Docker-образ, либо адаптируйте Dockerfile под источник сертификата.
- Запустите сервисы:
docker compose -f solution/docker-compose.yaml up --build- Проверьте health-check endpoints:
curl http://localhost:5011/health
curl http://localhost:5012/health
curl http://localhost:5013/healthОжидаемый ответ:
healthy
Репозиторий представляет собой учебную/портфолио-реализацию потоковой платформы обработки данных.
Проект сохраняет исходную архитектуру и основной подход к обработке данных, при этом содержит более подробную документацию и структуру, ориентированную на демонстрацию инженерных навыков.
Проект не позиционируется как полностью готовое production-решение.
Перед использованием в production необходимо усилить следующие компоненты:
- автоматизированное тестирование;
- валидацию схем сообщений;
- обработку ошибочных сообщений и Dead Letter Queue;
- мониторинг и метрики;
- сквозную идемпотентность обработки событий;
- управление секретами;
- CI/CD;
- транзакционную обработку и оптимизацию batch-загрузок.
Python · SQL · PostgreSQL · Apache Kafka · Redis · Docker · Docker Compose · Helm · Kubernetes · ETL · Data Vault · STG/DDS/CDM · событийно-ориентированная архитектура · микросервисы · интеграция с облачными сервисами · моделирование данных
Исходный репозиторий: https://github.com/verydnobl337/cloud-technologies