- По умолчанию bash ./scripts/load_kafka_data.sh Resetting topics: browser_events location_events device_events geo_events Loading mode: full (LIMIT=unset) Bootstrap (inside container): kafka:29092 Publishing: data/browser_events.jsonl -> browser_events (full) Publishing: data/device_events.jsonl -> device_events (full) Publishing: data/geo_events.jsonl -> geo_events (full) Publishing: data/location_events.jsonl -> location_events (full) Done. загружает все записи из файлов (вместо 50 строк) - Для ограничения используется bash ./scripts/load_kafka_data.sh - Удалён устаревший параметр Изменённые файлы: - scripts/load_kafka_data.sh — обновлена логика и документация - plans/runbook.md — обновлены примеры использования - plans/kafka_ingest_plan.md — обновлён план реализации Теперь: - bash ./scripts/load_kafka_data.sh Resetting topics: browser_events location_events device_events geo_events Loading mode: full (LIMIT=unset) Bootstrap (inside container): kafka:29092 Publishing: data/browser_events.jsonl -> browser_events (full) Publishing: data/device_events.jsonl -> device_events (full) Publishing: data/geo_events.jsonl -> geo_events (full) Publishing: data/location_events.jsonl -> location_events (full) Done. — все записи (4000 сообщений) - bash ./scripts/load_kafka_data.sh Resetting topics: browser_events location_events device_events geo_events Loading mode: slice (LIMIT=50) Bootstrap (inside container): kafka:29092 Publishing: data/browser_events.jsonl -> browser_events (first 50 lines) Publishing: data/device_events.jsonl -> device_events (first 50 lines) Publishing: data/geo_events.jsonl -> geo_events (first 50 lines) Publishing: data/location_events.jsonl -> location_events (first 50 lines) Done. — 50 строк каждого типа (200 сообщений)
163 lines
8.9 KiB
Markdown
163 lines
8.9 KiB
Markdown
# План реализации: загрузка демо‑данных в 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=<int>` — сколько строк брать из каждого файла. Если не указан — загружаются все записи.
|
||
- `BOOTSTRAP_SERVER=<host:port>` — bootstrap изнутри Kafka‑контейнера (по умолчанию `kafka:29092`).
|
||
|
||
Это обеспечивает:
|
||
|
||
- полная загрузка по умолчанию: `make data` (все записи);
|
||
- быстрый debug: `LIMIT=50 make data`.
|
||
|
||
---
|
||
|
||
## 5) Детальный план работ (что я изменю/добавлю)
|
||
|
||
### 5.1 Добавить/доработать скрипт загрузки
|
||
|
||
Файл: `scripts/load_kafka_data.sh`
|
||
|
||
Что сделаю:
|
||
|
||
1) В начале скрипта зафиксирую defaults:
|
||
- `KAFKA_SERVICE=kafka`
|
||
- `BOOTSTRAP_SERVER=kafka:29092`
|
||
- `RESET_TOPICS=1`, `LIMIT=` (пустое — все записи)
|
||
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) Реализую загрузку:
|
||
- если `LIMIT` не задан → `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;
|
||
- по умолчанию грузит все записи (весь файл);
|
||
- в режиме `LIMIT=50` грузит первые 50 строк (быстрый тест);
|
||
- в режиме `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`) после заливки — опционально.
|
||
|