Skip to content

Latest commit

 

History

8 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Cloud Technologies — потоковая платформа обработки данных

Портфолио-версия проекта 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/

Поток обработки данных

1. STG

service_stg получает событие о заказе из Kafka, сохраняет исходный payload в stg.order_events, получает справочные данные о пользователе и ресторане из Redis, нормализует структуру заказа и публикует обогащённое событие в следующий Kafka-топик.

2. DDS

service_dds получает обогащённые события и преобразует их в модель данных в стиле Data Vault:

  • Hubs: пользователь, ресторан, заказ, продукт, категория;
  • Links: заказ–пользователь, заказ–продукт, продукт–категория, продукт–ресторан;
  • Satellites: стоимость заказа, статус заказа, названия продуктов, названия ресторанов, имена пользователей.

Бизнес-ключи преобразуются в детерминированные идентификаторы UUID5. При вставке используется обработка конфликтов, что делает повторную обработку событий более безопасной.

3. CDM

service_cdm получает сообщения, сформированные на основе DDS, и поддерживает аналитические счётчики для связей пользователь–продукт и пользователь–категория.

Агрегированные значения обновляются атомарно с использованием PostgreSQL ON CONFLICT ... DO UPDATE.

Конфигурация

Все параметры подключения передаются через переменные окружения. Учётные данные не должны храниться непосредственно в Git-репозитории.

Основные группы переменных:

  • KAFKA_* — брокер, параметры аутентификации, consumer group и топики;
  • REDIS_* — параметры подключения к Redis;
  • PG_WAREHOUSE_* — параметры подключения к PostgreSQL.

Для реального развёртывания секреты рекомендуется передавать через менеджер секретов или Kubernetes Secrets, а не хранить в открытом виде в переменных окружения или конфигурационных файлах.

Локальный запуск

  1. Создайте необходимые переменные окружения на основе параметров вашей облачной инфраструктуры.
  2. Поместите CA-сертификат Yandex Cloud в расположение, которое ожидает Docker-образ, либо адаптируйте Dockerfile под источник сертификата.
  3. Запустите сервисы:
docker compose -f solution/docker-compose.yaml up --build
  1. Проверьте 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

About

Event-driven streaming data platform built with Apache Kafka, PostgreSQL, Redis, and Docker. Implements STG, Data Vault-style DDS, and CDM layers with microservices-based data processing and incremental aggregation.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages