Изучаем Apache Kafka с нуля. Урок 29. kafka-consumer-perf-test.sh

Изучаем Apache Kafka с нуля. Урок 29. kafka-consumer-perf-test.sh

 

В уроке 28 мы разобрали kafka-producer-perf-test.sh: как гнать синтетическую нагрузку в топик, читать перцентили латентности и подбирать параметры продюсера. Получили цифры на стороне записи. Но у любого потока данных есть вторая сторона — чтение. И там своя картина. Узкое место системы не всегда продюсер. Часто потребитель не успевает за входящим потоком, и накапливается лаг. Понять, с какой скоростью конкретный кластер отдаёт данные, и проверить, как ведут себя настройки консьюмера при разных сценариях — для этого есть kafka-consumer-perf-test.sh. Это симметричный инструмент к producer-perf-test: та же идеология, но для стороны чтения.

Для тех, кто занимается администрированием и хочет понять, как интерпретировать результаты обоих тестов и выстраивать методологию нагрузочного тестирования кластера, эти темы подробно разбираются в курсе «Администрирование кластера Kafka».

 

Чем kafka-consumer-perf-test.sh отличается от producer-perf-test

Принципиальное отличие одно: consumer-perf-test читает уже существующие данные из топика, а не генерирует их. Это значит, что перед запуском теста в топике должны быть сообщения — обычно их заливают через kafka-producer-perf-test.sh в предыдущем шаге.

Второе отличие — формат вывода. Если producer-perf-test пишет строки в реальном времени по мере отправки, то consumer-perf-test по умолчанию выдаёт одну итоговую строку в CSV-формате с именованными полями. Это удобнее для автоматизации: вывод проще парсить скриптом.

Ещё одна особенность: утилита по умолчанию читает с начала топика (—from-latest отключён). То есть каждый запуск теста начинает с самого первого сообщения, что делает тесты воспроизводимыми без дополнительных сбросов offset’ов.

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Параметры команды kafka-consumer-perf-test.sh

Три параметра обязательны — без них утилита не запустится:

Флаг Значение Описание
—bootstrap-server host:port Адрес брокера. Обязателен в KRaft-режиме
—topic имя топика Топик, из которого читаем
—messages целое число Сколько сообщений прочитать и завершить тест

Дополнительные параметры, которые используются в реальных сценариях:

Флаг Значение Описание
—group строка ID consumer group. По умолчанию генерируется случайный, каждый запуск — новая группа
—threads целое число Количество потоков-консьюмеров. По умолчанию 10
—fetch-size байты Размер fetch-запроса. По умолчанию 1 048 576 (1 МБ)
—from-latest флаг Читать не с начала топика, а с конца (только новые сообщения)
—reporting-interval мс Как часто выводить промежуточную статистику. По умолчанию 5000 (5 секунд)
—show-detailed-stats флаг Включает промежуточный вывод по каждому интервалу. Без него — только финальная строка
—consumer.config путь к файлу Файл .properties с параметрами консьюмера
—timeout мс Таймаут ожидания новых сообщений, после которого тест завершается. По умолчанию 10 000

Параметр —threads создаёт несколько независимых потоков-консьюмеров в рамках одной consumer group. Это грубая эмуляция параллельного чтения: каждый поток берёт свои партиции. Для корректного распределения число партиций топика должно быть не меньше числа потоков.

 

Как читать вывод утилиты kafka-consumer-perf-test.sh

По умолчанию (без —show-detailed-stats) утилита выводит заголовок и одну итоговую строку в CSV-формате:

# Пример вывода kafka-consumer-perf-test.sh
start.time, end.time, data.consumed.in.MB, MB.sec, data.consumed.in.nMsg, nMsg.sec, rebalance.time.ms, fetch.time.ms, fetch.MB.sec, fetch.nMsg.sec
2025-04-10 11:02:00:000, 2025-04-10 11:02:08:412, 976.5625, 115.96, 1000000, 118838.12, 312, 8100, 120.56, 123456.78

