From cbc2f728716871ce2234cae3fdbef9adbe94a2f8 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 14 Feb 2026 12:33:34 +0300 Subject: [PATCH] =?UTF-8?q?docs(docs):=20=D0=B0=D0=BA=D1=82=D1=83=D0=B0?= =?UTF-8?q?=D0=BB=D0=B8=D0=B7=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=D0=BE=20?= =?UTF-8?q?=D0=A2=D0=97=20=D0=B8=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=D0=B8=20?= =?UTF-8?q?ingest?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - убрать рассинхрон между кратким ТЗ, архитектурой и планом генератора - Что: - сокращен 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 по измененным файлам - в коммит включены только мои документационные изменения --- README.md | 6 ++- docs/ARCHITECTURE.md | 58 +++++++++++++++++------ docs/DE-task.md | 72 ++++++++++++----------------- plans/generator_demo_stream_plan.md | 27 ++++++----- 4 files changed, 95 insertions(+), 68 deletions(-) diff --git a/README.md b/README.md index 63e2f85..3a6e5a5 100644 --- a/README.md +++ b/README.md @@ -9,8 +9,10 @@ витринами. Поток данных коротко: -`data/*.jsonl → Airflow (kafka_load) → Kafka → ClickHouse (слой STG) → Airflow -(etl_pipeline: STG → ODS → DDS → DM) → Superset`. +- **bootstrap**: `data/*.jsonl → Airflow (kafka_load) → Kafka → ClickHouse (слой STG) → + Airflow (etl_pipeline: STG → ODS → DDS → DM) → Superset`. +- **steady-stream**: `generator-service → Kafka → ClickHouse (STG) → Airflow (etl_pipeline) + → Superset`. ## Куда дальше diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index c61195c..e55c15d 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -27,10 +27,14 @@ flowchart LR subgraph AF["Airflow"] DAG1["ddl_init"] - DAG2["kafka_load"] + DAG2["kafka_load (bootstrap)"] DAG3["etl_pipeline"] end + subgraph GEN["Generator"] + G["generator-service (steady-stream)"] + end + subgraph Kafka["Kafka"] K[4 топика] end @@ -53,7 +57,9 @@ flowchart LR V[витрины VIEW] end - DAG2 -->|JSONL| K -->|MV| S + DAG2 -->|bootstrap JSONL| K + G -->|stream events| K + K -->|MV| S S -->|batch| O O -->|argMax + JOIN| D1 & D2 O -.->|ошибки| OE @@ -63,16 +69,25 @@ flowchart LR DAG3 -.->|batch| ODS & DDS ``` +В учебном стенде предусмотрены два пути ingest: + +- `bootstrap`: DAG `kafka_load` для разового/контрольного прогона из `data/*.jsonl`; +- `steady-stream`: автономный генератор, публикующий события в Kafka непрерывно. + ### Слои и их назначение ```mermaid flowchart LR subgraph AF["Airflow"] DAG1["ddl_init"] - DAG2["kafka_load"] + DAG2["kafka_load (bootstrap)"] DAG3["etl_pipeline"] end + subgraph GEN["Generator"] + G["generator-service"] + end + subgraph L1["STG"] KAFKA["Kafka Engine"] STG_T["*_raw таблицы"] @@ -91,7 +106,9 @@ flowchart LR DM_T["VIEW"] end - DAG2 -->|JSONL| KAFKA -->|MV| STG_T + DAG2 -->|bootstrap JSONL| KAFKA + G -->|steady-stream| KAFKA + KAFKA -->|MV| STG_T STG_T -->|batch| ODS_T ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T ODS_T -.->|ошибки| DQ @@ -347,6 +364,7 @@ sequenceDiagram participant User as Пользователь participant Compose as Docker Compose participant Airflow as Airflow + participant Gen as generator-service participant K as Kafka participant CH as ClickHouse participant STG as stg.*_raw @@ -358,6 +376,7 @@ sequenceDiagram Compose->>K: docker compose up -d kafka Compose->>CH: docker compose up -d clickhouse Compose->>Airflow: docker compose up -d airflow-* + Compose->>Gen: docker compose up -d generator Compose-->>User: ✅ Инфраструктура готова User->>Airflow: Trigger ddl_init @@ -368,14 +387,22 @@ sequenceDiagram Airflow->>CH: sql/ddl/dm/40_dm.sql CH-->>User: ✅ Структура БД создана - User->>Airflow: Trigger kafka_load - Airflow->>K: precheck + prepare_topics - loop 4 файла - Airflow->>K: KafkaProducer.send(topic, json_line) + alt Bootstrap режим + User->>Airflow: Trigger kafka_load + Airflow->>K: precheck + prepare_topics + loop 4 файла + Airflow->>K: KafkaProducer.send(topic, json_line) + end + K-->>User: ✅ Данные в Kafka + else Streaming режим + loop каждые 1-10 секунд + Gen->>K: send N_t (Poisson) в 4 топика + end + K-->>User: ✅ Непрерывный поток в Kafka end + K->>CH: Потребление сообщений CH->>STG: INSERT через MV - K-->>User: ✅ Данные в Kafka CH-->>User: ✅ Данные в STG User->>Airflow: Trigger etl_pipeline @@ -607,20 +634,25 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; ### Airflow-оркестрация -Инфраструктура Airflow развёрнута и готова к использованию: +Инфраструктура Airflow развёрнута и отвечает за DDL/ETL. +Генератор работает отдельно и не управляется через Airflow DAG-и. ```python # 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) # Учебный формат: # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; # - SQL-файлы вызываются по фиксированным путям; -# - загрузка данных в Kafka выполняется через DAG `kafka_load`. +# - ingest может идти двумя путями: +# 1) bootstrap через DAG `kafka_load`; +# 2) непрерывный поток через автономный `generator-service`. # -# Основной demo-сценарий: +# Базовый demo-сценарий: # ddl_init -> kafka_load -> etl_pipeline +# Расширенный учебный сценарий: +# generator-service (continuous) + периодический etl_pipeline ``` **DAG `kafka_load`**: diff --git a/docs/DE-task.md b/docs/DE-task.md index 357fba9..c0aa514 100644 --- a/docs/DE-task.md +++ b/docs/DE-task.md @@ -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** -* Postgresql + dagster + Grafana -* dlt + duckdb + Airflow + [Panel](https://panel.holoviz.org/) -* Clickhouse + Airflow + Redis + [Streamlit](https://streamlit.io/) -* sqlitedb + скрипты питона + Flask + d3.js +- "Грязные" записи не должны валить пайплайн: ошибки фиксируются в ODS/DQ-слое. +- Генератор работает отдельно от потребителей и Airflow-триггеров. +- Решение должно быть воспроизводимым и пригодным для учебной демонстрации. -Начните с маленького среза данных, чтобы ознакомиться, что за данные вам предоставляют: +## Границы MVP -* browser_events.jsonl — данные с информацией о просмотрах браузера -* device_events.jsonl — данные об устройствах, с которых пользователи пользовались сайтом -* geo_events.jsonl — данные о локации пользователя -* location_events.jsonl — вопреки названию, это данные не о локации пользователя, а непосредственно о положении пользователя на сайте и информации о том, откуда пользователь попал на страницу +- Один простой режим генератора (`steady-stream`) без сложных сценариев. +- Без production-гарантий уровня exactly-once. +- Фокус на рабочем end-to-end контуре, а не на полной симуляции реального продакшена. -Смысл задачи на **своей ВМ** развернуть инструменты и загрузить данные в Kafka и из нее прогнать их в Clickhouse /другая DB по слоям. +## Критерии готовности -В Airflow взять данные из таблиц и преобразовать, потом выгрузить в Clickhouse / другая DB в агрегированном формате под дашборд. +1. Стенд стабильно запускается и работает несколько часов. +2. Данные корректно проходят путь Kafka -> STG -> ODS -> DDS -> DM. +3. Поток событий в Kafka идёт постепенно малыми порциями, без искусственного минутного burst. +4. Есть базовый дашборд и метрики для разбора работы стенда. -Дашборд не оценивается, но будет плюсом, если получиться собрать удобный визуал. +## Где детали реализации -Записать можно на видео и отправить контактному лицу или продемонстрировать на собеседовании - -Можно использовать любые инструменты. - -Данные могут быть(скорее всего) с ошибками и компания которая их предоставляем заранее об этом сообщила т.к. они не хотят делиться 100% качественными данными полноценно с продакшена. - -### ИНФРАСТРУКТУРА - -* инфраструктура клевая и работает = 18 баллов -* инфраструктура работает, но есть баги = 12 баллов -* инфраструктура почти есть, но не совсем = 6 баллов -* инфраструктуры нет = 0 баллов - -### ДАШБОРДЫ - -* дашборды есть и удобные и работают = 12 баллов -* дашборды работают, но неудобные = 8 баллов -* дашборды почти работают = 4 балла -* дашбордов нет = 0 баллов +- Техническая архитектура: `docs/ARCHITECTURE.md` +- План по генератору: `plans/generator_demo_stream_plan.md` +- Эксплуатация и запуск: `docs/OPERATIONS.md` diff --git a/plans/generator_demo_stream_plan.md b/plans/generator_demo_stream_plan.md index ed547c6..4f5f32e 100644 --- a/plans/generator_demo_stream_plan.md +++ b/plans/generator_demo_stream_plan.md @@ -1,4 +1,4 @@ -# План реализации: простой автономный генератор (rev4) +# План реализации: простой автономный генератор (rev5) Дата ревизии: 14 февраля 2026. @@ -44,11 +44,13 @@ --- -## 4) Единственный режим генерации: `steady` +## 4) Единственный режим генерации: `steady-stream` Логика режима: -- каждую минуту публикуем фиксированный объём событий; +- публикуем в Kafka **постепенно**, короткими тиками (например, каждые 1-10 секунд); +- на каждом тике отправляем небольшую порцию сообщений; +- держим целевую интенсивность в `events/min`, а не крупный минутный batch; - распределяем события по 4 топикам: - `browser_events` - `location_events` @@ -59,15 +61,16 @@ - обновляем `event_timestamp`; - сохраняем реалистичные связи `event_id <-> location`, `click_id <-> device/geo`. -Этого достаточно, чтобы стенд жил часами и данные выглядели не статично. +Этого достаточно, чтобы стенд жил часами, выглядел как реальный streaming и не создавал искусственных «пакетов раз в минуту». ### 4.1 Минимальная статистическая модель (Poisson) -Чтобы линия не была «ровной», используем простую интенсивность событий: +Чтобы линия не была «ровной», используем простую интенсивность: -- число событий на тик: `N_t ~ Poisson(lambda_t)`; -- базовая интенсивность: `lambda_base` (например, 200 событий/мин); -- плавный профиль времени: `lambda_t = lambda_base * hour_factor(t)`. +- базовая интенсивность в минуту: `lambda_minute` (например, 200 событий/мин); +- для текущего тика: + - `lambda_tick = lambda_minute * hour_factor(t) * tick_seconds / 60`; + - `N_t ~ Poisson(lambda_tick)`. Где `hour_factor(t)` можно сделать очень простым: @@ -82,6 +85,8 @@ Итог: поведение уже похоже на живой поток, но код остаётся компактным. +Важно: режим «раз в минуту» не удаляем полностью, но рассматриваем только как опцию (`GEN_TICK_SECONDS=60`) для контролируемых демо. + --- ## 5) Как упростить реализацию @@ -104,8 +109,7 @@ Через env: -- `GEN_TICK_SECONDS` (по умолчанию `60`) -- `GEN_EVENTS_PER_TICK` (например, `200`) +- `GEN_TICK_SECONDS` (по умолчанию `5`) - `GEN_JITTER_PCT` (например, `20`) - `GEN_SEED` (для воспроизводимости) - `KAFKA_BOOTSTRAP_SERVERS` @@ -117,6 +121,7 @@ - `GEN_ENABLED` (быстро включать/выключать цикл) - `GEN_HOUR_PROFILE` (простая карта коэффициентов по часам) +- `GEN_TICK_SECONDS=60` (опциональный режим «раз в минуту») --- @@ -188,7 +193,7 @@ MVP успешен, если: 1. Генератор автономно работает 2-4 часа. 2. Потребители можно останавливать/запускать отдельно, генератор продолжает работу. 3. Данные остаются совместимыми с текущим пайплайном. -4. Есть базовые метрики и понятные логи. +4. Поток публикуется равномерно малыми порциями (без искусственного минутного burst). 5. Есть история batch-ов для учебного разбора. ---