docs(generator): закрыты находки финального ревью цепочки
- Зачем: - финальный review должен видеть согласованные PRD, issue, курс, Superset и архитектурные документы. - Что: - обновлены PRD, чекбоксы закрытых issue и журнал coordinator-loop. - синхронизированы архитектура, карта репозитория, CONTEXT и курс со startup-history-путём. - убраны старые маркеры Superset-геокарты после перехода на Top Countries. - Проверка: - rg-проверки финального review по PRD, issue и Superset-маркерам. - git diff --cached --check.
This commit is contained in:
+32
-33
@@ -27,12 +27,12 @@
|
||||
flowchart LR
|
||||
subgraph AF["Airflow"]
|
||||
DAG1["ddl_init"]
|
||||
DAG2["kafka_load (bootstrap)"]
|
||||
DAG2["generator_control"]
|
||||
DAG3["etl_pipeline"]
|
||||
end
|
||||
|
||||
subgraph GEN["Generator"]
|
||||
G["generator-service (steady-stream)"]
|
||||
G["generator-service"]
|
||||
end
|
||||
|
||||
subgraph Kafka["Kafka"]
|
||||
@@ -57,8 +57,8 @@ flowchart LR
|
||||
V[витрины VIEW]
|
||||
end
|
||||
|
||||
DAG2 -->|bootstrap JSONL| K
|
||||
G -->|stream events| K
|
||||
DAG2 -->|startup-history backfill/import| K
|
||||
G -->|live после make generator-continue| K
|
||||
K -->|MV| S
|
||||
S -->|batch| O
|
||||
O -->|argMax + JOIN| D1 & D2
|
||||
@@ -69,10 +69,10 @@ flowchart LR
|
||||
DAG3 -.->|batch| ODS & DDS
|
||||
```
|
||||
|
||||
В учебном стенде предусмотрены два пути ingest:
|
||||
В учебном стенде предусмотрены два пути загрузки:
|
||||
|
||||
- `bootstrap`: DAG `kafka_load` для разового/контрольного прогона из `data/*.jsonl`;
|
||||
- `steady-stream`: автономный генератор, публикующий события в Kafka непрерывно.
|
||||
- `startup-history`: DAG `generator_control` создаёт или импортирует историю, запускает ETL и проверяет витрины;
|
||||
- `live`: генератор запускается явно через `make generator-continue`, когда нужна непрерывная подача новых событий.
|
||||
|
||||
### Слои и их назначение
|
||||
|
||||
@@ -80,7 +80,7 @@ flowchart LR
|
||||
flowchart LR
|
||||
subgraph AF["Airflow"]
|
||||
DAG1["ddl_init"]
|
||||
DAG2["kafka_load (bootstrap)"]
|
||||
DAG2["generator_control"]
|
||||
DAG3["etl_pipeline"]
|
||||
end
|
||||
|
||||
@@ -106,8 +106,8 @@ flowchart LR
|
||||
DM_T["VIEW"]
|
||||
end
|
||||
|
||||
DAG2 -->|bootstrap JSONL| KAFKA
|
||||
G -->|steady-stream| KAFKA
|
||||
DAG2 -->|startup-history| KAFKA
|
||||
G -->|live| KAFKA
|
||||
KAFKA -->|MV| STG_T
|
||||
STG_T -->|batch| ODS_T
|
||||
ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T
|
||||
@@ -376,7 +376,6 @@ 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
|
||||
@@ -387,14 +386,13 @@ sequenceDiagram
|
||||
Airflow->>CH: sql/ddl/dm/40_dm.sql
|
||||
CH-->>User: ✅ Структура БД создана
|
||||
|
||||
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 режим
|
||||
alt Startup-history режим
|
||||
User->>Airflow: Trigger generator_control (backfill/import)
|
||||
Airflow->>K: события стартовой истории
|
||||
Airflow->>Airflow: trigger etl_pipeline + check
|
||||
K-->>User: ✅ История в Kafka и витринах
|
||||
else Live режим
|
||||
User->>Compose: make generator-continue
|
||||
loop каждые 1-10 секунд
|
||||
Gen->>K: send N_t (Poisson) в 4 топика
|
||||
end
|
||||
@@ -634,32 +632,33 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
|
||||
|
||||
### Airflow-оркестрация
|
||||
|
||||
Инфраструктура Airflow развёрнута и отвечает за DDL/ETL.
|
||||
Генератор работает отдельно и не управляется через Airflow DAG-и.
|
||||
Инфраструктура Airflow развёрнута и отвечает за DDL/ETL и стартовую историю.
|
||||
Живой генератор контейнеров запускается отдельно через Makefile.
|
||||
|
||||
```python
|
||||
# airflow/dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap)
|
||||
# airflow/dags/kafka_load_dag.py — bootstrap-загрузка JSONL в Kafka (через kafka-python)
|
||||
# airflow/dags/ddl_init_dag.py — создание баз/таблиц
|
||||
# airflow/dags/generator_control_dag.py — backfill/import/check стартовой истории
|
||||
# airflow/dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
|
||||
# airflow/dags/kafka_load_dag.py — архивный ручной путь из JSONL, не основной контур
|
||||
|
||||
# Учебный формат:
|
||||
# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
|
||||
# - SQL-файлы вызываются по фиксированным путям;
|
||||
# - ingest может идти двумя путями:
|
||||
# 1) bootstrap через DAG `kafka_load`;
|
||||
# 2) непрерывный поток через автономный `generator-service`.
|
||||
# - загрузка может идти двумя путями:
|
||||
# 1) startup-history через DAG `generator_control`;
|
||||
# 2) live-поток через явный `make generator-continue`.
|
||||
#
|
||||
# Базовый demo-сценарий:
|
||||
# ddl_init -> kafka_load -> etl_pipeline
|
||||
# ddl_init -> generator_control(backfill/import) -> etl_pipeline -> check
|
||||
# Расширенный учебный сценарий:
|
||||
# generator-service (continuous) + периодический etl_pipeline
|
||||
# make generator-continue + периодический etl_pipeline
|
||||
```
|
||||
|
||||
**DAG `kafka_load`**:
|
||||
- Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python`
|
||||
- TaskGroup `precheck`: проверка Kafka, файлов, параметров
|
||||
- TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка
|
||||
- Параметры: `limit` (0 = все), `reset_topics`
|
||||
**DAG `generator_control`**:
|
||||
- `backfill`: создаёт стартовую историю через генератор
|
||||
- `import`: импортирует портативный артефакт стартовой истории
|
||||
- `check`: сверяет ClickHouse с manifest стартовой истории
|
||||
- После `backfill` и `import` запускает `etl_pipeline` с `full_refresh`
|
||||
|
||||
**Подключение к ClickHouse:**
|
||||
- Connection: `clickhouse_default`
|
||||
|
||||
Reference in New Issue
Block a user