Разберём каждое поле:

  • data.consumed.in.MB — сколько мегабайт прочитано всего. Позволяет проверить, что тест действительно получил ожидаемый объём данных.
  • MB.sec — общая пропускная способность чтения в мегабайтах в секунду, включая время ребалансировки.
  • data.consumed.in.nMsg — сколько сообщений прочитано. Должно совпасть с —messages.
  • nMsg.sec — скорость в сообщениях в секунду, тоже включая rebalance.
  • rebalance.time.ms — время, которое ушло на ребалансировку consumer group до начала фактического чтения. Это overhead, который нужно учитывать: в production с большим числом партиций он может быть заметным.
  • fetch.time.ms — чистое время чтения без ребалансировки. Именно это значение показывает реальную скорость IO.
  • fetch.MB.sec и fetch.nMsg.sec — пропускная способность только за время фактического чтения. Это главные метрики для оценки производительности: они не замутнены временем на rebalance.

Когда сравниваете два запуска, смотрите на fetch.MB.sec, а не на MB.sec. Разница между ними — это цена ребалансировки, и при частых перезапусках группы она может быть существенной.

 

Практическое использование kafka-producer-perf-test.sh

Шаг 1. Подготовка данных и базовый тест

Для теста нужны данные в топике. Используем kafka-producer-perf-test.sh из предыдущего урока, чтобы заполнить топик:

# Шаг 1а. Создаём топик с 6 партициями для параллельного чтения
kafka-topics.sh \
  --bootstrap-server localhost:9092 \
  --create \
  --topic consumer-perf-test \
  --partitions 6 \
  --replication-factor 1

# Шаг 1б. Заливаем 1 000 000 записей по 1024 байта
kafka-producer-perf-test.sh \
  --topic consumer-perf-test \
  --num-records 1000000 \
  --record-size 1024 \
  --throughput -1 \
  --producer-props bootstrap.servers=localhost:9092 acks=1

Теперь базовый тест чтения: один поток (—threads 1), читаем 1 миллион сообщений:

# Базовый тест скорости чтения одним консьюмером
kafka-consumer-perf-test.sh \
  --bootstrap-server localhost:9092 \
  --topic consumer-perf-test \
  --messages 1000000 \
  --threads 1 \
  --group perf-group-single

Значение fetch.MB.sec в этом тесте — это верхняя граница скорости одного консьюмера при текущих настройках сети и дисков. Запомните эту цифру: она станет точкой отсчёта для всех следующих экспериментов.

 

Шаг 2. Параллельное чтение несколькими потоками

Один консьюмер упирается в одно TCP-соединение и один поток обработки. Чтобы использовать все партиции топика, нужно несколько консьюмеров. Флаг —threads запускает их в рамках одной группы:

# Тест параллельного чтения: 6 потоков на 6 партиций
# Каждый поток получит ровно по одной партиции
kafka-consumer-perf-test.sh \
  --bootstrap-server localhost:9092 \
  --topic consumer-perf-test \
  --messages 1000000 \
  --threads 6 \
  --group perf-group-parallel

При числе потоков, равном числу партиций, каждый поток работает с одной партицией без конкуренции. Это оптимальный вариант. Если потоков больше, чем партиций — лишние будут простаивать: Kafka не может назначить одну партицию двум консьюмерам одной группы одновременно.

Сравните fetch.MB.sec этого теста с однопоточным. Линейный рост (6 потоков = 6x скорости) встречается редко: обычно есть потолок на стороне сети или дисков брокера. Если прирост непропорционально мал, узкое место уже не консьюмер, а инфраструктура.

 

Шаг 3. Промежуточная статистика с —show-detailed-stats

Флаг —show-detailed-stats включает вывод по каждому интервалу, заданному в —reporting-interval. Это полезно, чтобы увидеть не среднее за весь тест, а динамику: стабильна ли скорость или она скачет:

# Детальная статистика каждые 2 секунды
kafka-consumer-perf-test.sh \
  --bootstrap-server localhost:9092 \
  --topic consumer-perf-test \
  --messages 1000000 \
  --threads 6 \
  --group perf-group-detailed \
  --show-detailed-stats \
  --reporting-interval 2000

Вывод будет содержать строку-заголовок и отдельную строку за каждый 2-секундный интервал. Если скорость ровная — кластер работает стабильно. Провалы в отдельных интервалах указывают на GC-паузы на брокере, перегрузку сети или неравномерное распределение данных по партициям.

 

