End-to-end ML-пайплайн для предсказания оттока клиентов финтех-компании.
Данные обрабатываются через PySpark, хранятся в ClickHouse, оркестрация через Apache Airflow, визуализация в Grafana. С помощью FastAPI Реализована возможность получения предсказаний моделей по запросу.
- Задача
- Датасет
- Технологии
- Структура DAG
- Описание тасков
- Структура проекта
- Запуск
- Результаты моделей
- Использование результатов работы моделей
- Grafana
Бинарная классификация: предсказать, уйдёт клиент (churn) или останется.
На основе демографических данных, транзакционной активности и поведенческих метрик строятся признаки, обучаются три модели, результаты сравниваются и визуализируются.
COFINFAD (Colombian Fintech Financial Analytics Dataset) — синтетические данные финтех-компании из Колумбии.
| Таблица | Записей | Описание |
|---|---|---|
customer_data.csv |
48 723 клиента, 54 колонки | Демография, продукты, поведение, NPS, целевая переменная churn_probability |
transactions_data.csv |
~3.1 млн транзакций | customer_id, дата, сумма, тип (deposit, payment, transfer, withdrawal) |
| Технология | Роль в проекте |
|---|---|
| PySpark | ETL: чтение CSV, агрегация транзакций, feature engineering, обучение моделей (ML) |
| ClickHouse | Хранение сырых данных, фичей, предсказаний и метрик моделей |
| Apache Airflow | Оркестрация: DAG с 7 тасками, еженедельный запуск |
| Grafana | Дашборды: метрики моделей, распределение предсказаний, аналитика клиентов |
| FastAPI | Эндпоинт для получения предсказаний |
| Docker | Инфраструктура: все сервисы в docker-compose (Airflow, ClickHouse, Grafana) |
| Python | pandas, clickhouse-connect, Jupyter Notebooks |
ingest_raw_data (Task 1)
│
▼
build_features_spark (Task 2)
│
▼
quality_check_and_split (Task 3)
┌────┼────┐
▼ ▼ ▼
train_ train_ train_ (Task 4a / 4b / 4c)
logreg rf gbt
└────┼────┘
▼
collect_predictions (Task 5)
Task 1: Ingest Raw Data
Читает CSV-файлы через PySpark, эмулирует еженедельный батч (случайная выборка ~15% клиентов с детерминированным seed из даты), записывает сырые данные в ClickHouse (raw_customers, raw_transactions).
Task 2: Build Features
Читает сырые данные из ClickHouse, через PySpark считает транзакционные агрегаты (count, sum, avg, stddev по клиенту, pivot по типам транзакций), кодирует категориальные признаки (StringIndexer), бинаризует target (churn_probability → churn_label), записывает feature-таблицу в ClickHouse (features).
Task 3: Quality Check & Split
Проверяет и чистит данные: удаление дубликатов по customer_id, заполнение NULL (медиана для числовых, мода для категориальных), клипирование выбросов (IQR). Делает train/test split (80/20), сохраняет в parquet. Метаинформация о разбиении — в split_info.json.
Task 4a/4b/4c: Train Models (параллельно)
Три параллельных таска, каждый обучает свою модель через PySpark ML:
| Модель | Гиперпараметры |
|---|---|
| Logistic Regression | maxIter=100, regParam=0.01 |
| Random Forest | numTrees=5, maxDepth=10 |
| Gradient Boosted Trees | maxIter=10, maxDepth=5 |
Каждый таск: обучение на train → оценка на test (AUC-ROC, AUC-PR, F1, Accuracy, Precision, Recall) → скоринг всех клиентов → сохранение предсказаний в parquet и метрик в JSON.
Task 5: Collect Predictions
Собирает предсказания трёх моделей из parquet, объединяет по customer_id в одну таблицу, записывает в ClickHouse (predictions). Собирает метрики всех моделей → таблица model_metrics. Данные готовы для Grafana.
churn_predictor/
├── dags/
│ └── churn_prediction_dag.py # DAG для Airflow
├── scripts/
│ ├── task1_ingest.py # Загрузка сырых данных
│ ├── task2_features.py # Feature engineering
│ ├── task3_quality_split.py # Очистка + split
│ ├── task4a_train_logreg.py # Logistic Regression
│ ├── task4b_train_rf.py # Random Forest
│ ├── task4c_train_boosting.py # Gradient Boosting
│ └── task5_collect_predictions.py # Сбор результатов → ClickHouse
├── notebooks/ # Jupyter-ноутбуки (разработка)
├── config/ # Конфигурация Airflow
├── logs/ # Логи Airflow
├── docker-compose.yml # Airflow + ClickHouse + Grafana
├── Dockerfile # Образ Airflow с PySpark и Java
├── default-user.xml # Конфигурация доступа ClickHouse
└── .env # Переменные окружения
Данные и модели (вне репозитория)
churn_prediction_project/
├── data/
│ ├── customer_data.csv
│ ├── transactions_data.csv
│ ├── splits/
│ │ ├── train.parquet
│ │ ├── test.parquet
│ │ └── split_info.json
│ └── predictions/
│ ├── logreg_predictions.parquet
│ ├── random_forest_predictions.parquet
│ └── boosting_predictions.parquet
├── models/
│ ├── logreg/metrics.json
│ ├── random_forest/metrics.json
│ └── boosting/metrics.json
└── churn_predictor/ # ← репозиторий
1. Клонировать репозиторий
git clone https://github.com/jack1591/churn_predictor.git
cd churn_predictor2. Подготовить данные
Скачать COFINFAD dataset и положить CSV-файлы в ../data/.
3. Запустить инфраструктуру
docker-compose build
docker-compose up -d4. Открыть сервисы
| Сервис | URL | Логин |
|---|---|---|
| Airflow | localhost:8082 | airflow / airflow |
| Grafana | localhost:3000 | admin / admin |
| ClickHouse | localhost:8123 | — |
5. Запустить DAG
В Airflow включить DAG churn_prediction_pipeline и запустить вручную или дождаться расписания.
| Модель | AUC-ROC | AUC-PR | Accuracy | F1 | Precision | Recall |
|---|---|---|---|---|---|---|
| Logistic Regression | 0.6607 | 0.7649 | 0.6364 | 0.6424 | 0.6909 | 0.6364 |
| Random Forest | 0.7500 | 0.8516 | 0.6364 | 0.6424 | 0.6909 | 0.6364 |
| Gradient Boosting | 0.7500 | 0.8788 | 0.7273 | 0.7319 | 0.7485 | 0.7273 |
Для быстрого доступа к предсказаниям моделей используется FastAPI. Один эндпоинт позволяет получить результаты всех трёх моделей по конкретному клиенту.
Эндпоинт:
| Метод | Путь | Описание |
|---|---|---|
| GET | /predict/{customer_id} |
Возвращает предсказания Logistic Regression, Random Forest и Gradient Boosted Trees для клиента с указанным customer_id. |
Пример запроса:
http://localhost:8000/predict/16783564
Пример ответа:
{"customer_id":16783564,"logreg_pred":1,"forest_pred":1,"boosting_pred":1}
Дашборд включает панели:
- Сравнение метрик моделей
- Распределение возраста и собственности по доходам
- Сравнение профилей лояльного клиента и клиента, который ушел
Datasource: ClickHouse → clickhouse:9000, user default, без пароля.


