Платформа мониторинга окружающей среды и курсов валют с обнаружением аномалий, построенная на современном стеке data engineering. Является моим учебным полигоном для Data Enginerring
Полностью контейнеризованный пайплайн данных, который непрерывно отслеживает погодные условия и качество воздуха в заданных пользователем локациях, автоматически выявляет аномалии путём сравнения измерений с соседними станциями и предоставляет аналитические возможности через выделенный OLAP-движок.
Основной рабочий процесс: Локации задаются через REST API → сбор данных по расписанию → структурированное хранение → обнаружение аномалий → синхронизация в аналитическую БД → дашборды
Источники данных:
- Open-Meteo API — погодные данные (температура, давление, влажность, ветер)
- OpenAQ API — данные о качестве воздуха (PM2.5, PM10, NO₂, O₃, SO₂, CO)
- [cbr.xml API] - данные о курсах валют к рублю
| Слой | Технология | Назначение |
|---|---|---|
| Оркестрация | Apache Airflow 2.9 | Планирование пайплайнов на основе DAG, повторные попытки, мониторинг |
| API | FastAPI | Управление локациями, доступ к данным, Swagger UI |
| OLTP БД | PostgreSQL 13 | Оперативное хранение — сырые и структурированные данные |
| OLAP БД | ClickHouse 24.12 | Колоночная аналитика — агрегации на больших объёмах |
| DataLake | MinIO собирает все сырые данные | |
| Визуализация | Grafana | Дашборды в реальном времени для погоды и качества воздуха |
| Контейнеризация | Docker Compose | Оркестрация всего стека, воспроизводимое окружение |
| Языки | Python 3.14, SQL | Логика ETL, модели данных, аналитические запросы |
| Kafka(Redpanda) | Полная связь всех сервисов |
- Погода из Open-Meteo: температура, давление, влажность, скорость ветра
- Качество воздуха из OpenAQ: PM2.5, PM10, NO₂, O₃, SO₂, CO
- Сырой JSON сохраняется для аудита, структурированные таблицы — для запросов
- Сравнение каждой локации с соседними станциями (пространственный радиус)
- Статистическое обнаружение на основе пороговых значений (отклонение от среднего по соседям)
- Защита от дубликатов: повторяющиеся аномалии не создаются
- Разные типы аномалий: скачки температуры, падение влажности, пики загрязнений
- Уровень качества воздуха вычисляется на уровне базы данных
- Шесть уровней: good → moderate → unhealthy_sensitive → unhealthy → very_unhealthy → hazardous
- Используются пороговые значения ВОЗ для всех 6 критериальных загрязнителей
- Доступен в SQL-запросах, отображается в дашбордах Grafana
- Синхронизация PostgreSQL → ClickHouse каждые 5 минут
- Идемпотентность: используется метка max(id), дубли исключены
- Быстрые агрегации на растущих объёмах данных
- Swagger UI по адресу
/docs— интерактивная документация API - Эндпоинты: добавление/удаление локаций, просмотр сырых и структурированных данных, запрос аномалий, аналитика
- Pydantic-схемы для валидации запросов
- Airflow UI: история запусков DAG, логи задач, отслеживание повторных попыток
- Дашборды Grafana: графики температуры, карты качества воздуха, частота аномалий
- Все компоненты с проверками работоспособности (healthcheck) в Docker Compose
- Docker и Docker Compose v2
- Python 3.12 (для локальной разработки)
- Ключ API OpenAQ (зарегистрироваться здесь)
- Клонируйте репозиторий
git clone https://github.com/cheburek 4535/DataEngineerPlatform.git
cd DataEngineerPlatform- Задайте переменные окружения
# Создайте файл .env
echo -e "AIRFLOW_UID=$(id -u)\nAIRFLOW_GID=0" > .env
# Добавьте ваш ключ OpenAQ API в services/extract/extract_air_quality.py
# OPENAQ_API_KEY = "ваш-ключ-сюда"- Инициализируйте Airflow
docker compose up airflow-init- Запустите платформу
docker compose up -d-
Доступ к сервисам | Сервис | URL | Учётные данные | |---|---|---| | Airflow | http://localhost:8080 | airflow / airflow | | FastAPI Swagger | http://localhost:8000/docs | — | | Grafana | http://localhost:3000 | admin / admin | | ClickHouse HTTP | http://localhost:8123 | default / ch_password | | PostgreSQL | localhost:5433 | weather_user / weather_pass |
-
Добавьте первую локацию
curl -X POST http://localhost:8000/locations/add \
-H "Content-Type: application/json" \
-d '{"lat": 55.75, "lon": 37.62, "check_interval": 900}'- Включите DAG в Airflow
- Откройте Airflow UI → включите
weather_pipeline,air_quality_pipeline,check_locations,sync_to_clickhouse - Запустите вручную или дождитесь срабатывания по расписанию
- Откройте Airflow UI → включите
| DAG ID | Расписание | Описание |
|---|---|---|
check_locations |
Каждые 6 часов | Обходит все отслеживаемые локации, запускает полный пайплайн |
sync_to_clickhouse |
Каждый час | Инкрементально синхронизирует погоду, курсы валют и качество воздуха в ClickHouse |
check_currencies |
Каждый час | Проверяет курсы всех валют |
check_air_quality |
Каждые 3 часа | Обходит все отслеживаемые локации, запускает полный пайплайн aq |
Для каждой локации система:
- Получает текущие данные о погоде/качестве воздуха
- Находит все станции в настраиваемом радиусе
- Вычисляет средние значения по соседям для каждого параметра
- Отмечает измерения, превышающие порог × среднее отклонение
- Сохраняет записи об аномалиях с измеренными и ожидаемыми значениями
Это позволяет отлавливать локальные аномалии: резкое падение температуры в одной точке при стабильности у соседей.
Разделение ответственности:
- PostgreSQL обрабатывает операционные нагрузки (вставки, точечные запросы, CRUD)
- ClickHouse обрабатывает аналитические нагрузки (агрегации, временные ряды, тренды)
- Airflow оркестрирует потоки данных, но не содержит бизнес-логики
Почему колоночное хранение (ClickHouse):
- Данные о погоде быстро накапливаются: 5000 локаций × 2 проверки/день × 365 дней = 365000 строк/год
- ClickHouse сжимает колоночные данные в 5–10 раз, запросы читают только нужные колонки
- Агрегации (AVG температуры по месяцам) векторизованы, а не построчны
Почему Airflow, а не cron:
- DAG явно определяют зависимости (извлечь → преобразовать → обнаружить)
- Встроенные повторные попытки с экспоненциальной задержкой
- История выполнения и логи из коробки
- Возможность заполнения исторических данных (backfill)
- Основы Data Engineering: проектирование ETL/ELT-пайплайнов, OLTP против OLAP, паттерны инкрементальной синхронизации
- Оркестрация: проектирование DAG, зависимости задач, планирование, стратегии повторных попыток, мониторинг
- Интеграция API: работа с REST API, ограничение частоты запросов, обработка ошибок, контракты данных
- Моделирование данных: сырые и структурированные таблицы, JSONB для полуструктурированных данных, гибридные свойства SQLAlchemy
- Контейнеризация: многосервисный Docker Compose, кастомные Dockerfile, управление томами, healthcheck-проверки
- Аналитическая инженерия: колоночные базы данных, партиционирование, материализованные представления, проектирование дашбордов Grafana
- Экосистема Python: SQLAlchemy ORM, FastAPI, Pydantic, clickhouse-connect, логирование
- Интеграция Golang: Отдельный сервис на Go связывается с проектом через gRPC, для ускорения обработки больших данных
- dbt — трансформации данных как код с тестированием и документацией
- Great Expectations — наборы проверок качества данных
- CI/CD — GitHub Actions для тестирования и деплоя DAG
- Kubernetes — миграция с Docker Compose на K8s для масштабируемости
- Оповещения — уведомления в Telegram/Slack о критических аномалиях
Создано как самообразовательный проект для входа в data engineering.
<img width="1574" height="728" alt="{E4354A01-207D-4A35-B64C-E545308A2DDF}" src="https://github.com/user-attachments/assets/7af57633-c861-4365-ae8b-16a591e10f69" />
<img width="1587" height="392" alt="{7E04E419-B6BC-48B7-91F8-9B374763A10D}" src="https://github.com/user-attachments/assets/4ca64a55-9213-47d2-b732-d707c4c08321" />
<img width="1573" height="804" alt="{62B85150-53D6-424D-9E34-53C0F7CD19BB}" src="https://github.com/user-attachments/assets/3dd4bcf4-0ea5-4090-b879-e60ddfe6f852" />