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.
9.0 KiB
9.0 KiB
План реализации: загрузка демо‑данных в Kafka (из data/*.jsonl)
Цель: реализовать воспроизводимую загрузку демо‑данных в Kafka из файлов *.jsonl так, чтобы дальше их мог забрать ClickHouse (или любой другой потребитель) в режиме “1 Kafka message = 1 event (JSON)”.
Документ описывает что именно я буду делать в репозитории, какие решения принимаются и как это будет проверяться.
1) Цели и ограничения
Цели
- Поддержать команду
make data, которая:- (по умолчанию) пересоздаёт топики с теми же именами (сброс данных для повторяемых прогонов);
- загружает данные из
data/*.jsonlв Kafka.
- Сделать загрузку устойчивой к “грязным” данным: скрипт не парсит JSON, он публикует строки как есть.
- Дать два режима:
- debug‑срез (для отладки) — первые
Nстрок каждого файла; - full — весь файл целиком.
- debug‑срез (для отладки) — первые
- Не завязывать загрузку на ClickHouse/DDL: если пользователь не делал
make ddl,make dataвсё равно должен работать.
Ограничения (из правил репо)
- Не грузить
*.jsonlцеликом “по умолчанию”. Полная заливка — отдельный осознанный режим.
2) Контракт данных (что кладём в Kafka)
- Сообщение Kafka:
value= ровно одна строка из.jsonl(один JSON‑объект).- Никаких “обёрток” (массивов), никаких доп.полей от скрипта.
- Кодировка: UTF‑8 (как в исходных файлах).
- Разбиение на топики (фиксированное):
data/browser_events.jsonl→browser_eventsdata/location_events.jsonl→location_eventsdata/device_events.jsonl→device_eventsdata/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>— сколько строк брать из каждого файла в debug‑режиме (по умолчанию50).FULL=1|0— грузить ли весь файл (по умолчанию0). ЕслиFULL=1, параметрLIMITигнорируется.BOOTSTRAP_SERVER=<host:port>— bootstrap изнутри Kafka‑контейнера (по умолчаниюkafka:29092).
Это обеспечивает:
- быстрый debug:
LIMIT=50 make data; - полная заливка осознанно:
FULL=1 make data.
5) Детальный план работ (что я изменю/добавлю)
5.1 Добавить/доработать скрипт загрузки
Файл: scripts/load_kafka_data.sh
Что сделаю:
- В начале скрипта зафиксирую defaults:
KAFKA_SERVICE=kafkaBOOTSTRAP_SERVER=kafka:29092RESET_TOPICS=1,LIMIT=50,FULL=0
- Добавлю проверку наличия исходных файлов
data/*_events.jsonlи понятные ошибки. - Реализую обнаружение Kafka CLI внутри контейнера (через
command -v), чтобы не завязываться на конкретный путь. - Реализую “reset” топиков:
- удалить (
--delete --if-exists) 4 фиксированных топика; - создать (
--create --if-not-exists) эти же топики сpartitions=1,replication-factor=1.
- удалить (
- Реализую загрузку:
- если
FULL=1→cat file | producer; - иначе →
head -n LIMIT file | producer; - один топик на файл по маппингу.
- если
- Добавлю минимальные, но полезные сообщения в 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и “дозаливка”, либо ручная очистка.
- скрипт сообщит об ошибке; workaround для демо:
- Повторная загрузка без 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) после заливки — опционально.