10 KiB
Гайд для менти: как работают примеры flow на JSON (NiFi → Kafka → Postgres)
Последнее обновление: 2026-01-08 | Версия NiFi: 1.x.x | ← Вернуться в README
Этот гайд про два flow definition файла:
nifi-templates/Sample2Kafka.json— генерирует записи и публикует их в Kafka topicSample2Kafkaв виде JSON.nifi-templates/SampleKafka2Postgres.json— читает JSON из Kafka topicSample2Kafkaи пишет данные в Postgres.
0) Что мы строим (общая картинка)
Цепочка очень простая:
- NiFi генерирует запись (как “строчку данных”).
- NiFi отправляет эту запись в Kafka (как сообщение в topic).
- Другой flow в NiFi читает эти сообщения из Kafka.
- NiFi фильтрует/готовит данные и вставляет их в Postgres.
Это учебный пример, который показывает базовые роли технологий:
- NiFi — визуальный конвейер (pipeline) данных.
- Kafka — "почтовый ящик"/шина сообщений между системами.
- Postgres — база данных, где данные хранятся и доступны SQL-запросами.
Визуализация потока данных
graph TD
subgraph Flow1["Flow 1: NiFi → Kafka"]
GR["GenerateRecord<br/>(генерация)"]
PK["PublishKafkaRecord_2_6<br/>(публикация в Kafka)"]
KT["Kafka Topic: Sample2Kafka"]
GR --> PK --> KT
end
subgraph Flow2["Flow 2: Kafka → Postgres"]
CK["ConsumeKafkaRecord_2_6<br/>(чтение из Kafka)"]
QR["QueryRecord<br/>(фильтрация по GOOD_DATE)"]
MC["MergeContent<br/>(micro-batch)"]
PDR["PutDatabaseRecord<br/>(вставка в stg)"]
ES["ExecuteSQL<br/>(вызов процедуры + очистка)"]
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 — шаг обработки (прочитать/преобразовать/записать).
- Controller Service — общая настройка, которую используют процессоры (например, “как читать JSON” или “как подключаться к базе”).
2) Какой JSON у нас в примере
Запись состоит из двух полей:
dttm— дата/время (timestamp).txt— строка (text).
Пример JSON-записи:
{
"dttm": "2026-01-08T11:49:52Z",
"txt": "Пример текста"
}
В первом flow запись создаётся процессором GenerateRecord по схеме (она задана в настройке процессора).
Важно: в этом учебном варианте мы не используем Avro и не используем Schema Registry. Сообщения в Kafka — это обычный JSON.
3) Перед стартом: что проверить
-
Стенд поднят:
docker compose up -d -
Открывается NiFi: http://localhost:18443/nifi/
💡 Примечание: В текущей конфигурации аутентификация отключена для упрощения обучения. Логин и пароль не требуются.
-
Открывается Kafka UI: http://localhost:8082/
-
Примените 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.
Шаги внутри flow
GenerateRecord- Создаёт запись по схеме (в настройке
schema-text). - Частота генерации задаётся расписанием процессора (например, раз в 2 секунды).
- Создаёт запись по схеме (в настройке
PublishKafkaRecord_2_6- Берёт записи и публикует в Kafka.
- Важные настройки:
bootstrap.servers:kafka:29092(имя сервиса внутри Docker).topic:Sample2Kafkarecord-reader:JsonTreeReader(как “прочитать запись” перед отправкой)record-writer:JsonRecordSetWriter(как сериализовать в JSON)
Как проверить, что flow работает
- В NiFi у
PublishKafkaRecord_2_6растёт счётчик “out”. - В Kafka UI в topic
Sample2Kafkaпоявляются сообщения.
Если сообщений нет:
- Проверьте, что процессоры запущены и нет очереди/ошибок на
failure. - Проверьте, что внутри NiFi указано
kafka:29092, а неlocalhost:9092.
5) Flow 2 — SampleKafka2Postgres.json (Kafka → Postgres)
Цель: читать сообщения из Sample2Kafka, отфильтровать записи и записать их в Postgres.
Шаги внутри flow
ConsumeKafkaRecord_2_6- Читает сообщения из Kafka topic
Sample2Kafka. - Важные настройки:
bootstrap.servers:kafka:29092topic:Sample2Kafkarecord-reader:JsonTreeReader(как разобрать JSON из Kafka)
- Читает сообщения из Kafka topic
QueryRecord- Выполняет простой “фильтр” по данным.
- В примере есть запрос
GOOD_DATE(это имя relationship), который отбирает записи по условию наdttm. - На выходе получаются записи, которые прошли фильтр.
MergeContent(micro-batch)- Собирает несколько FlowFile в "пачку" перед загрузкой в БД.
- Зачем: много маленьких вставок (каждое сообщение → отдельный insert) может перегружать БД; batching снижает накладные расходы.
- Важные настройки:
Minimum Number of Entries:100— минимальный размер пачкиMaximum Number of Entries:1000— максимальный размер пачкиMax Bin Age:30 seconds— максимальное ожидание, чтобы пачка всё равно отправилась даже если данных мало
- Компромисс: чем больше batch, тем меньше нагрузка на БД, но больше задержка (latency).
PutDatabaseRecord- Вставляет записи в таблицу
stg.samplekafka2postgres. - Использует подключение к Postgres через
DBCPConnectionPool.
- Вставляет записи в таблицу
ExecuteSQL- Вызывает хранимую процедуру
ods.load_samplekafka2postgres()и чистит staging. - Это шаг “переложить в витрину/ODS”, чтобы показать типичный паттерн.
- Вызывает хранимую процедуру
Как проверить результат в Postgres
В контейнере Postgres выполните:
docker compose exec -it postgres bash -c \"export PGPASSWORD=postgres; psql -U postgres -d app\"
И затем SQL:
select * from ods.samplekafka2postgres order by id desc limit 10;
Если таблица пустая:
- Убедитесь, что в Kafka действительно есть сообщения (проверьте Kafka UI).
- Посмотрите, нет ли ошибок на
parse.failure(у Kafka consumer) илиfailure(у DB).
6) Типовые проблемы (и куда смотреть)
- В NiFi указали
localhostвместо имени сервиса: из контейнераlocalhost— это сам контейнер, не Kafka/Postgres. - Controller Services выключены: внутри process group откройте
Controller Servicesи включите нужные сервисы. - Нет таблиц/процедуры: после “чистого сброса” нужно снова выполнить
SampleKafka2Postgres.sql.
7) Мини-упражнения для закрепления (5–10 минут)
-
Изменить частоту генерации
- В
Sample2Kafka.jsonпоменяйте расписаниеGenerateRecord(например, 2 sec → 5 sec) и проверьте скорость появления сообщений в Kafka UI.
- В
-
Изменить фильтр
- В
QueryRecordизмените фильтрGOOD_DATE(например, более жёсткое/мягкое условие) и проверьте, как меняется количество строк в Postgres.
- В
-
Добавить новое поле (расширенное упражнение)
- Добавьте новое поле в схему
GenerateRecordи протащите его до Postgres. - После изменения не забудьте обновить SQL-схему/таблицы.
- Подсказка: потребуются изменения в схеме JSON, настройках процессоров и SQL-скрипте.
- Добавьте новое поле в схему