Шаг 4. Тюнинг консьюмера через —consumer.config

Параметры консьюмера напрямую влияют на скорость чтения. Три самых важных для производительности:

  • fetch.min.bytes — минимальный объём данных, который брокер должен накопить перед ответом на fetch-запрос. По умолчанию 1 байт (отвечает немедленно). Увеличение до 65536-1048576 снижает число roundtrip-запросов и повышает throughput при потоках с небольшими сообщениями.
  • fetch.max.wait.ms — как долго брокер ждёт, чтобы набрать fetch.min.bytes. По умолчанию 500 мс. Работает в паре с fetch.min.bytes: брокер отвечает, когда либо накопил нужный объём, либо истёк таймаут.
  • max.partition.fetch.bytes — максимум данных, которые брокер возвращает за один fetch по одной партиции. По умолчанию 1 048 576 (1 МБ). При крупных сообщениях стоит увеличить.

Пример файла конфигурации для теста на высокий throughput:

# consumer-perf.properties
bootstrap.servers=localhost:9092
fetch.min.bytes=1048576
fetch.max.wait.ms=500
max.partition.fetch.bytes=10485760

Передаём файл через —consumer.config:

# Тест с настройками из файла конфигурации
kafka-consumer-perf-test.sh \
  --bootstrap-server localhost:9092 \
  --topic consumer-perf-test \
  --messages 1000000 \
  --threads 6 \
  --group perf-group-tuned \
  --consumer.config /tmp/consumer-perf.properties

Сравните fetch.MB.sec с результатами Шага 2 (те же данные и потоки, но дефолтные настройки). Рост throughput при увеличении fetch.min.bytes особенно заметен на топиках с мелкими сообщениями — там число fetch-запросов снижается значительно.

 

Связка producer-perf и consumer-perf: полный сценарий тестирования

На практике оба инструмента запускают в связке. Сначала заливают данные с известными параметрами, потом измеряют скорость чтения. Это позволяет ответить на конкретный вопрос: сколько консьюмеров нужно, чтобы успевать за продюсером с заданным throughput.


flowchart TD
    A["kafka-producer-perf-test.sh\n--num-records 1 000 000\n--record-size 1024\n--throughput -1"] --> B["Топик: consumer-perf-test\n6 партиций\n~976 МБ данных"]
    B --> C["kafka-consumer-perf-test.sh\n--threads 1\n--group single"]
    B --> D["kafka-consumer-perf-test.sh\n--threads 6\n--group parallel"]
    C --> E["fetch.MB.sec: ~116\nrebalance.time.ms: ~250"]
    D --> F["fetch.MB.sec: ~580\nrebalance.time.ms: ~350"]
    E --> G{"throughput\nпродюсера > fetch.MB.sec?"}
    F --> G
    G -- "Да" --> H["Увеличить threads\nили добавить партиций"]
    G -- "Нет" --> I["Консьюмер справляется\nс текущей нагрузкой"]

Логика проверки простая: если fetch.MB.sec консьюмера ниже, чем MB.sec продюсера из предыдущего теста — лаг будет расти. Нужно либо увеличивать параллелизм чтения, либо масштабировать сам кластер.

 

Альтернативы и сравнение подходов

Для тестирования скорости чтения есть и другие инструменты. kcat (урок 33) умеет читать из топика с замером времени через системные утилиты, но не даёт готовой статистики throughput. Для точных цифр придётся считать вручную через pv или аналог.

Java-консьюмер с ручным замером времени через System.currentTimeMillis() даёт максимальный контроль: можно измерить время от первого poll() до обработки последнего сообщения, включить свою логику десериализации. Но это полноценный код, а не CLI-инструмент. Для быстрой оценки состояния кластера kafka-consumer-perf-test.sh в разы удобнее — одна команда, CSV на выходе.

Ещё один вариант — kafka-verifiable-consumer.sh из урока 27. Он показывает каждое сообщение, но не считает throughput: для нагрузочного тестирования это не тот инструмент.

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Что дальше

В уроке 30 переходим к kafka-mirror-maker.sh — утилите для репликации данных между кластерами. Это уже не тестирование, а production-инструмент для катастрофоустойчивости и географического распределения данных.

