Содержание
- Чем kafka-consumer-perf-test.sh отличается от producer-perf-test
- Параметры команды kafka-consumer-perf-test.sh
- Как читать вывод утилиты kafka-consumer-perf-test.sh
- Практическое использование kafka-producer-perf-test.sh
- Шаг 1. Подготовка данных и базовый тест
- Шаг 2. Параллельное чтение несколькими потоками
- Шаг 3. Промежуточная статистика с --show-detailed-stats
- Шаг 4. Тюнинг консьюмера через --consumer.config
- Связка producer-perf и consumer-perf: полный сценарий тестирования
- Альтернативы и сравнение подходов
- Что дальше
- Референсы
- Все уроки курса
В уроке 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 для инженеров данных» разбирает настройку консьюмеров, управление лагом и паттерны чтения на реальных кейсах.
Референсы
- Apache Kafka 4.2 Documentation. Performance Tools (официальная документация, 2025)
- Apache Kafka Consumer Configurations (fetch.min.bytes, fetch.max.wait.ms, max.partition.fetch.bytes)
- Confluent Developer. Optimizing Consumer Throughput (2025)
- Confluent Blog. Kafka Consumer Performance Tuning (2024-2025)
