diff --git a/MENTEE_GUIDE_JSON_FLOWS.md b/MENTEE_GUIDE_JSON_FLOWS.md index d4315cb..42fd808 100644 --- a/MENTEE_GUIDE_JSON_FLOWS.md +++ b/MENTEE_GUIDE_JSON_FLOWS.md @@ -1,5 +1,7 @@ # Гайд для менти: как работают примеры flow на JSON (NiFi → Kafka → Postgres) +> **Последнее обновление:** 2026-01-08 | **Версия NiFi:** 1.x.x | ← [Вернуться в README](README.md) + Этот гайд про два flow definition файла: - `nifi-templates/Sample2Kafka.json` — генерирует записи и публикует их в Kafka topic `Sample2Kafka` в виде JSON. - `nifi-templates/SampleKafka2Postgres.json` — читает JSON из Kafka topic `Sample2Kafka` и пишет данные в Postgres. @@ -13,9 +15,36 @@ Это учебный пример, который показывает базовые роли технологий: - **NiFi** — визуальный конвейер (pipeline) данных. -- **Kafka** — “почтовый ящик”/шина сообщений между системами. +- **Kafka** — "почтовый ящик"/шина сообщений между системами. - **Postgres** — база данных, где данные хранятся и доступны SQL-запросами. +### Визуализация потока данных + +```mermaid +graph TD + subgraph Flow1["Flow 1: NiFi → Kafka"] + GR["GenerateRecord
(генерация)"] + PK["PublishKafkaRecord_2_6
(публикация в Kafka)"] + KT["Kafka Topic: Sample2Kafka"] + + GR --> PK --> KT + end + + subgraph Flow2["Flow 2: Kafka → Postgres"] + CK["ConsumeKafkaRecord_2_6
(чтение из Kafka)"] + QR["QueryRecord
(фильтрация по GOOD_DATE)"] + MC["MergeContent
(micro-batch)"] + PDR["PutDatabaseRecord
(вставка в stg)"] + ES["ExecuteSQL
(вызов процедуры + очистка)"] + + KT --> CK --> QR --> MC --> PDR --> ES + end + + style Flow1 fill:#e1f5ff,stroke:#01579b + style Flow2 fill:#fff3e0,stroke:#e65100 + style KT fill:#fff9c4,stroke:#fbc02d +``` + ## 1) Мини-словарь NiFi (3 понятия) - **FlowFile** — “сообщение”, которое течёт по графу. У FlowFile есть *content* (тело) и *attributes* (метаданные). - **Processor** — шаг обработки (прочитать/преобразовать/записать). @@ -26,16 +55,27 @@ - `dttm` — дата/время (timestamp). - `txt` — строка (text). +Пример JSON-записи: +```json +{ + "dttm": "2026-01-08T11:49:52Z", + "txt": "Пример текста" +} +``` + В первом flow запись создаётся процессором `GenerateRecord` по схеме (она задана в настройке процессора). Важно: в этом учебном варианте **мы не используем Avro** и не используем Schema Registry. Сообщения в Kafka — это обычный JSON. ## 3) Перед стартом: что проверить 1) Стенд поднят: `docker compose up -d` -2) Открывается NiFi: http://localhost:18443/nifi/ (логин `admin`, пароль `Password123456`) +2) Открывается NiFi: http://localhost:18443/nifi/ + > 💡 **Примечание:** В текущей конфигурации аутентификация отключена для упрощения обучения. Логин и пароль не требуются. 3) Открывается Kafka UI: http://localhost:8082/ -4) (Если это первый запуск после `docker compose down -v`) примените SQL-скрипт для демо-таблиц: +4) **Примените SQL-скрипт для демо-таблиц** (требуется при первом запуске стенда или после `docker compose down -v`): - `docker compose exec -T postgres psql -U postgres -d app -f /nifi-templates/SampleKafka2Postgres.sql` + + > ⚠️ **Важно:** Скрипт не идемпотентный — при повторном запуске будут ошибки про существующие схемы/таблицы. ## 4) Flow 1 — `Sample2Kafka.json` (NiFi → Kafka) Цель: регулярно публиковать новые JSON-сообщения в Kafka topic `Sample2Kafka`. @@ -75,11 +115,12 @@ - В примере есть запрос `GOOD_DATE` (это имя relationship), который отбирает записи по условию на `dttm`. - На выходе получаются записи, которые прошли фильтр. 3) `MergeContent` (micro-batch) - - Собирает несколько FlowFile в “пачку” перед загрузкой в БД. + - Собирает несколько FlowFile в "пачку" перед загрузкой в БД. - Зачем: много маленьких вставок (каждое сообщение → отдельный insert) может перегружать БД; batching снижает накладные расходы. - Важные настройки: - - `Minimum Number of Entries` / `Maximum Number of Entries` — размер пачки (сколько записей объединяем). - - `Max Bin Age` — максимальное ожидание, чтобы пачка всё равно отправилась даже если данных мало. + - `Minimum Number of Entries`: `100` — минимальный размер пачки + - `Maximum Number of Entries`: `1000` — максимальный размер пачки + - `Max Bin Age`: `30 seconds` — максимальное ожидание, чтобы пачка всё равно отправилась даже если данных мало - Компромисс: чем больше batch, тем меньше нагрузка на БД, но больше задержка (latency). 4) `PutDatabaseRecord` - Вставляет записи в таблицу `stg.samplekafka2postgres`. @@ -105,6 +146,14 @@ - **Нет таблиц/процедуры**: после “чистого сброса” нужно снова выполнить `SampleKafka2Postgres.sql`. ## 7) Мини-упражнения для закрепления (5–10 минут) -1) В `Sample2Kafka.json` поменяйте расписание `GenerateRecord` (например, 2 sec → 5 sec) и проверьте скорость появления сообщений. -2) В `QueryRecord` измените фильтр `GOOD_DATE` (например, более жёсткое/мягкое условие) и проверьте, как меняется количество строк в Postgres. -3) Добавьте новое поле в схему `GenerateRecord` и протащите его до Postgres (после изменения не забудьте обновить SQL-схему/таблицы). + +1) **Изменить частоту генерации** + - В `Sample2Kafka.json` поменяйте расписание `GenerateRecord` (например, 2 sec → 5 sec) и проверьте скорость появления сообщений в Kafka UI. + +2) **Изменить фильтр** + - В `QueryRecord` измените фильтр `GOOD_DATE` (например, более жёсткое/мягкое условие) и проверьте, как меняется количество строк в Postgres. + +3) **Добавить новое поле (расширенное упражнение)** + - Добавьте новое поле в схему `GenerateRecord` и протащите его до Postgres. + - После изменения не забудьте обновить SQL-схему/таблицы. + - *Подсказка:* потребуются изменения в схеме JSON, настройках процессоров и SQL-скрипте. diff --git a/README.md b/README.md index 69c6034..76b6f1e 100644 --- a/README.md +++ b/README.md @@ -1,7 +1,35 @@ -# nifi-docker +> **Последнее обновление:** 2026-01-08 | **Версия NiFi:** 1.x.x Учебный стенд для знакомства с Apache NiFi и Kafka: NiFi + NiFi Registry + Postgres + Kafka + Kafka UI. +## Содержание + +- [Содержание](#содержание) +- [Быстрый старт](#быстрый-старт) + - [Что сохраняется между перезапусками](#что-сохраняется-между-перезапусками) +- [Адреса и доступы](#адреса-и-доступы) +- [Важно про адреса (внутри Docker и с хоста)](#важно-про-адреса-внутри-docker-и-с-хоста) +- [PostgreSQL](#postgresql) + - [Подключение через DBeaver (удобнее всего)](#подключение-через-dbeaver-удобнее-всего) + - [Через консоль (если нужно)](#через-консоль-если-нужно) +- [Kafka](#kafka) + - [Работа с Kafka](#работа-с-kafka) + - [Базовые команды CLI (опционально)](#базовые-команды-cli-опционально) +- [Примеры flow (шаблоны)](#примеры-flow-шаблоны) + - [Памятка: как импортировать Process Group / Flow в NiFi](#памятка-как-импортировать-process-group--flow-в-nifi) +- [Shared folder (общая папка)](#shared-folder-общая-папка) +- [Полезные команды](#полезные-команды) +- [Проверка работоспособности](#проверка-работоспособности) + - [Checklist](#checklist) + - [Проверка flow](#проверка-flow) +- [Если что-то не работает](#если-что-то-не-работает) + - [Проблемы с запуском контейнеров](#проблемы-с-запуском-контейнеров) + - [Проблемы с подключением в NiFi](#проблемы-с-подключением-в-nifi) + - [Проблемы с Kafka](#проблемы-с-kafka) + - [Проблемы с Postgres](#проблемы-с-postgres) + - [Проблемы с правами доступа](#проблемы-с-правами-доступа) + - [Дополнительная диагностика](#дополнительная-диагностика) + ## Быстрый старт ```sh docker compose up -d @@ -40,7 +68,8 @@ docker compose down -v Примечание: в этой конфигурации Kafka-сообщения/топики не сохраняются между `docker compose down` → `up` (чтобы не копить дисковое пространство). Между `stop` → `start` Kafka сохраняется. Пример подключения тома (volume) для Kafka есть в `docker-compose.yml`. ## Адреса и доступы -- NiFi: http://localhost:18443/nifi/ (логин `admin`, пароль `Password123456`) +- NiFi: http://localhost:18443/nifi/ + > 💡 **Примечание:** В текущей конфигурации аутентификация отключена для упрощения обучения. Логин и пароль не требуются. - NiFi docs: https://nifi.apache.org/documentation/ - Registry: http://localhost:18080/nifi-registry - Registry docs: https://nifi.apache.org/docs/nifi-registry-docs/ @@ -48,7 +77,8 @@ docker compose down -v - Kafka UI docs: https://docs.kafka-ui.provectus.io/ ## Важно про адреса (внутри Docker и с хоста) -Если вы настраиваете подключение *в NiFi*, то `localhost` почти всегда будет неправильным (NiFi живёт в контейнере). + +> ⚠️ **Критически важно:** Если вы настраиваете подключение *в NiFi*, то `localhost` почти всегда будет неправильным (NiFi живёт в контейнере). Используйте имена сервисов из `docker-compose.yml`: - Postgres (из NiFi): `jdbc:postgresql://postgres:5432/app` @@ -96,6 +126,40 @@ select * from ods.samplekafka2postgres order by id desc limit 10; - С локальной машины (например, для консольных утилит): `localhost:9092` - Из NiFi (внутри Docker): `kafka:29092` +### Работа с Kafka + +> 💡 **Рекомендация:** Для большинства задач используйте **Kafka UI** (http://localhost:8082/) — это удобный и наглядный веб-интерфейс для просмотра топиков, сообщений и consumer groups. + +**Kafka CLI** (командная строка) полезен для: +- Автоматизации и скриптов +- Продвинутых операций (например, изменение конфигурации топиков) +- Быстрой проверки без браузера + +#### Базовые команды CLI (опционально) + +Все команды выполняются через `docker compose exec kafka`: + +**Чтение сообщений:** +```sh +docker compose exec kafka kafka-console-consumer.sh \ + --bootstrap-server localhost:9092 \ + --topic Sample2Kafka \ + --from-beginning +``` + +**Запись сообщений:** +```sh +docker compose exec kafka kafka-console-producer.sh \ + --bootstrap-server localhost:9092 \ + --topic Sample2Kafka +``` +Пример JSON-сообщения: `{"dttm": 1704698992000, "txt": "Тестовое сообщение"}` + +**Просмотр списка топиков:** +```sh +docker compose exec kafka kafka-topics.sh --bootstrap-server localhost:9092 --list +``` + ## Примеры flow (шаблоны) Короткий гайд для менти по JSON-версии потоков: `MENTEE_GUIDE_JSON_FLOWS.md`. @@ -138,7 +202,99 @@ docker compose exec postgres bash Для доступа из NiFi к сервисам на локальной машине используйте `host.docker.internal` вместо `localhost`. +## Проверка работоспособности + +После запуска стенда проверьте, что все компоненты работают корректно: + +### Checklist + +- [ ] **NiFi UI открывается** по адресу http://localhost:18443/nifi/ (аутентификация не требуется) +- [ ] **Kafka UI показывает кластер** по адресу http://localhost:8082/ +- [ ] **В Kafka UI виден topic `Sample2Kafka`** (если flow уже запущен) +- [ ] **Postgres отвечает на подключение** через DBeaver или консоль: + ```sh + docker compose exec -it postgres bash -c "export PGPASSWORD=postgres; psql -U postgres -d app -c 'SELECT version();'" + ``` +- [ ] **Контейнеры запущены** (все в статусе `Up`): + ```sh + docker compose ps + ``` + +### Проверка flow + +Если вы импортировали и запустили примеры flow: +- [ ] В NiFi процессоры запущены (зелёный индикатор) +- [ ] В Kafka UI в topic `Sample2Kafka` появляются сообщения +- [ ] В Postgres таблица `ods.samplekafka2postgres` заполняется данными: + ```sh + docker compose exec -it postgres bash -c "export PGPASSWORD=postgres; psql -U postgres -d app -c 'SELECT COUNT(*) FROM ods.samplekafka2postgres;'" + ``` + ## Если что-то не работает -- NiFi может запускаться 1–3 минуты; смотрите `docker compose logs -f nifi`. -- Если процессор в NiFi не подключается к Kafka/Postgres, проверьте, что используете адреса “из NiFi” (см. раздел про Docker). -- Если порты заняты, измените проброс портов в `docker-compose.yml`. + +### Проблемы с запуском контейнеров + +**Симптом:** Контейнеры не запускаются или сразу падают. +- **Решение:** Проверьте логи: + ```sh + docker compose logs -f nifi + docker compose logs -f kafka + docker compose logs -f postgres + ``` +- NiFi может запускаться 1–3 минуты — это нормально. + +**Симптом:** Ошибка "port is already allocated". +- **Решение:** Порты заняты. Измените проброс портов в `docker-compose.yml` или остановите процессы, использующие эти порты. + +### Проблемы с подключением в NiFi + +**Симптом:** Процессор в NiFi не подключается к Kafka или Postgres. +- **Решение:** Проверьте, что используете адреса "из NiFi" (см. раздел [Важно про адреса](#важно-про-адреса-внутри-docker-и-с-хоста)): + - Postgres: `postgres:5432` (не `localhost:5437`) + - Kafka: `kafka:29092` (не `localhost:9092`) + +**Симптом:** Controller Services выключены. +- **Решение:** Внутри process group откройте `Controller Services` → выберите сервисы → `Enable`. + +### Проблемы с Kafka + +**Симптом:** Сообщения не появляются в Kafka UI. +- **Решение:** + - Проверьте, что процессор `PublishKafkaRecord_2_6` запущен и растёт счётчик "out" + - Убедитесь, что topic `Sample2Kafka` существует (создаётся автоматически при первой записи) + - Проверьте логи Kafka: `docker compose logs -f kafka` + +**Симптом:** Consumer в NiFi не читает сообщения. +- **Решение:** + - Проверьте, что используете правильный `bootstrap.servers`: `kafka:29092` + - Убедитесь, что consumer group не зафиксирован на другом offset (можно создать новый consumer group) + +### Проблемы с Postgres + +**Симптом:** Ошибка при вставке данных в Postgres. +- **Решение:** + - Проверьте, что таблицы существуют: выполните SQL-скрипт `SampleKafka2Postgres.sql` + - Убедитесь, что JDBC-драйвер подключён: `/opt/nifi/nifi-current/drivers/postgresql-42.7.4.jar` + - Проверьте параметры подключения: `jdbc:postgresql://postgres:5432/app`, пользователь `postgres`, пароль `postgres` + +**Симптом:** Таблица пустая, хотя flow работает. +- **Решение:** + - Проверьте, что сообщения проходят фильтр `QueryRecord` (relationship `GOOD_DATE`) + - Посмотрите логи процессора `PutDatabaseRecord` на наличие ошибок + - Убедитесь, что процедура `ods.load_samplekafka2postgres()` существует + +### Проблемы с правами доступа + +**Симптом:** Ошибки доступа к `shared-folder/`. +- **Решение:** Исправьте права: + ```sh + sudo chown -R $USER shared-folder + ``` + +### Дополнительная диагностика + +Если проблема не решена: +1. Посмотрите логи всех сервисов: `docker compose logs` +2. Перезапустите проблемный контейнер: `docker compose restart nifi` (или другой сервис) +3. Попробуйте полный сброс: `docker compose down -v && docker compose up -d` + > ⚠️ **Внимание:** Это удалит все данные из Postgres, NiFi и Registry.