Если тема производительности для вас не только про инструменты, но и про архитектуру — курс «Apache Kafka для инженеров данных» разбирает настройку консьюмеров, управление лагом и паттерны чтения на реальных кейсах.

Референсы

Все уроки курса

Тема URL
1 Установка Kafka с Zookeeper https://test.bigdataschool.ru/blog/news/lesson1-kafka-zookeeper-install/
2 Установка Kafka в режиме KRaft https://test.bigdataschool.ru/blog/news/lesson2-kafka-kraft-install/
3 Docker KRaft: однонодовый кластер https://test.bigdataschool.ru/blog/news/lesson3-kafka-docker-single/
4 Docker KRaft: 3-нодовый кластер https://test.bigdataschool.ru/blog/news/lesson4-kafka-docker-cluster/
5 Утилиты bin/: переменные окружения и основы https://test.bigdataschool.ru/blog/news/lesson5-kafka-bin-intro/
6 kafka-topics.sh: управление топиками https://test.bigdataschool.ru/blog/news/lesson6-kafka-topics/
7 kafka-console-producer.sh https://test.bigdataschool.ru/blog/news/lesson7-kafka-console-consumer/
8 kafka-console-consumer.sh https://test.bigdataschool.ru/blog/news/lesson8-kafka-console-consumer/
9 kafka-server-start.sh / kafka-server-stop.sh https://test.bigdataschool.ru/blog/news/lesson9-kafka-server-start-stop/
10 kafka-storage.sh https://test.bigdataschool.ru/blog/news/lesson10-kafka-storage/
11 kafka-cluster.sh https://test.bigdataschool.ru/blog/news/lesson11-kafka-cluster/
12 kafka-metadata-quorum.sh https://test.bigdataschool.ru/blog/news/lesson12-kafka-metadata-quorum/
13 kafka-metadata-shell.sh https://test.bigdataschool.ru/blog/news/lesson13-kafka-metadata-shell/
14 kafka-features.sh https://test.bigdataschool.ru/blog/news/lesson14-kafka-features/
15 kafka-configs.sh https://test.bigdataschool.ru/blog/news/lesson15-kafka-configs/
16 kafka-log-dirs.sh https://test.bigdataschool.ru/blog/news/lesson16-kafka-log-dirs/
17 kafka-dump-log.sh https://test.bigdataschool.ru/blog/news/lesson17-kafka-dump-log/
18 kafka-delete-records.sh https://test.bigdataschool.ru/blog/news/lesson18-kafka-delete-records/
19 kafka-consumer-groups.sh https://test.bigdataschool.ru/blog/news/lesson19-kafka-consumer-groups/
20 kafka-streams-application-reset.sh https://test.bigdataschool.ru/blog/news/lesson20-kafka-streams-reset/
21 kafka-leader-election.sh https://test.bigdataschool.ru/blog/news/lesson21-kafka-leader-election/
22 kafka-reassign-partitions.sh https://test.bigdataschool.ru/blog/news/lesson22-kafka-reassign-partitions/
23 kafka-replica-verification.sh https://test.bigdataschool.ru/blog/news/lesson23-kafka-replica-verification/
24 kafka-acls.sh https://test.bigdataschool.ru/blog/news/lesson24-kafka-acls/
25 kafka-broker-api-versions.sh https://test.bigdataschool.ru/blog/news/lesson25-kafka-broker-api-versions/
26 kafka-get-offsets.sh https://test.bigdataschool.ru/blog/news/lesson26-kafka-get-offsets/
27 kafka-verifiable-producer/consumer.sh https://test.bigdataschool.ru/blog/news/lesson27-kafka-verifiable/
28 kafka-producer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson28-kafka-producer-perf/
29 kafka-consumer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson29-kafka-consumer-perf/
30 kafka-mirror-maker.sh https://test.bigdataschool.ru/blog/news/lesson30-kafka-mirror-maker/
31 connect-standalone.sh https://test.bigdataschool.ru/blog/news/lesson31-kafka-connect-standalone/
32 connect-distributed.sh https://test.bigdataschool.ru/blog/news/lesson32-kafka-connect-distributed/
33 kcat. Альтернативный CLI https://test.bigdataschool.ru/blog/news/lesson33-kcat/