docs(docs): актуализировано ТЗ и потоки ingest
- Зачем: - убрать рассинхрон между кратким ТЗ, архитектурой и планом генератора - Что: - сокращен docs/DE-task.md до формата краткого ТЗ проекта - обновлены docs/ARCHITECTURE.md и README.md: bootstrap через kafka_load и steady-stream через generator-service - обновлен plans/generator_demo_stream_plan.md: режим steady-stream и тик-публикация - Проверка: - просмотрен git diff по измененным файлам - в коммит включены только мои документационные изменения
This commit is contained in:
@@ -9,8 +9,10 @@
|
|||||||
витринами.
|
витринами.
|
||||||
|
|
||||||
Поток данных коротко:
|
Поток данных коротко:
|
||||||
`data/*.jsonl → Airflow (kafka_load) → Kafka → ClickHouse (слой STG) → Airflow
|
- **bootstrap**: `data/*.jsonl → Airflow (kafka_load) → Kafka → ClickHouse (слой STG) →
|
||||||
(etl_pipeline: STG → ODS → DDS → DM) → Superset`.
|
Airflow (etl_pipeline: STG → ODS → DDS → DM) → Superset`.
|
||||||
|
- **steady-stream**: `generator-service → Kafka → ClickHouse (STG) → Airflow (etl_pipeline)
|
||||||
|
→ Superset`.
|
||||||
|
|
||||||
## Куда дальше
|
## Куда дальше
|
||||||
|
|
||||||
|
|||||||
+41
-9
@@ -27,10 +27,14 @@
|
|||||||
flowchart LR
|
flowchart LR
|
||||||
subgraph AF["Airflow"]
|
subgraph AF["Airflow"]
|
||||||
DAG1["ddl_init"]
|
DAG1["ddl_init"]
|
||||||
DAG2["kafka_load"]
|
DAG2["kafka_load (bootstrap)"]
|
||||||
DAG3["etl_pipeline"]
|
DAG3["etl_pipeline"]
|
||||||
end
|
end
|
||||||
|
|
||||||
|
subgraph GEN["Generator"]
|
||||||
|
G["generator-service (steady-stream)"]
|
||||||
|
end
|
||||||
|
|
||||||
subgraph Kafka["Kafka"]
|
subgraph Kafka["Kafka"]
|
||||||
K[4 топика]
|
K[4 топика]
|
||||||
end
|
end
|
||||||
@@ -53,7 +57,9 @@ flowchart LR
|
|||||||
V[витрины VIEW]
|
V[витрины VIEW]
|
||||||
end
|
end
|
||||||
|
|
||||||
DAG2 -->|JSONL| K -->|MV| S
|
DAG2 -->|bootstrap JSONL| K
|
||||||
|
G -->|stream events| K
|
||||||
|
K -->|MV| S
|
||||||
S -->|batch| O
|
S -->|batch| O
|
||||||
O -->|argMax + JOIN| D1 & D2
|
O -->|argMax + JOIN| D1 & D2
|
||||||
O -.->|ошибки| OE
|
O -.->|ошибки| OE
|
||||||
@@ -63,16 +69,25 @@ flowchart LR
|
|||||||
DAG3 -.->|batch| ODS & DDS
|
DAG3 -.->|batch| ODS & DDS
|
||||||
```
|
```
|
||||||
|
|
||||||
|
В учебном стенде предусмотрены два пути ingest:
|
||||||
|
|
||||||
|
- `bootstrap`: DAG `kafka_load` для разового/контрольного прогона из `data/*.jsonl`;
|
||||||
|
- `steady-stream`: автономный генератор, публикующий события в Kafka непрерывно.
|
||||||
|
|
||||||
### Слои и их назначение
|
### Слои и их назначение
|
||||||
|
|
||||||
```mermaid
|
```mermaid
|
||||||
flowchart LR
|
flowchart LR
|
||||||
subgraph AF["Airflow"]
|
subgraph AF["Airflow"]
|
||||||
DAG1["ddl_init"]
|
DAG1["ddl_init"]
|
||||||
DAG2["kafka_load"]
|
DAG2["kafka_load (bootstrap)"]
|
||||||
DAG3["etl_pipeline"]
|
DAG3["etl_pipeline"]
|
||||||
end
|
end
|
||||||
|
|
||||||
|
subgraph GEN["Generator"]
|
||||||
|
G["generator-service"]
|
||||||
|
end
|
||||||
|
|
||||||
subgraph L1["STG"]
|
subgraph L1["STG"]
|
||||||
KAFKA["Kafka Engine"]
|
KAFKA["Kafka Engine"]
|
||||||
STG_T["*_raw таблицы"]
|
STG_T["*_raw таблицы"]
|
||||||
@@ -91,7 +106,9 @@ flowchart LR
|
|||||||
DM_T["VIEW"]
|
DM_T["VIEW"]
|
||||||
end
|
end
|
||||||
|
|
||||||
DAG2 -->|JSONL| KAFKA -->|MV| STG_T
|
DAG2 -->|bootstrap JSONL| KAFKA
|
||||||
|
G -->|steady-stream| KAFKA
|
||||||
|
KAFKA -->|MV| STG_T
|
||||||
STG_T -->|batch| ODS_T
|
STG_T -->|batch| ODS_T
|
||||||
ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T
|
ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T
|
||||||
ODS_T -.->|ошибки| DQ
|
ODS_T -.->|ошибки| DQ
|
||||||
@@ -347,6 +364,7 @@ sequenceDiagram
|
|||||||
participant User as Пользователь
|
participant User as Пользователь
|
||||||
participant Compose as Docker Compose
|
participant Compose as Docker Compose
|
||||||
participant Airflow as Airflow
|
participant Airflow as Airflow
|
||||||
|
participant Gen as generator-service
|
||||||
participant K as Kafka
|
participant K as Kafka
|
||||||
participant CH as ClickHouse
|
participant CH as ClickHouse
|
||||||
participant STG as stg.*_raw
|
participant STG as stg.*_raw
|
||||||
@@ -358,6 +376,7 @@ sequenceDiagram
|
|||||||
Compose->>K: docker compose up -d kafka
|
Compose->>K: docker compose up -d kafka
|
||||||
Compose->>CH: docker compose up -d clickhouse
|
Compose->>CH: docker compose up -d clickhouse
|
||||||
Compose->>Airflow: docker compose up -d airflow-*
|
Compose->>Airflow: docker compose up -d airflow-*
|
||||||
|
Compose->>Gen: docker compose up -d generator
|
||||||
Compose-->>User: ✅ Инфраструктура готова
|
Compose-->>User: ✅ Инфраструктура готова
|
||||||
|
|
||||||
User->>Airflow: Trigger ddl_init
|
User->>Airflow: Trigger ddl_init
|
||||||
@@ -368,14 +387,22 @@ sequenceDiagram
|
|||||||
Airflow->>CH: sql/ddl/dm/40_dm.sql
|
Airflow->>CH: sql/ddl/dm/40_dm.sql
|
||||||
CH-->>User: ✅ Структура БД создана
|
CH-->>User: ✅ Структура БД создана
|
||||||
|
|
||||||
|
alt Bootstrap режим
|
||||||
User->>Airflow: Trigger kafka_load
|
User->>Airflow: Trigger kafka_load
|
||||||
Airflow->>K: precheck + prepare_topics
|
Airflow->>K: precheck + prepare_topics
|
||||||
loop 4 файла
|
loop 4 файла
|
||||||
Airflow->>K: KafkaProducer.send(topic, json_line)
|
Airflow->>K: KafkaProducer.send(topic, json_line)
|
||||||
end
|
end
|
||||||
|
K-->>User: ✅ Данные в Kafka
|
||||||
|
else Streaming режим
|
||||||
|
loop каждые 1-10 секунд
|
||||||
|
Gen->>K: send N_t (Poisson) в 4 топика
|
||||||
|
end
|
||||||
|
K-->>User: ✅ Непрерывный поток в Kafka
|
||||||
|
end
|
||||||
|
|
||||||
K->>CH: Потребление сообщений
|
K->>CH: Потребление сообщений
|
||||||
CH->>STG: INSERT через MV
|
CH->>STG: INSERT через MV
|
||||||
K-->>User: ✅ Данные в Kafka
|
|
||||||
CH-->>User: ✅ Данные в STG
|
CH-->>User: ✅ Данные в STG
|
||||||
|
|
||||||
User->>Airflow: Trigger etl_pipeline
|
User->>Airflow: Trigger etl_pipeline
|
||||||
@@ -607,20 +634,25 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
|
|||||||
|
|
||||||
### Airflow-оркестрация
|
### Airflow-оркестрация
|
||||||
|
|
||||||
Инфраструктура Airflow развёрнута и готова к использованию:
|
Инфраструктура Airflow развёрнута и отвечает за DDL/ETL.
|
||||||
|
Генератор работает отдельно и не управляется через Airflow DAG-и.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
# airflow/dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap)
|
# airflow/dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap)
|
||||||
# airflow/dags/kafka_load_dag.py — загрузка JSONL в Kafka (через kafka-python)
|
# airflow/dags/kafka_load_dag.py — bootstrap-загрузка JSONL в Kafka (через kafka-python)
|
||||||
# airflow/dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
|
# airflow/dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
|
||||||
|
|
||||||
# Учебный формат:
|
# Учебный формат:
|
||||||
# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
|
# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
|
||||||
# - SQL-файлы вызываются по фиксированным путям;
|
# - SQL-файлы вызываются по фиксированным путям;
|
||||||
# - загрузка данных в Kafka выполняется через DAG `kafka_load`.
|
# - ingest может идти двумя путями:
|
||||||
|
# 1) bootstrap через DAG `kafka_load`;
|
||||||
|
# 2) непрерывный поток через автономный `generator-service`.
|
||||||
#
|
#
|
||||||
# Основной demo-сценарий:
|
# Базовый demo-сценарий:
|
||||||
# ddl_init -> kafka_load -> etl_pipeline
|
# ddl_init -> kafka_load -> etl_pipeline
|
||||||
|
# Расширенный учебный сценарий:
|
||||||
|
# generator-service (continuous) + периодический etl_pipeline
|
||||||
```
|
```
|
||||||
|
|
||||||
**DAG `kafka_load`**:
|
**DAG `kafka_load`**:
|
||||||
|
|||||||
+30
-42
@@ -1,55 +1,43 @@
|
|||||||
**Дашборд для e-commerce кликстрима**
|
# DE-task: краткое ТЗ проекта
|
||||||
|
|
||||||
Представьте, что вы работаете в e-commerce проекте. Бизнес-юнит в процессе обсуждения, закупать ли вам кликстрим логов посещения некоторого портала другого магазина. Для первичного анализа вам предоставили данные, чтобы оценить применимость и выводы, которые можно сделать из них.
|
Дата актуализации: 14 февраля 2026.
|
||||||
|
|
||||||
Вам нужно подготовить данные для первичного анализа и собрать несколько графиков в дашборд. Финальный результат вашей работы — развернутая инфраструктура и дашборд с графиками на которых "бизнес" может сделать выводы о данных. Вам дали полный карт-бланш на тему того, как эти данные можно обработать.
|
## Цель
|
||||||
|
|
||||||
## Вам нужно:
|
Собрать учебный DE-стенд, который показывает полный путь данных:
|
||||||
|
`Kafka -> ClickHouse (STG/ODS/DDS/DM) -> BI`, с устойчивой обработкой "грязных" данных и понятной наблюдаемостью.
|
||||||
|
|
||||||
* спроектировать хранилище в /PostgreSQL/Clickhouse/ со слоями и объектами, которые вы посчитаете нужными
|
## Что нужно сделать
|
||||||
* написать все нужные расчёты и алгоритмы для создания витрин
|
|
||||||
* реализовать расчёт с помощью какого-то регулярного процесса
|
|
||||||
* сделать на основе этих данных дашборд который можно будет посмотреть в браузере
|
|
||||||
|
|
||||||
### Какие инструменты выбрать
|
- Развернуть стек в `docker compose`: Kafka, ClickHouse, Airflow, Superset, Prometheus, Grafana.
|
||||||
|
- Организовать ingest в Kafka в двух режимах:
|
||||||
|
- `bootstrap`: загрузка JSONL через Airflow DAG `kafka_load`;
|
||||||
|
- `steady-stream`: автономный генератор, работающий независимо от потребителей.
|
||||||
|
- Построить слои хранилища `STG -> ODS -> DDS -> DM` в ClickHouse.
|
||||||
|
- Настроить регулярные трансформации в Airflow (`ddl_init`, `etl_pipeline`).
|
||||||
|
- Подготовить витрины и дашборд для первичного анализа.
|
||||||
|
|
||||||
Пара вариантов:
|
## Ключевые требования
|
||||||
|
|
||||||
* **Kafka + Clickhouse + Airflow + Superset/Power BI**
|
- "Грязные" записи не должны валить пайплайн: ошибки фиксируются в ODS/DQ-слое.
|
||||||
* Postgresql + dagster + Grafana
|
- Генератор работает отдельно от потребителей и Airflow-триггеров.
|
||||||
* dlt + duckdb + Airflow + [Panel](https://panel.holoviz.org/)
|
- Решение должно быть воспроизводимым и пригодным для учебной демонстрации.
|
||||||
* Clickhouse + Airflow + Redis + [Streamlit](https://streamlit.io/)
|
|
||||||
* sqlitedb + скрипты питона + Flask + d3.js
|
|
||||||
|
|
||||||
Начните с маленького среза данных, чтобы ознакомиться, что за данные вам предоставляют:
|
## Границы MVP
|
||||||
|
|
||||||
* browser_events.jsonl — данные с информацией о просмотрах браузера
|
- Один простой режим генератора (`steady-stream`) без сложных сценариев.
|
||||||
* device_events.jsonl — данные об устройствах, с которых пользователи пользовались сайтом
|
- Без production-гарантий уровня exactly-once.
|
||||||
* geo_events.jsonl — данные о локации пользователя
|
- Фокус на рабочем end-to-end контуре, а не на полной симуляции реального продакшена.
|
||||||
* location_events.jsonl — вопреки названию, это данные не о локации пользователя, а непосредственно о положении пользователя на сайте и информации о том, откуда пользователь попал на страницу
|
|
||||||
|
|
||||||
Смысл задачи на **своей ВМ** развернуть инструменты и загрузить данные в Kafka и из нее прогнать их в Clickhouse /другая DB по слоям.
|
## Критерии готовности
|
||||||
|
|
||||||
В Airflow взять данные из таблиц и преобразовать, потом выгрузить в Clickhouse / другая DB в агрегированном формате под дашборд.
|
1. Стенд стабильно запускается и работает несколько часов.
|
||||||
|
2. Данные корректно проходят путь Kafka -> STG -> ODS -> DDS -> DM.
|
||||||
|
3. Поток событий в Kafka идёт постепенно малыми порциями, без искусственного минутного burst.
|
||||||
|
4. Есть базовый дашборд и метрики для разбора работы стенда.
|
||||||
|
|
||||||
Дашборд не оценивается, но будет плюсом, если получиться собрать удобный визуал.
|
## Где детали реализации
|
||||||
|
|
||||||
Записать можно на видео и отправить контактному лицу или продемонстрировать на собеседовании
|
- Техническая архитектура: `docs/ARCHITECTURE.md`
|
||||||
|
- План по генератору: `plans/generator_demo_stream_plan.md`
|
||||||
Можно использовать любые инструменты.
|
- Эксплуатация и запуск: `docs/OPERATIONS.md`
|
||||||
|
|
||||||
Данные могут быть(скорее всего) с ошибками и компания которая их предоставляем заранее об этом сообщила т.к. они не хотят делиться 100% качественными данными полноценно с продакшена.
|
|
||||||
|
|
||||||
### ИНФРАСТРУКТУРА
|
|
||||||
|
|
||||||
* инфраструктура клевая и работает = 18 баллов
|
|
||||||
* инфраструктура работает, но есть баги = 12 баллов
|
|
||||||
* инфраструктура почти есть, но не совсем = 6 баллов
|
|
||||||
* инфраструктуры нет = 0 баллов
|
|
||||||
|
|
||||||
### ДАШБОРДЫ
|
|
||||||
|
|
||||||
* дашборды есть и удобные и работают = 12 баллов
|
|
||||||
* дашборды работают, но неудобные = 8 баллов
|
|
||||||
* дашборды почти работают = 4 балла
|
|
||||||
* дашбордов нет = 0 баллов
|
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# План реализации: простой автономный генератор (rev4)
|
# План реализации: простой автономный генератор (rev5)
|
||||||
|
|
||||||
Дата ревизии: 14 февраля 2026.
|
Дата ревизии: 14 февраля 2026.
|
||||||
|
|
||||||
@@ -44,11 +44,13 @@
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## 4) Единственный режим генерации: `steady`
|
## 4) Единственный режим генерации: `steady-stream`
|
||||||
|
|
||||||
Логика режима:
|
Логика режима:
|
||||||
|
|
||||||
- каждую минуту публикуем фиксированный объём событий;
|
- публикуем в Kafka **постепенно**, короткими тиками (например, каждые 1-10 секунд);
|
||||||
|
- на каждом тике отправляем небольшую порцию сообщений;
|
||||||
|
- держим целевую интенсивность в `events/min`, а не крупный минутный batch;
|
||||||
- распределяем события по 4 топикам:
|
- распределяем события по 4 топикам:
|
||||||
- `browser_events`
|
- `browser_events`
|
||||||
- `location_events`
|
- `location_events`
|
||||||
@@ -59,15 +61,16 @@
|
|||||||
- обновляем `event_timestamp`;
|
- обновляем `event_timestamp`;
|
||||||
- сохраняем реалистичные связи `event_id <-> location`, `click_id <-> device/geo`.
|
- сохраняем реалистичные связи `event_id <-> location`, `click_id <-> device/geo`.
|
||||||
|
|
||||||
Этого достаточно, чтобы стенд жил часами и данные выглядели не статично.
|
Этого достаточно, чтобы стенд жил часами, выглядел как реальный streaming и не создавал искусственных «пакетов раз в минуту».
|
||||||
|
|
||||||
### 4.1 Минимальная статистическая модель (Poisson)
|
### 4.1 Минимальная статистическая модель (Poisson)
|
||||||
|
|
||||||
Чтобы линия не была «ровной», используем простую интенсивность событий:
|
Чтобы линия не была «ровной», используем простую интенсивность:
|
||||||
|
|
||||||
- число событий на тик: `N_t ~ Poisson(lambda_t)`;
|
- базовая интенсивность в минуту: `lambda_minute` (например, 200 событий/мин);
|
||||||
- базовая интенсивность: `lambda_base` (например, 200 событий/мин);
|
- для текущего тика:
|
||||||
- плавный профиль времени: `lambda_t = lambda_base * hour_factor(t)`.
|
- `lambda_tick = lambda_minute * hour_factor(t) * tick_seconds / 60`;
|
||||||
|
- `N_t ~ Poisson(lambda_tick)`.
|
||||||
|
|
||||||
Где `hour_factor(t)` можно сделать очень простым:
|
Где `hour_factor(t)` можно сделать очень простым:
|
||||||
|
|
||||||
@@ -82,6 +85,8 @@
|
|||||||
|
|
||||||
Итог: поведение уже похоже на живой поток, но код остаётся компактным.
|
Итог: поведение уже похоже на живой поток, но код остаётся компактным.
|
||||||
|
|
||||||
|
Важно: режим «раз в минуту» не удаляем полностью, но рассматриваем только как опцию (`GEN_TICK_SECONDS=60`) для контролируемых демо.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## 5) Как упростить реализацию
|
## 5) Как упростить реализацию
|
||||||
@@ -104,8 +109,7 @@
|
|||||||
|
|
||||||
Через env:
|
Через env:
|
||||||
|
|
||||||
- `GEN_TICK_SECONDS` (по умолчанию `60`)
|
- `GEN_TICK_SECONDS` (по умолчанию `5`)
|
||||||
- `GEN_EVENTS_PER_TICK` (например, `200`)
|
|
||||||
- `GEN_JITTER_PCT` (например, `20`)
|
- `GEN_JITTER_PCT` (например, `20`)
|
||||||
- `GEN_SEED` (для воспроизводимости)
|
- `GEN_SEED` (для воспроизводимости)
|
||||||
- `KAFKA_BOOTSTRAP_SERVERS`
|
- `KAFKA_BOOTSTRAP_SERVERS`
|
||||||
@@ -117,6 +121,7 @@
|
|||||||
|
|
||||||
- `GEN_ENABLED` (быстро включать/выключать цикл)
|
- `GEN_ENABLED` (быстро включать/выключать цикл)
|
||||||
- `GEN_HOUR_PROFILE` (простая карта коэффициентов по часам)
|
- `GEN_HOUR_PROFILE` (простая карта коэффициентов по часам)
|
||||||
|
- `GEN_TICK_SECONDS=60` (опциональный режим «раз в минуту»)
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -188,7 +193,7 @@ MVP успешен, если:
|
|||||||
1. Генератор автономно работает 2-4 часа.
|
1. Генератор автономно работает 2-4 часа.
|
||||||
2. Потребители можно останавливать/запускать отдельно, генератор продолжает работу.
|
2. Потребители можно останавливать/запускать отдельно, генератор продолжает работу.
|
||||||
3. Данные остаются совместимыми с текущим пайплайном.
|
3. Данные остаются совместимыми с текущим пайплайном.
|
||||||
4. Есть базовые метрики и понятные логи.
|
4. Поток публикуется равномерно малыми порциями (без искусственного минутного burst).
|
||||||
5. Есть история batch-ов для учебного разбора.
|
5. Есть история batch-ов для учебного разбора.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|||||||
Reference in New Issue
Block a user