From 380e8fcff3d5e74cffa3cbd1b7fd15f83425ad92 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 5 Feb 2026 21:49:09 +0300 Subject: [PATCH] docs(plans): fix kafka broker address and add data loading docs 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. --- plans/clickhouse_ddl.md | 10 ++- plans/kafka_ingest_plan.md | 163 +++++++++++++++++++++++++++++++++++++ plans/runbook.md | 82 +++++++++++++++++++ 3 files changed, 251 insertions(+), 4 deletions(-) create mode 100644 plans/kafka_ingest_plan.md create mode 100644 plans/runbook.md diff --git a/plans/clickhouse_ddl.md b/plans/clickhouse_ddl.md index f0c02c4..33c75c4 100644 --- a/plans/clickhouse_ddl.md +++ b/plans/clickhouse_ddl.md @@ -182,13 +182,15 @@ ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); ```sql -- Пример: одна колонка raw, один message = одна строка. -- Замените broker/topic/group под вашу инфраструктуру. +-- Если ClickHouse запущен в docker compose в одной сети с Kafka — обычно это `kafka:29092`. +-- Если ClickHouse подключается к Kafka с хоста — обычно это `localhost:9092`. CREATE TABLE IF NOT EXISTS stg.kafka_browser_raw ( raw String ) ENGINE = Kafka SETTINGS - kafka_broker_list = 'kafka:9092', + kafka_broker_list = 'kafka:29092', kafka_topic_list = 'browser_events', kafka_group_name = 'ch_stg_browser', kafka_format = 'JSONAsString', @@ -214,7 +216,7 @@ FROM stg.kafka_browser_raw; CREATE TABLE IF NOT EXISTS stg.kafka_location_raw (raw String) ENGINE = Kafka SETTINGS - kafka_broker_list = 'kafka:9092', + kafka_broker_list = 'kafka:29092', kafka_topic_list = 'location_events', kafka_group_name = 'ch_stg_location', kafka_format = 'JSONAsString', @@ -224,7 +226,7 @@ SETTINGS CREATE TABLE IF NOT EXISTS stg.kafka_device_raw (raw String) ENGINE = Kafka SETTINGS - kafka_broker_list = 'kafka:9092', + kafka_broker_list = 'kafka:29092', kafka_topic_list = 'device_events', kafka_group_name = 'ch_stg_device', kafka_format = 'JSONAsString', @@ -234,7 +236,7 @@ SETTINGS CREATE TABLE IF NOT EXISTS stg.kafka_geo_raw (raw String) ENGINE = Kafka SETTINGS - kafka_broker_list = 'kafka:9092', + kafka_broker_list = 'kafka:29092', kafka_topic_list = 'geo_events', kafka_group_name = 'ch_stg_geo', kafka_format = 'JSONAsString', diff --git a/plans/kafka_ingest_plan.md b/plans/kafka_ingest_plan.md new file mode 100644 index 0000000..cf1759a --- /dev/null +++ b/plans/kafka_ingest_plan.md @@ -0,0 +1,163 @@ +# План реализации: загрузка демо‑данных в Kafka (из `data/*.jsonl`) + +Цель: реализовать воспроизводимую загрузку демо‑данных в Kafka из файлов `*.jsonl` так, чтобы дальше их мог забрать ClickHouse (или любой другой потребитель) в режиме “1 Kafka message = 1 event (JSON)”. + +Документ описывает **что именно я буду делать** в репозитории, какие решения принимаются и как это будет проверяться. + +--- + +## 1) Цели и ограничения + +### Цели + +- Поддержать команду `make data`, которая: + - (по умолчанию) **пересоздаёт** топики с теми же именами (сброс данных для повторяемых прогонов); + - загружает данные из `data/*.jsonl` в Kafka. +- Сделать загрузку **устойчивой к “грязным” данным**: скрипт не парсит JSON, он публикует строки как есть. +- Дать два режима: + - **debug‑срез** (для отладки) — первые `N` строк каждого файла; + - **full** — весь файл целиком. +- Не завязывать загрузку на ClickHouse/DDL: если пользователь не делал `make ddl`, `make data` всё равно должен работать. + +### Ограничения (из правил репо) + +- Не грузить `*.jsonl` целиком “по умолчанию”. Полная заливка — отдельный осознанный режим. + +--- + +## 2) Контракт данных (что кладём в Kafka) + +- Сообщение Kafka: + - `value` = **ровно одна строка** из `.jsonl` (один JSON‑объект). + - Никаких “обёрток” (массивов), никаких доп.полей от скрипта. +- Кодировка: UTF‑8 (как в исходных файлах). +- Разбиение на топики (фиксированное): + - `data/browser_events.jsonl` → `browser_events` + - `data/location_events.jsonl` → `location_events` + - `data/device_events.jsonl` → `device_events` + - `data/geo_events.jsonl` → `geo_events` + +Почему так: это напрямую соответствует описанию в `plans/clickhouse_ddl.md` и позволяет ClickHouse читать поток “raw String”. + +--- + +## 3) Kafka‑доступ и утилиты (как будем публиковать) + +Решение: **Kafka CLI запускаем внутри Kafka‑контейнера**, а управляем процессом с хоста bash‑скриптом. + +Почему: + +- не нужно ставить `kcat`/Java на хост; +- меньше “локальных зависимостей”; +- bootstrap внутри сети compose стабильный: `kafka:29092`. + +Технически: + +- Внутри контейнера используем: + - `kafka-topics.sh` — создание/удаление топиков; + - `kafka-console-producer.sh` — публикация сообщений. +- Хостовый скрипт будет делать `docker compose exec -T kafka ...` и подавать файл/`head` в stdin producer’а. + +--- + +## 4) UX/интерфейс (как пользователь запускает) + +### Команды + +- `make up` — поднимает стек. +- `make data` — загружает данные в Kafka (должна работать при поднятой Kafka, без ClickHouse). + +### Параметры (через env) + +Будут поддержаны: + +- `RESET_TOPICS=1|0` — пересоздавать ли топики перед загрузкой (по умолчанию `1`). +- `LIMIT=` — сколько строк брать из каждого файла в debug‑режиме (по умолчанию `50`). +- `FULL=1|0` — грузить ли весь файл (по умолчанию `0`). Если `FULL=1`, параметр `LIMIT` игнорируется. +- `BOOTSTRAP_SERVER=` — bootstrap изнутри Kafka‑контейнера (по умолчанию `kafka:29092`). + +Это обеспечивает: + +- быстрый debug: `LIMIT=50 make data`; +- полная заливка осознанно: `FULL=1 make data`. + +--- + +## 5) Детальный план работ (что я изменю/добавлю) + +### 5.1 Добавить/доработать скрипт загрузки + +Файл: `scripts/load_kafka_data.sh` + +Что сделаю: + +1) В начале скрипта зафиксирую defaults: + - `KAFKA_SERVICE=kafka` + - `BOOTSTRAP_SERVER=kafka:29092` + - `RESET_TOPICS=1`, `LIMIT=50`, `FULL=0` +2) Добавлю проверку наличия исходных файлов `data/*_events.jsonl` и понятные ошибки. +3) Реализую обнаружение Kafka CLI внутри контейнера (через `command -v`), чтобы не завязываться на конкретный путь. +4) Реализую “reset” топиков: + - удалить (`--delete --if-exists`) 4 фиксированных топика; + - создать (`--create --if-not-exists`) эти же топики с `partitions=1`, `replication-factor=1`. +5) Реализую загрузку: + - если `FULL=1` → `cat file | producer`; + - иначе → `head -n LIMIT file | producer`; + - один топик на файл по маппингу. +6) Добавлю минимальные, но полезные сообщения в stdout (что грузим, сколько строк, в какой топик) без лишнего шума. + +### 5.2 Подключить в `make data` + +Файл: `Makefile` + +Что сделаю: + +- Добавлю (или оставлю) таргет `data`, который вызывает `bash ./scripts/load_kafka_data.sh`. + +### 5.3 Документация “для пользователя” + +Файлы: + +- `plans/runbook.md` — оставить короткую инструкцию “как запустить”. +- `plans/kafka_ingest_plan.md` (этот файл) — план реализации/поддержки. + +Что сделаю: + +- В runbook’е зафиксирую: + - что `make data` не зависит от `make ddl`; + - какие env‑параметры есть и как включить `FULL=1`. + +--- + +## 6) Краевые случаи и как их обработаем + +- **Docker daemon недоступен** (permission denied к `/var/run/docker.sock`): + - скрипты не чинят это автоматически, но выводят понятную ошибку; в runbook указать, что нужно настроить доступ. +- **Топики не удаляются** (настройки брокера): + - скрипт сообщит об ошибке; workaround для демо: `RESET_TOPICS=0` и “дозаливка”, либо ручная очистка. +- **Повторная загрузка без reset**: + - это осознанно добавит сообщения (возможны дубликаты), что нормально для демо. +- **“Грязный” JSON / невалидные строки**: + - мы не валидируем JSON, публикуем строку как есть; обработка ошибок — ответственность downstream (например, ODS в ClickHouse). + +--- + +## 7) Проверка результата (acceptance criteria) + +Минимальный набор проверок после `make data`: + +- В Kafka UI (`localhost:8082`) видны 4 топика и в них появились сообщения. +- `make data`: + - работает без ClickHouse/DDL; + - по умолчанию грузит небольшой срез (не весь файл); + - в режиме `FULL=1` грузит весь файл; + - в режиме `RESET_TOPICS=1` даёт воспроизводимый “чистый” прогон. + +--- + +## 8) Улучшения (опционально, позже) + +- Поддержка key для партиционирования (например, `event_id`/`click_id`) через `kafka-console-producer` с `parse.key=true` и разделителем — **но это потребует менять формат input**, поэтому не делаем в MVP. +- “Согласованный срез” (фильтрация связанных потоков по `event_id/click_id`), чтобы на маленьком N джойны лучше сходились — опционально. +- Автоматический smoke‑check consumer’ом (`kafka-console-consumer --max-messages`) после заливки — опционально. + diff --git a/plans/runbook.md b/plans/runbook.md new file mode 100644 index 0000000..f9c576b --- /dev/null +++ b/plans/runbook.md @@ -0,0 +1,82 @@ +# Runbook: запуск демо и загрузка данных + +Этот документ фиксирует порядок действий и `make`‑таргеты. Он не описывает внутренности ClickHouse‑слоёв (это в `plans/clickhouse_ddl.md`). + +## Предпосылки + +- Docker + Docker Compose. +- Доступ к Docker daemon (если `docker compose ...` пишет `permission denied ... /var/run/docker.sock`, добавьте пользователя в группу `docker` или запускайте команды с правами, принятыми в вашей среде). + +## Быстрый сценарий + +1) Поднять инфраструктуру: + +```bash +make up +``` + +2) (Опционально) Залить данные в Kafka: + +```bash +make data +``` + +`make data` не зависит от ClickHouse/DDL — достаточно, чтобы Kafka была поднята. + +3) Применить DDL в ClickHouse: + +```bash +make ddl +``` + +## Make таргеты + +- `make up` — `docker compose up -d` (поднимает весь стек из `docker-compose.yml`). +- `make ddl` — применяет SQL из `plans/clickhouse_ddl.md` в контейнер ClickHouse (извлекает все блоки ```sql``` и исполняет их через `clickhouse-client`). +- `make data` — пересоздаёт топики (по умолчанию) и публикует события из `data/*.jsonl` в Kafka (1 строка = 1 Kafka message value). + +## Загрузка данных в Kafka (`make data`) + +### Топики + +Скрипт использует фиксированный маппинг: + +- `data/browser_events.jsonl` → `browser_events` +- `data/location_events.jsonl` → `location_events` +- `data/device_events.jsonl` → `device_events` +- `data/geo_events.jsonl` → `geo_events` + +### Режимы загрузки + +- По умолчанию — “debug срез”: первые 50 строк каждого файла. +- Полная загрузка — весь файл. + +Параметры (env): + +- `LIMIT` — сколько строк брать из каждого `.jsonl` (по умолчанию `50`). `LIMIT=0` трактуется как “весь файл”. +- `FULL` — если `FULL=1`, грузит весь файл независимо от `LIMIT`. +- `RESET_TOPICS` — если `RESET_TOPICS=1` (по умолчанию), топики удаляются и создаются заново с теми же именами. +- `BOOTSTRAP_SERVER` — bootstrap для Kafka *изнутри kafka‑контейнера* (по умолчанию `kafka:29092`). + +Примеры: + +```bash +# 100 строк на поток +LIMIT=100 make data + +# Полная заливка всех строк +FULL=1 make data + +# Дозалить данные без пересоздания топиков +RESET_TOPICS=0 make data +``` + +## Применение DDL в ClickHouse (`make ddl`) + +Скрипт исполняет SQL из `plans/clickhouse_ddl.md`. Для `ENGINE = Kafka` важно, чтобы `kafka_broker_list` был доступен из контейнера ClickHouse. + +В текущем compose: + +- для соединений “контейнер → Kafka” используйте `kafka:29092`; +- `localhost:9092` подходит только для клиентов на хосте. +