From 4819e10cf89aa01729fb031f203dd95f45ea765f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Feb 2026 17:21:22 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D0=B8=D0=B7=D0=BC=D0=B5=D0=BD=D1=91?= =?UTF-8?q?=D0=BD=20=D1=80=D0=B5=D0=B6=D0=B8=D0=BC=20=D0=B7=D0=B0=D0=B3?= =?UTF-8?q?=D1=80=D1=83=D0=B7=D0=BA=D0=B8=20=D0=B4=D0=B0=D0=BD=D0=BD=D1=8B?= =?UTF-8?q?=D1=85=20=D0=BF=D0=BE=20=D1=83=D0=BC=D0=BE=D0=BB=D1=87=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D1=8E=20=E2=80=94=20=D1=82=D0=B5=D0=BF=D0=B5=D1=80?= =?UTF-8?q?=D1=8C=20=D0=B2=D1=81=D0=B5=20=D0=B7=D0=B0=D0=BF=D0=B8=D1=81?= =?UTF-8?q?=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - По умолчанию 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 сообщений) --- plans/kafka_ingest_plan.md | 17 ++++++++--------- plans/runbook.md | 15 +++++++++------ scripts/load_kafka_data.sh | 20 +++++++------------- 3 files changed, 24 insertions(+), 28 deletions(-) diff --git a/plans/kafka_ingest_plan.md b/plans/kafka_ingest_plan.md index cf1759a..0b6b89f 100644 --- a/plans/kafka_ingest_plan.md +++ b/plans/kafka_ingest_plan.md @@ -72,14 +72,13 @@ Будут поддержаны: - `RESET_TOPICS=1|0` — пересоздавать ли топики перед загрузкой (по умолчанию `1`). -- `LIMIT=` — сколько строк брать из каждого файла в debug‑режиме (по умолчанию `50`). -- `FULL=1|0` — грузить ли весь файл (по умолчанию `0`). Если `FULL=1`, параметр `LIMIT` игнорируется. +- `LIMIT=` — сколько строк брать из каждого файла. Если не указан — загружаются все записи. - `BOOTSTRAP_SERVER=` — bootstrap изнутри Kafka‑контейнера (по умолчанию `kafka:29092`). Это обеспечивает: -- быстрый debug: `LIMIT=50 make data`; -- полная заливка осознанно: `FULL=1 make data`. +- полная загрузка по умолчанию: `make data` (все записи); +- быстрый debug: `LIMIT=50 make data`. --- @@ -94,14 +93,14 @@ 1) В начале скрипта зафиксирую defaults: - `KAFKA_SERVICE=kafka` - `BOOTSTRAP_SERVER=kafka:29092` - - `RESET_TOPICS=1`, `LIMIT=50`, `FULL=0` + - `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) Реализую загрузку: - - если `FULL=1` → `cat file | producer`; + - если `LIMIT` не задан → `cat file | producer` (весь файл); - иначе → `head -n LIMIT file | producer`; - один топик на файл по маппингу. 6) Добавлю минимальные, но полезные сообщения в stdout (что грузим, сколько строк, в какой топик) без лишнего шума. @@ -149,9 +148,9 @@ - В Kafka UI (`localhost:8082`) видны 4 топика и в них появились сообщения. - `make data`: - работает без ClickHouse/DDL; - - по умолчанию грузит небольшой срез (не весь файл); - - в режиме `FULL=1` грузит весь файл; - - в режиме `RESET_TOPICS=1` даёт воспроизводимый “чистый” прогон. + - по умолчанию грузит все записи (весь файл); + - в режиме `LIMIT=50` грузит первые 50 строк (быстрый тест); + - в режиме `RESET_TOPICS=1` даёт воспроизводимый "чистый" прогон. --- diff --git a/plans/runbook.md b/plans/runbook.md index c191645..dd3b5d2 100644 --- a/plans/runbook.md +++ b/plans/runbook.md @@ -56,19 +56,22 @@ make ddl Параметры (env): -- `LIMIT` — сколько строк брать из каждого `.jsonl` (по умолчанию `50`). `LIMIT=0` трактуется как “весь файл”. -- `FULL` — если `FULL=1`, грузит весь файл независимо от `LIMIT`. +- `LIMIT` — сколько строк брать из каждого `.jsonl`. По умолчанию загружаются все записи (весь файл). + Для ограничения используйте `LIMIT=50` или `LIMIT=100`. - `RESET_TOPICS` — если `RESET_TOPICS=1` (по умолчанию), топики удаляются и создаются заново с теми же именами. - `BOOTSTRAP_SERVER` — bootstrap для Kafka *изнутри kafka‑контейнера* (по умолчанию `kafka:29092`). Примеры: ```bash -# 100 строк на поток -LIMIT=100 make data +# Загрузить все данные (по умолчанию) +make data -# Полная заливка всех строк -FULL=1 make data +# Быстрый тест — 50 строк на поток +LIMIT=50 make data + +# Ограниченная загрузка — 100 строк на поток +LIMIT=100 make data # Дозалить данные без пересоздания топиков RESET_TOPICS=0 make data diff --git a/scripts/load_kafka_data.sh b/scripts/load_kafka_data.sh index a416156..59d5265 100644 --- a/scripts/load_kafka_data.sh +++ b/scripts/load_kafka_data.sh @@ -9,9 +9,9 @@ set -euo pipefail # - Kafka CLI запускаем внутри docker‑контейнера `kafka`, а сами файлы читаем на хосте. # # Как запускать: -# make data +# make data # залить файлы целиком (по умолчанию) # LIMIT=100 make data # взять первые 100 строк каждого файла -# FULL=1 make data # залить файлы целиком +# LIMIT=50 make data # взять первые 50 строк каждого файла (быстрый тест) # RESET_TOPICS=0 make data # не пересоздавать топики, а дописать сообщения # # Требования: @@ -22,17 +22,11 @@ COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" KAFKA_SERVICE="${KAFKA_SERVICE:-kafka}" BOOTSTRAP_SERVER="${BOOTSTRAP_SERVER:-kafka:29092}" -# Поведение по умолчанию: пересоздать топики и залить небольшой “срез” данных. +# Поведение по умолчанию: пересоздать топики и залить все данные. RESET_TOPICS="${RESET_TOPICS:-1}" -FULL="${FULL:-0}" -LIMIT="${LIMIT:-50}" - -# Валидация параметров: лучше упасть с понятной ошибкой, чем молча сделать “не то”. -if [[ "${FULL}" != "0" && "${FULL}" != "1" ]]; then - echo "FULL must be 0 or 1 (got: ${FULL})" >&2 - exit 1 -fi +LIMIT="${LIMIT:-}" +# Валидация параметров: лучше упасть с понятной ошибкой, чем молча сделать "не то". if [[ "${RESET_TOPICS}" != "0" && "${RESET_TOPICS}" != "1" ]]; then echo "RESET_TOPICS must be 0 or 1 (got: ${RESET_TOPICS})" >&2 exit 1 @@ -131,11 +125,11 @@ fi # - full: весь файл # - slice: первые N строк (быстрее для отладки) mode="slice" -if [[ "$FULL" == "1" || "$LIMIT" == "0" || -z "$LIMIT" ]]; then +if [[ "$LIMIT" == "0" || -z "$LIMIT" ]]; then mode="full" fi -echo "Loading mode: ${mode} (FULL=${FULL}, LIMIT=${LIMIT:-unset})" +echo "Loading mode: ${mode} (LIMIT=${LIMIT:-unset})" echo "Bootstrap (inside container): ${BOOTSTRAP_SERVER}" for f in "${files[@]}"; do