- Зачем:
- нужно зафиксировать архитектуру перехода от full refresh к инкрементальной загрузке
- план служит референсом для реализации и code review
- Что:
- добавлен plans/incremental-etl-v2.md с полным описанием:
* архитектура watermark (event_ts + lookback 5min)
* DDL для meta.etl_watermarks_history с TTL 30 дней
* SQL шаблоны для всех слоев (ODS, DDS, DM)
* структура DAG с параллельной загрузкой ODS
* функции get_watermark и save_watermark
* алерты Grafana для late arrivals
* демонстрация late arrivals через комментарии и логи
* сценарии тестирования
* оценка трудозатрат (15-16 часов)
- Проверка:
- файл создан: plans/incremental-etl-v2.md
- структура соответствует принятой архитектуре
- все параметры согласованы (TTL, lookback, параллельность)
- Зачем:
- зафиксировать решение об observability генератора на уровне MVP-плана
- Что:
- добавлены требования по /metrics, scrape_config и target generator:9109
- добавлен env-параметр GEN_METRICS_PORT
- обновлены шаг внедрения и критерии успеха
- Проверка:
- проверен diff только для plans/generator_demo_stream_plan.md
- Зачем:
- убрать рассинхрон между кратким ТЗ, архитектурой и планом генератора
- Что:
- сокращен 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 по измененным файлам
- в коммит включены только мои документационные изменения
- Зачем:
- зафиксировать реалистичный MVP без переусложнения
- Что:
- оставлен один режим steady для автономного генератора
- добавлена минимальная статистическая модель потока на базе Poisson
- уточнены минимальные метрики, история batch и короткий roadmap внедрения
- Проверка:
- проверен diff и итоговое содержимое plans/generator_demo_stream_plan.md
- Why:
- Airflow task metrics were mapped to non-emitted StatsD keys
- reload-monitoring did not restart statsd-exporter after mapping changes
- What:
- update StatsD mapping for Airflow 2.10.5 metric names
- remove problematic catch-all mapping that produced inconsistent series
- restart statsd-exporter in reload-monitoring flow
- sync operations runbook and airflow monitoring plan with actual metrics
- Check:
- make reload-monitoring
- Prometheus targets: airflow/clickhouse/kafka are UP
- trigger ddl_init and verify airflow_task_duration_seconds_count
- verify airflow_task_success_total and airflow_task_failures_total in Prometheus
- Add statsd-exporter service to docker-compose.yml (prom/statsd-exporter:v0.27.1)
- Add StatsD env vars to airflow-default-env for metrics export
- Add airflow job to prometheus.yml scrape configs
- Add Airflow Overview dashboard (Grafana provisioning)
- Add Airflow alert rules: scheduler down, queue backlog, failures, parse time
- Add configs/statsd_mapping.yml for StatsD → Prometheus conversion
- Use Prometheus naming convention (_total for counters, _seconds for timers)
- Add monitoring plan at plans/monitoring_airflow_plan.md
- Update OPERATIONS.md and Makefile for airflow monitoring
Tested: all 3 jobs (airflow, clickhouse, kafka) showing UP in Prometheus,
metrics flowing (dagbag_size=3, executor slots, heartbeats with _total suffix),
all 4 alert rules loaded in Grafana
- Add kafka-exporter service to docker-compose.yml
- Add kafka job to prometheus.yml scrape configs
- Add Kafka Overview dashboard (Grafana provisioning)
- Add Kafka alert rules (broker down, consumer lag, etc.)
- Add make reload-monitoring command for easy updates
- Update OPERATIONS.md with TL;DR and troubleshooting
API verified via Context7:
- /danielqsj/kafka_exporter for exporter config
- /prometheus/docs for scrape_configs format
- Why:
- keep Airflow artifacts under a single airflow/ directory
- align repository layout with intended project structure
- What:
- move dags/ to airflow/dags/ and update compose mounts
- make SQL root resolution work in container and local runs
- update DAG path references in README, AGENTS, ARCHITECTURE, and plans
- remove tracked Python cache artifacts from old DAG location
- Check:
- airflow dags list
- airflow dags list-import-errors
- e2e success: ddl_init, kafka_load(limit=50), etl_pipeline
- Why:\n - User-facing docs mixed Airflow and legacy CLI ingest paths and caused confusion\n- What:\n - Rework README quick start and status to use DAG chain ddl_init -> kafka_load -> etl_pipeline\n - Rewrite runbook as canonical Airflow-first execution flow\n - Sync architecture diagrams/sequence and DQ wording with current SQL and DAG behavior\n- Check:\n - Verified updated sections and removed stale markers with rg in README.md, docs/ARCHITECTURE.md, plans/runbook.md
- Why:
- For DE task we only need full ingest or limit-based sample.
- load_* and full_load params were redundant and unclear in current flow.
- What:
- Remove full_load and load_* params from kafka_load DAG contract.
- Simplify kafka helpers (validate/check files) to fixed 4-stream ingest.
- Sync AGENTS, README, runbook, architecture and airflow plan docs.
- Check:
- python3 -m py_compile dags/kafka_load_dag.py dags/utils/kafka_helpers.py
- Airflow smoke/full runs: ddl_init -> kafka_load -> etl_pipeline (all success).
- Legacy path: make data && make transform (success).
Move DDL files from flat ddl/ directory to sql/ddl/ with layer-based
subdirectories (stg, ods, dds, dm). Move batch transformation SQL from
jobs/ to sql/ layer directories. Update scripts and documentation to
reflect new paths for improved organization and Airflow integration.
Split Kafka ingestion into a dedicated `kafka_load` DAG to enable independent
experimentation with data loading without triggering the full ETL pipeline.
Restructure implementation phases: Stage 1 uses `make data` for MVP, Stage 2
adds the standalone Kafka DAG, Stage 3 adds monitoring. Update DAG numbering,
parameters, task groups, and acceptance criteria to reflect the new
architecture.
Consolidate the orchestration strategy by merging `kafka_load` and
`etl_batch_transform` into a unified `etl_pipeline`. Replace BashOperator
dependencies on Kafka CLI with PythonOperators utilizing `kafka-python`.
Add detailed technical specifications for helper functions, MVP stages,
and validation checks to align with current infrastructure constraints.
Detail the architecture for migrating ETL orchestration from make to
Airflow. Define DAG structures for database initialization, Kafka data
ingestion, batch transformation, and quality monitoring. Include
technical specifications, operator details, and implementation phases.
Replace materialized view joins with batch SQL transformations to avoid
consistency issues with out-of-order data. Document the reasoning for
using batch processing for ODS to DDS layer, including handling of
eventual consistency and versioning in ReplacingMergeTree. Update
data flow diagrams and remove MV creation DDL for DDS tables. Add
documentation for error handling tables and batch transformation jobs.
Add comprehensive plan for migrating executable DDL statements from markdown
to separate SQL files organized by layer. The plan outlines artifact structure,
execution requirements via make/Airflow, and environment parameters.
Existing inline DDL content is now marked as legacy in an appendix section,
providing clear separation between planned implementation and current state.
Add build automation via Makefile with targets for docker compose
management, DDL application, and data ingestion. Implement a robust bash
script for loading JSONL demo data into Kafka topics with configurable
options for limits, full dataset loading, and topic reset behavior.
Update Kafka broker address to use internal Docker network port (29092)
instead of external port (9092). Add kafka_ingest_plan.md with detailed
implementation strategy and runbook.md with user instructions.
- Описаны слои STG/ODS/DDS/DM и связи потоков (`event_id`/`click_id`)
- Добавлены DDL и MV-пайплайн для ingestion из Kafka (ClickHouse) + типизация/дедуп/DQ
- Добавлены витрины/VIEW для BI (Superset), mermaid-диаграмма и операционные заметки