
Введение
«Миллионы сообщений в секунду» - величина, лишённая смысла без профиля нагрузки: кластер, прокачивающий 2 млн сообщений по 200 байт, и кластер на 50 тыс. сообщений по 100 КБ упираются в совершенно разные ресурсы. Прежде чем крутить параметры, зафиксируйте средний и максимальный размер сообщения, целевой объём в МБ/с, требуемый p99 end-to-end и допустимость потерь. Без этих четырёх чисел настройка превращается в угадывание.
Вторая частая причина неудачного тюнинга - перенос советов из статей 2016–2020 годов: с тех пор сменились дефолты acks, session.timeout.ms, replica.lag.time.max.ms и режим работы кластера целиком. Все значения ниже приводятся для Kafka 3.x-4.x, и там, где дефолт менялся, это отмечено по месту.
За рамками сознательно оставлены Kafka Streams, Kafka Connect, Schema Registry и мультирегиональные развёртывания.
Настройки брокера, влияющие на пропускную способность и задержку

Потоки и очереди запросов
Kafka-брокер держит два пула: num.network.threads (по умолчанию 3) - сетевые потоки: читают и пишут в сокеты и складывают запросы в очередь. num.io.threads (по умолчанию 8) - обработчики: пишут в лог и обслуживают чтения.
Сетевые потоки почти не выполняют полезной работы, поэтому масштабировать их по числу ядер бессмысленно - на 64-ядерной машине три десятка сетевых потоков будут простаивать. Для num.io.threads ориентир другой: нижняя граница - число дисков под данные, верхняя - число ядер.
Насыщение видно по JMX: NetworkProcessorAvgIdlePercent (kafka.network:type=SocketServer) и RequestHandlerAvgIdlePercent (kafka.server:type=KafkaRequestHandlerPool), доля простоя от 0 до 1. Приближается к нулю - потоки загружены полностью, и добавление даст эффект, если есть свободный CPU и диски не упёрлись. Держится выше 0.3 - добавлять бесполезно, узкое место в другом.
Оба параметра динамические (KIP-226): меняются через kafka-configs.sh --alter --entity-type brokers без перезапуска брокера. Эксперимент стоит секунды, а не окна обслуживания.
queued.max.requests (по умолчанию 500) ограничивает очередь между пулами. Когда она заполнена, сетевые потоки перестают читать из сокетов - каналы «замьючиваются», и это работает как встроенный backpressure для продюсеров и реплик. Увеличение сглаживает кратковременные всплески, но растит задержку и потребление памяти, следить по RequestQueueSize.
Ограничения размера сообщений
socket.request.max.bytes - максимальный размер одного запроса к брокеру, по умолчанию 104857600 (100 МиБ). Крупные батчи или большие отдельные сообщения упрутся в него и будут отвергнуты, поднимать его стоит с оглядкой на память брокера. message.max.bytes на брокере и max.message.bytes на топике задают максимальный размер записи, по умолчанию около 1 МиБ.
Ключевой нюанс, о котором часто забывают: лимит применяется к сжатому record batch целиком, а не к отдельному сообщению. Именно поэтому RecordTooLargeException чаще всего прилетает не из-за одного большого сообщения, а из-за подросшего батча.
Чтобы принимать более крупные сообщения, вместе с message.max.bytes согласованно поднимают socket.request.max.bytes, replica.fetch.max.bytes (объём, который реплика-последователь запрашивает у лидера, по умолчанию 1048576), клиентский max.request.size на продюсере (1048576) и max.partition.fetch.bytes на консьюмере.
Здесь же - устаревший совет, который до сих пор кочует по статьям: «если replica.fetch.max.bytes меньше message.max.bytes, репликация встанет». Это верно только для версий до 0.10.1. Начиная с KIP-74 брокер всегда возвращает как минимум один record batch, даже если тот превышает лимит выборки, так что ни репликация, ни потребление на большом сообщении не застревают. Согласовывать значения по-прежнему стоит, но ради предсказуемости и корректного расчёта памяти, а не из страха дедлока.
Память под выборки считается как произведение размера выборки на число партиций и число fetcher-потоков, поэтому щедрые лимиты на кластере с тысячами партиций съедают heap незаметно. И стратегически: передача крупных объектов через Kafka - антипаттерн. Правильнее claim-check, когда в топик кладётся ссылка на объект во внешнем хранилище, а не сам объект.
Лог-сегменты и flush
log.segment.bytes задаёт размер сегмента, после которого брокер переворачивает лог, по умолчанию 1 Гб. Сегмент переключается не только по размеру, но и по времени - за это отвечают log.roll.ms и log.roll.hours (по умолчанию 7 дней), что существенно для топиков с редкой записью: там сегмент может не добраться до гигабайта неделями.
Крупные сегменты снижают накладные расходы на управление файлами, но удлиняют восстановление после нештатного завершения - при старте брокер проверяет последний сегмент каждой партиции. Ускорить эту фазу помогает num.recovery.threads.per.data.dir: по умолчанию он равен 1, и на многодисковых узлах его имеет смысл поднять до числа дисков. Мелкие сегменты, наоборот, восстанавливаются быстрее и быстрее удаляются по ретеншну ценой более частого открытия и закрытия файлов.
Отдельно про сайзинг диска: ретеншн (log.retention.ms, log.retention.bytes) удаляет только целые неактивные сегменты. При сегменте в 1 Гб фактический объём на диске всегда заметно больше номинально настроенного, и запас надо закладывать с этим расчётом.
Kafka не выполняет fsync при каждой записи - данные буферизуются в page cache и сбрасываются фоновыми механизмами ОС. log.flush.interval.messages (по умолчанию Long.MAX_VALUE) и log.flush.interval.ms (не задан) позволяют форсировать flush, но фактически отключены, и документация прямо рекомендует их не трогать, полагаясь на репликацию. Частый принудительный fsync заметно бьёт и по задержке, и по пропускной способности. В продакшене эти параметры оставляют по умолчанию.
Сетевые буферы
socket.send.buffer.bytes и socket.receive.buffer.bytes задают размер TCP-буферов на брокере, по умолчанию 102400 байт (100 КиБ) - этого достаточно для локальной сети. На каналах с большим RTT буфер увеличивают, ориентируясь на произведение пропускной способности на задержку (bandwidth-delay product), для репликации через WAN есть отдельный replica.socket.receive.buffer.bytes.
Критичная деталь, без которой совет не сработает: увеличение буфера в конфиге Kafka бесполезно, если не подняты системные лимиты net.core.wmem_max и net.core.rmem_max - ядро молча обрежет запрошенный размер. Значение -1 означает «использовать дефолт ОС», и в современных ядрах с TCP autotuning это часто лучший выбор, чем ручной подбор.
Здесь же уместно вспомнить про zero-copy (sendfile), на котором держится производительность Kafka: данные уходят из page cache в сокет минуя пользовательское пространство. С включённым SSL/TLS zero-copy перестаёт работать - каждый байт проходит через JVM для шифрования, и нагрузка на CPU заметно растёт. Это закладывают в сайзинг заранее, а не обнаруживают после включения TLS в проде.
Потоки репликации
num.replica.fetchers задаёт число потоков, которыми брокер-последователь вытягивает данные с каждого лидера, и по умолчанию равен 1. На кластерах с большим числом партиций именно он чаще всего оказывается узким местом репликации: один поток на источник не успевает обслуживать сотни партиций, а следствием становятся ненулевые UnderReplicatedPartitions и растущий лаг реплик - при том что CPU, сеть и диски выглядят недогруженными.
Разумный диапазон для нагруженных кластеров - от 2 до 8 потоков, с оглядкой на число ядер. Сопутствующие replica.fetch.min.bytes и replica.fetch.wait.max.ms работают для репликации так же, как fetch.min.bytes и fetch.max.wait.ms для консьюмера, позволяя укрупнять выборки.
Настраивать репликацию нужно вместе с числом партиций: рост партиций без роста fetcher-потоков предсказуемо ухудшает ISR-стабильность.
Квоты и защита кластера
В многопользовательском кластере квоты - единственный встроенный механизм, защищающий брокеры от клиента, который сошёл с ума. Задаются через kafka-configs.sh --alter --entity-type clients (или users) и включают producer_byte_rate и consumer_byte_rate (полоса в байт/с), request_percentage (доля времени потоков брокера, которую клиент вправе занять) и controller_mutation_rate (защита от лавины создания и удаления топиков и партиций).
Kafka не отбрасывает запросы при превышении квоты, а задерживает ответ, и эта задержка видна в метрике ThrottleTimeMs. Её надо смотреть первой, когда клиент жалуется на необъяснимо выросшую latency.
Сюда же относится ограничение числа подключений: max.connections, max.connections.per.ip и connections.max.idle.ms (по умолчанию 10 минут). При тысячах клиентов лимиты соединений и файловых дескрипторов упираются раньше, чем пулы потоков.
Конфигурация продюсера для эффективной записи

Подтверждения (acks)
acks определяет, сколько реплик должны подтвердить приём, прежде чем лидер вернёт продюсеру ответ. acks=0 - не ждать вообще, acks=1 - только лидера, acks=all (или -1) - все реплики из ISR.
Начиная с Kafka 3.0 значение по умолчанию - acks=all, а не acks=1, как было раньше: вместе с KIP-679 по умолчанию включилась идемпотентность продюсера (enable.idempotence=true), а она требует acks=all. Многие статьи и конспекты до сих пор указывают acks=1 как дефолт - это устаревшая информация.
При acks=1 лидер записывает сообщение в свой локальный лог и подтверждает сразу, не дожидаясь репликации. Задержка подтверждения меньше, но если лидер выйдет из строя до репликации, сообщение потеряно. При этом снижение acks не ускоряет доставку данных читателям: момент, с которого запись становится видимой консьюмеру, определяется продвижением high watermark по всем репликам ISR и от acks не зависит.
Практический вывод: не понижайте acks «на всякий случай». На современном железе разница между acks=all и acks=1 по пропускной способности обычно измеряется единицами процентов, а не разами - сначала измерьте её на своём профиле нагрузки. Если нужна максимальная надёжность, оставляйте acks=all в связке с min.insync.replicas, acks=1 и тем более acks=0 допустимы только там, где потеря части данных приемлема по бизнес-требованиям.
Батчинг: batch.size и linger.ms
Продюсер отправляет сообщения пачками. batch.size (по умолчанию 16384, то есть 16 Кб) задаёт целевой размер батча в байтах, linger.ms (по умолчанию 0) - максимальную задержку перед отправкой, чтобы набрать больше сообщений. В Kafka 4.0 дефолт linger.ms пересматривался (KIP-1030), сверьтесь с документацией своей версии.
Принципиальная деталь, которую обычно упускают: batch.size - это лимит на партицию, а не на запрос целиком. Продюсер накапливает отдельный батч для каждой партиции назначения и упаковывает готовые батчи в один запрос.
Отсюда два следствия. Пиковое потребление памяти продюсером - порядка «число активных партиций × batch.size», и оно должно помещаться в buffer.memory. И суммарный размер запроса ограничен max.request.size (по умолчанию 1048576 байт), поэтому поднять batch.size до 200 КБ и не тронуть max.request.size - типовая ошибка первого тюнинга.
Для высокой пропускной способности значения увеличивают: batch.size до 100-200 тыс. байт, linger.ms до 5-50 мс. Больший батч даёт больше throughput ценой добавленной задержки для первых сообщений в пакете.
Проверять результат надо не на глаз, а по клиентским метрикам batch-size-avg и record-queue-time-avg. Если средний размер батча заметно меньше batch.size, ограничителем является не размер, а темп поступления данных, и увеличивать batch.size дальше бессмысленно.
Сжатие
compression.type задаёт алгоритм: gzip, snappy, lz4, zstd или none (по умолчанию none). lz4 даёт наилучший баланс скорости, zstd - заметно лучшую степень сжатия при сопоставимых затратах CPU, в свежих версиях уровень регулируется параметром compression.zstd.level.
Обязательное условие, без которого сжатие превращается в антиоптимизацию: на брокере и топике compression.type должен быть равен producer (это дефолт). Если там выставлен конкретный кодек, отличный от того, которым пришли данные, брокер будет распаковывать и пережимать каждый батч - самая дорогая операция, которую можно ему навязать, и она надёжно съедает весь выигрыш.
Степень сжатия напрямую зависит от размера батча, потому что сжимается батч целиком. Поэтому compression.type нельзя тюнить в отрыве от batch.size и linger.ms: при linger.ms=0 и мелких батчах выигрыш будет минимальным. Контролировать эффект удобно по метрике compression-rate-avg.
Для сверхнизких задержек или очень мелких сообщений выигрыш может не окупить затрат CPU, но в большинстве высоконагруженных сценариев компрессию продюсера включают.
Буфер продюсера
buffer.memory (по умолчанию 33554432 байт, 32 Мб) определяет объём памяти под очередь отправки. При заполнении буфера вызовы send() блокируются, а по истечении max.block.ms (по умолчанию 60000 мс) бросается исключение.
Ориентир для расчёта конкретный: буфер должен быть не меньше «число активных партиций × batch.size» плюс запас в 1.5-2 раза на неотправленные запросы и накладные расходы. При 500 партициях и batch.size=100000 минимально осмысленное значение - уже около 50 МБ, то есть дефолта не хватает с большим отрывом.
Попали ли вы в целевой размер, показывают метрики buffer-available-bytes (не должна регулярно уходить в ноль) и waiting-threads (в норме нулевая). Именно они, а не догадки, говорят, упирается ли приложение в буфер.
Запросы в полёте
max.in.flight.requests.per.connection (по умолчанию 5) управляет тем, сколько запросов продюсер держит в полёте по одному соединению, не дожидаясь подтверждений. Большее значение лучше загружает сеть, особенно на каналах с заметным RTT.
Здесь необходимо развеять устойчивое заблуждение. Часто пишут, что идемпотентный или транзакционный продюсер вынужден работать с max.in.flight=1 - это неверно начиная с Kafka 1.0.0 (KAFKA-5494). Идемпотентный продюсер сохраняет порядок записей при значениях вплоть до 5 включительно: брокер отслеживает порядковые номера и сам отбрасывает либо переупорядочивает повторы.
Реальное ограничение - «не больше 5», и оно жёсткое. При включённой идемпотентности (а она включена по умолчанию с версии 3.0) значение больше 5 приводит к тому, что продюсер не стартует и падает с ConfigException. Поэтому распространённый совет «поднимите до 5-10 ради throughput» в современных версиях даёт не ускорение, а отказ приложения при инициализации.
Ставить max.in.flight=1 ради строгого порядка тоже больше не нужно - это легаси-рецепт для неидемпотентного продюсера, который сегодня лишь режет пропускную способность без выигрыша.
И стоит корректно понимать, что даёт идемпотентность: отсутствие дублей от повторных отправок в пределах сессии продюсера и одной партиции. Это не сквозная семантика exactly-once для всего конвейера - для неё нужны транзакции и соответствующая настройка консьюмера.
Таймауты и повторы
Пропускная способность в стабильном режиме - только половина задачи, вторая половина в том, как продюсер ведёт себя, когда брокер отвечает медленно или отваливается.
Верхняя граница жизни записи - delivery.timeout.ms (по умолчанию 120000 мс): суммарный бюджет от вызова send() до финального успеха или окончательной ошибки, включая ожидание в буфере, все повторы и сетевые задержки. retries в современных версиях по умолчанию равен Integer.MAX_VALUE и самостоятельного смысла почти не имеет - повторы обрываются именно по delivery.timeout.ms, поэтому настраивать нужно его, а не число попыток.
request.timeout.ms (по умолчанию 30000 мс) ограничивает ожидание ответа на один запрос, retry.backoff.ms (по умолчанию 100 мс) задаёт паузу между попытками. Соотношение должно быть осмысленным: delivery.timeout.ms обязан быть не меньше суммы linger.ms и request.timeout.ms, иначе конфигурация будет отвергнута.
Агрессивно короткий delivery.timeout.ms превращает кратковременную деградацию одного брокера в поток ошибок в приложении, чрезмерно длинный - в незаметное разрастание очереди и всплеск latency. Отслеживайте record-error-rate и record-retry-rate, чтобы видеть, живёт ли продюсер на повторах.
Партиционирование и наполняемость батчей
Реальный размер батча определяется не только batch.size и linger.ms, но и тем, как записи раскладываются по партициям. Для сообщений без ключа старые версии клиента раскидывали записи по кругу, из-за чего при большом числе партиций батчи не успевали наполняться и весь тюнинг батчинга сводился на нет.
Начиная с Kafka 2.4 (KIP-480) работает «липкий» партиционер: продюсер придерживается одной партиции, пока не наберёт полный батч, и только затем переключается. Это то самое изменение, которое зачастую даёт больший прирост throughput, чем ручной подбор batch.size. В версиях 3.3 и новее partitioner.class рекомендуется оставлять незаданным, чтобы использовалась встроенная реализация с равномерным липким распределением.
Для сообщений с ключом партиция вычисляется по хешу ключа, что даёт гарантию порядка в пределах ключа - и именно поэтому любое изменение числа партиций топика переносит ключи на другие партиции и эту гарантию ломает. Число партиций планируют заранее.
Конфигурация консьюмера для быстрой обработки сообщений

Консьюмер держит параллельные Fetch-запросы ко всем брокерам, на которых лежат его партиции, накапливает ответы в клиентском буфере и отдаёт их приложению порциями по max.poll.records. Границу «сколько приходит по сети» и границу «сколько отдаётся приложению» задают разные параметры, и путать их дорого.
Размер и задержка выборки
fetch.min.bytes задаёт минимальный объём данных, который брокер соберёт, прежде чем ответить на запрос выборки. По умолчанию он равен 1 байту - то есть брокер отвечает сразу, даже если данных почти нет. Увеличение до десятков или сотен килобайт заставляет накапливать пакет и снижает частоту запросов вместе с накладными расходами на сообщение, и у клиента, и у брокера.
fetch.max.wait.ms ограничивает ожидание сверху и по умолчанию равен 500 мс: если объём не собрался, брокер вернёт всё, что есть. Менять здесь имеет смысл прежде всего fetch.min.bytes - fetch.max.wait.ms уже равен 500, и совет «увеличьте до 500» ничего не меняет.
Плата за укрупнение выборок - хвостовые задержки: при неравномерном потоке p99 доставки поднимется практически до значения fetch.max.wait.ms, поэтому для latency-критичных потоков его, наоборот, снижают. Там, где сообщения идут плотно, увеличение fetch.min.bytes на задержку почти не влияет: данные для ответа есть всегда.
Максимальный объём выборки
max.partition.fetch.bytes ограничивает объём данных из одной партиции за запрос (по умолчанию 1048576, 1 Мб), fetch.max.bytes - объём всего ответа брокера (по умолчанию 52428800, около 50 МиБ). Вместе они не дают одной «горячей» партиции занять весь канал.
Как и в случае с репликацией, устарел совет «обязательно поднимите max.partition.fetch.bytes до размера максимального сообщения, иначе потребление встанет». После KIP-74 брокер возвращает как минимум один record batch, даже если тот больше лимита, поэтому консьюмер не застревает на большом сообщении. Увеличивать оба параметра стоит по другой причине - чтобы за один запрос получать больше данных и делать меньше сетевых раунд-трипов.
Оценивать память при этом надо аккуратно. fetch.max.bytes ограничивает объём ответа одного брокера, а консьюмер держит параллельные запросы ко всем брокерам, на которых лежат его партиции. Реалистичная верхняя оценка пикового буфера - «число брокеров с назначенными партициями × fetch.max.bytes», а не «число партиций × max.partition.fetch.bytes», как иногда пишут: вторая формула систематически занижает потребление. На каналах с высокой задержкой обратите внимание ещё и на receive.buffer.bytes консьюмера.
Автокоммит и смещения
enable.auto.commit (по умолчанию true) и auto.commit.interval.ms (по умолчанию 5000 мс) определяют, будет ли консьюмер сам фиксировать смещения и как часто.
Важна механика, а не сам факт: в классическом консьюмере автокоммит выполняется не отдельным фоновым потоком, а внутри вызова poll(). При очередном обращении к poll() клиент проверяет, истёк ли интервал, и фиксирует смещения предыдущей выборки. Отсюда неочевидное следствие: пока поток занят долгой обработкой и не вызывает poll(), коммита не происходит, сколько бы времени ни прошло. Именно из представления об автокоммите как о независимом фоновом процессе рождается большинство сюрпризов с потерями и дублями. В новом асинхронном консьюмере, появившемся вместе с протоколом KIP-848, часть работы действительно вынесена в фоновый поток - ещё одна причина фиксировать версию.
При сбоях автокоммит даёт неопределённость в обе стороны. Консьюмер погиб после обработки, но до коммита - сообщения считаются непрочитанными, и другой экземпляр группы обработает их повторно. Коммит прошёл до завершения обработки, а потребитель упал - часть сообщений помечена прочитанной и потеряна. Для критичных систем автокоммит отключают (enable.auto.commit=false) и коммитят сами после обработки батча. Слишком частый коммит нагружает внутренний топик __consumer_offsets, поэтому интервалы порядка нескольких секунд обычно оптимальны.
При ручном управлении есть commitSync() и commitAsync(): синхронный блокирует поток до подтверждения, но гарантирует сохранение смещений, асинхронный не задерживает обработку, но при сбое может не успеть подтвердить последние. Чаще всего применяют комбинацию - регулярные commitAsync() в цикле обработки и финальный commitSync() перед завершением или в обработчике ребаланса (ConsumerRebalanceListener.onPartitionsRevoked). Ни один из вариантов сам по себе не даёт exactly-once: если обработка имеет побочные эффекты, приложение должно быть идемпотентным либо использовать транзакции.
poll и таймаут сессии
Консьюмер должен регулярно вызывать poll() - и чтобы получать данные, и чтобы подтверждать свою жизнеспособность. Если poll() не вызывается чаще, чем max.poll.interval.ms (по умолчанию 300000 мс, 5 минут), консьюмер считает себя зависшим и сам покидает группу, отправляя LeaveGroup и инициируя ребаланс. Это делает клиент, а не брокер, поэтому искать причину в логах брокера бесполезно.
Когда обработка долгая, есть два выхода: поднять max.poll.interval.ms или уменьшить порцию через max.poll.records (по умолчанию 500). Здесь нужно чётко разделять две вещи, которые часто путают. fetch.* управляет тем, сколько данных приходит по сети и лежит в буфере клиента, max.poll.records лишь нарезает уже полученный буфер на порции для приложения. Снижение max.poll.records сокращает время одной итерации и риск просрочить max.poll.interval.ms, но не уменьшает ни сетевой трафик, ни потребление памяти - лечить им OOM бесполезно, для этого есть fetch.max.bytes и max.partition.fetch.bytes.
Обнаружение отказа настраивается отдельно. session.timeout.ms определяет, через сколько без heartbeat координатор группы признает потребителя упавшим. Начиная с Kafka 3.0 значение по умолчанию - 45000 мс (KIP-735), и оно было специально увеличено с прежних 10 секунд, чтобы паузы GC и сетевой джиттер не вызывали ложных ребалансов. Устаревшая рекомендация «поднимите session.timeout примерно до 30 секунд» сегодня означает уменьшение дефолта вдвое, то есть действие, обратное задуманному.
Heartbeat отправляется отдельным потоком с периодом heartbeat.interval.ms (по умолчанию 3000 мс), и правило «не больше трети session.timeout.ms» относится к настройке этого параметра, а не описывает поведение клиента. Верхняя граница задаётся брокерским group.max.session.timeout.ms. Оптимум - минимальное значение, при котором в нормальном режиме потребителей из группы не выкидывает.
Параллелизм потребления
Потребители с одним group.id делят партиции топика между собой, причём каждая партиция обслуживается ровно одним потребителем. Максимальный параллелизм ограничен числом партиций: экземпляры сверх этого числа данных не получают и работают как горячий резерв - не простаивают бессмысленно, а мгновенно принимают партиции на себя при сбое активного, что сокращает время восстановления.
При планировании учитывайте два ограничения. Число партиций топика нельзя уменьшить, а его увеличение перераспределяет ключи и ломает гарантию порядка в пределах ключа для уже записанных данных. Поэтому запас закладывают заранее.
Если обработка сообщения тяжёлая и упирается не в Kafka, а в саму бизнес-логику, наращивать партиции ради параллелизма не всегда правильно. Альтернатива - развязать потребление и обработку: вычитывать данные одним консьюмером и распределять работу по пулу воркеров, в том числе готовыми решениями вроде Parallel Consumer, с аккуратным управлением коммитами.
Ребалансировка
Для высоконагруженных групп настройки ребалансировки влияют на доступность сильнее, чем любые размеры выборок.
По умолчанию долгое время использовалась «жадная» (eager) стратегия: на время ребаланса все потребители отдают все свои партиции, и группа полностью останавливается - на большой группе это секунды простоя при каждом изменении состава. Начиная с Kafka 2.4 (KIP-429) доступна кооперативная стратегия: с partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor переназначаются только те партиции, которые реально меняют владельца, а остальные продолжают обрабатываться.
Второй по важности инструмент - статическое членство (KIP-345). Если задать каждому экземпляру уникальный и стабильный group.instance.id, плановый перезапуск в пределах session.timeout.ms вообще не вызовет ребаланса: координатор дождётся возвращения того же участника. Для развёртываний в Kubernetes, где перезапуски происходят постоянно, это устраняет целый класс «штормов ребаланса».
В Kafka 4.0 доступен новый протокол ребалансировки на стороне брокера (KIP-848), включаемый параметром group.protocol=consumer: он убирает stop-the-world-фазу и переносит вычисление назначений на координатора.
Если вы наблюдаете растущий лаг без роста нагрузки, начните диагностику с метрики rebalance-rate-per-hour, а не с размеров выборок.
Изоляция чтения и стартовая позиция
isolation.level по умолчанию равен read_uncommitted, и в этом режиме консьюмер видит записи незавершённых и даже прерванных транзакций. Если в системе есть транзакционные продюсеры, потребителя обязательно переводят в read_committed - иначе вся выстроенная на стороне записи семантика exactly-once не имеет смысла, читатель всё равно получит данные прерванных транзакций. Плата за это - дополнительная задержка: консьюмер не может читать дальше LSO (last stable offset), то есть дальше первой незавершённой транзакции.
auto.offset.reset (по умолчанию latest) определяет поведение при отсутствии сохранённого смещения или при его устаревании. latest означает «начать с конца», из-за чего новая группа молча пропустит всё уже накопленное, earliest - «прочитать всё с начала», что на большом топике даёт лавину чтения. Значение none заставляет клиента упасть с ошибкой, и для критичных конвейеров это нередко предпочтительнее молчаливого пропуска данных.
Топология кластера: партиции, репликация и балансировка нагрузки

Партиции топика раскладываются по брокерам, у каждой один лидер и RF-1 реплик-последователей. Лидер обслуживает запись и по умолчанию чтение, фолловеры вытягивают данные собственными fetch-запросами. При размещении в нескольких зонах реплики одной партиции разводят по разным зонам.
Количество партиций
Модель, которую часто описывают неверно: у брокера нет отдельного потока на каждую партицию. Запросы обслуживаются общим пулом num.io.threads, репликация - потоками num.replica.fetchers. Партиция является единицей параллелизма для клиентов и единицей распределения данных по узлам, но не единицей планирования потоков на брокере.
Отсюда и форма кривой: пропускная способность растёт с числом партиций примерно пропорционально лишь на начальном участке, дальше выходит на плато, определяемое дисками, сетью и CPU кластера, а после начинает падать из-за накладных расходов. Каждая партиция стоит ресурсов - файловые дескрипторы на сегменты и индексы, память под буферы, отдельные записи в метаданных, дополнительная работа fetcher-потоков.
Распространённое «нельзя больше нескольких сотен партиций на брокер» - сильно устаревшая цифра. В эпоху ZooKeeper практическим ориентиром были порядка 4000 партиций на брокер и около 200 000 на кластер, причём ограничение диктовалось не данными, а временем восстановления метаданных и перевыборов контроллера. В KRaft-режиме - единственном начиная с Kafka 4.0 - это ограничение снято: метаданные лежат в реплицируемом логе, и публично демонстрировались кластеры с миллионами партиций.
«Сколько влезет» при этом плохая стратегия. Исходите из требуемого параллелизма: нужно N параллельных обработчиков - нужно не меньше N партиций, плюс разумный запас на рост. Запас важен потому, что число партиций нельзя уменьшить, а его увеличение перераспределяет ключи и ломает гарантию порядка по ключу для уже записанных данных.
И не забывайте масштабировать вместе с партициями num.replica.fetchers: рост числа партиций без роста потоков репликации - самая частая причина деградации ISR.
Фактор репликации
Стандарт для продакшена - RF=3. Цена конкретная: каждый записываемый байт передаётся на RF-1 дополнительных брокеров, то есть при RF=3 суммарный объём записи на диски кластера и внутрикластерный сетевой трафик утраиваются относительно того, что поступает от продюсеров. Это снижает максимальный совокупный throughput, особенно при acks=all.
Снижать RF ради скорости тем не менее не стоит: RF=1 лишает систему отказоустойчивости, RF=2 несёт риск потери данных при падении одного узла во время планового обновления второго. Минимум из трёх брокеров и RF=3 - устойчивый компромисс. Для второстепенных топиков в кластерах с большими объёмами и невысокими требованиями к сохранности RF=2 допустим.
Отдельно: брокерский default.replication.factor по умолчанию равен 1, поэтому на новом кластере его нужно явно выставить в 3 - иначе автоматически созданные топики окажутся без реплик.
Если кластер размазан по нескольким зонам доступности, реплики раскладывают по зонам через rack-awareness.
Распределение лидерства
У каждой партиции один лидер, обслуживающий запросы клиентов, и несколько реплик-последователей. Фолловеры вытягивают данные с лидера собственными fetch-запросами, но называть эту репликацию просто «асинхронной» не вполне корректно: при acks=all подтверждение продюсеру выдаётся только после того, как запись получена всеми репликами ISR, то есть на пути подтверждения репликация фактически синхронна.
Разрыв в нагрузке между лидером и фолловером тоже меньше, чем принято думать: фолловер пишет на диск ровно те же байты, а начиная с Kafka 2.4 (KIP-392) может ещё и обслуживать чтения.
При создании топика Kafka распределяет партиции циклически и назначает для каждой «предпочитаемую» (preferred) реплику - первую в списке. Баланс нарушается при отказах: когда брокер выходит из строя, лидерство его партиций переходит на другие узлы, а после возвращения он поднимается уже в роли фолловера, и без вмешательства перекос сохраняется.
Восстанавливает баланс preferred leader election - возврат лидерства предпочитаемым репликам. Автоматически этим занимается контроллер: auto.leader.rebalance.enable (по умолчанию true), leader.imbalance.check.interval.seconds (300) и leader.imbalance.per.broker.percentage (10%) задают частоту проверки и допустимый перекос. Вручную - командой kafka-leader-election.sh с типом PREFERRED.
Автоматический режим не так однозначен, как кажется: массовый перенос лидерства сам по себе даёт всплеск задержек и кратковременные ошибки NOT_LEADER_OR_FOLLOWER у клиентов. Часть эксплуатантов нагруженных кластеров выключает auto.leader.rebalance.enable и проводит балансировку в контролируемое окно, в том числе средствами Cruise Control, который анализирует метрики и предлагает перераспределение с учётом реальной нагрузки на лидеров и фолловеров. Выбор зависит от того, что для вас дороже: редкие незапланированные всплески latency или ручной контроль.
Балансировка между брокерами
Перекос возникает, если крупный топик неудачно распределён или если после добавления новых брокеров на них не перенесли часть существующих партиций - новые узлы сами по себе данные не забирают. Проверять распределение можно через kafka-topics.sh --describe, брокерские метрики или Cruise Control, перераспределять - утилитой kafka-reassign-partitions.sh.
Обязательное условие безопасности, которое нельзя пропускать: перераспределение всегда запускают с ограничением скорости. Нетроттлированный реассайн на нагруженном кластере насыщает сеть и диски, ISR схлопывается, продюсеры с acks=all начинают получать ошибки - это один из самых надёжных способов уронить прод «плановой операцией». Ограничение задаётся флагом --throttle (он выставляет leader.replication.throttled.rate и follower.replication.throttled.rate), а после завершения обязательно снимается запуском с --verify, иначе троттлинг останется и будет тормозить штатную репликацию. Начинайте с консервативного значения и повышайте, наблюдая за UnderReplicatedPartitions.
Проверьте заодно дефолты для автоматически создаваемых топиков: num.partitions по умолчанию равен 1. Для кластера из 6 брокеров разумно задать num.partitions=6, чтобы новые топики сразу распределялись по всем узлам. Ещё лучше в продакшене выключить автосоздание вовсе (auto.create.topics.enable=false) и заводить топики явно с осознанными параметрами.
Не забывайте и о балансе внутри узла: при нескольких каталогах в log.dirs данные распределяются по дискам, и перекос приводит к тому, что один диск заполняется раньше остальных.
Rack-awareness и чтение с фолловеров
Если кластер развёрнут в нескольких зонах доступности или стойках, задайте на каждом брокере broker.rack. Kafka учтёт его при размещении и постарается разложить реплики одной партиции по разным зонам - отказ целой зоны не приведёт к потере всех копий.
Вторая, часто недооценённая настройка - чтение с ближайшей реплики (KIP-392, с версии 2.4). По умолчанию все чтения обслуживает лидер, из-за чего консьюмер в зоне A вычитывает данные у лидера в зоне B, генерируя межзональный трафик. У облачных провайдеров он тарифицируется отдельно и на больших объёмах становится заметной статьёй расходов. Установив на брокерах replica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector, а на консьюмерах client.rack, вы разрешаете читать с локальной реплики-фолловера.
Побочный эффект, который надо учитывать: фолловер отдаёт данные только до high watermark, поэтому такое чтение может отставать от чтения с лидера на время репликации.
Настройки надёжности и устойчивости к сбоям

Минимальное число синхронных реплик
min.insync.replicas задаёт минимальное количество реплик, включая лидера, которые должны быть в ISR, чтобы лидер принял запись при acks=all. Настраивается на уровне топика или брокера.
Критически важная деталь: значение по умолчанию равно 1. То есть один только acks=all без явной настройки min.insync.replicas не защищает ни от чего - при усохшем до одного лидера ISR запись будет успешно подтверждена.
В кластере с RF=3 выставляют min.insync.replicas=2: тогда продюсер с acks=all получит подтверждение, только если сообщение записано минимум на две реплики из трёх. Если синхронных реплик меньше, запись отклоняется ошибкой NotEnoughReplicasException или NotEnoughReplicasAfterAppendException - приложение обязано её корректно обрабатывать, а не считать фатальной.
Значение должно быть строго меньше RF. Если выставить его равным RF (3 из 3), выпадение любой единственной реплики остановит приём данных в топик до её восстановления: вы обмениваете доступность на прочность в пропорции, которая почти никогда не нужна.
Unclean leader election
По умолчанию при отказе лидера Kafka выбирает новым лидером только реплику из ISR - синхронную с лидером на момент последнего подтверждённого сообщения. Это гарантирует, что ни одно подтверждённое сообщение не потеряется. Но если в ISR не осталось ни одной живой реплики, Kafka будет ждать возвращения одного из участников, и партиция останется недоступной и для записи, и для чтения.
unclean.leader.election.enable (по умолчанию false начиная с версии 0.11) позволяет нарушить это правило и выбрать лидером отставшую реплику вне ISR. Сообщения, не успевшие реплицироваться, теряются, зато партиция становится доступной без ожидания упавшего узла.
Там, где простой недопустим, эту опцию иногда включают осознанно, принимая риск. Для критичных данных так делать не стоит - правильнее обеспечить достаточную избыточность и производительность репликации, чтобы чистое переключение всегда было возможно.
Если опция включена, обязательно мониторьте UncleanLeaderElectionsPerSec: каждое срабатывание означает состоявшуюся потерю данных, и об этом надо знать, а не узнавать постфактум от потребителей.
Отставание реплик
replica.lag.time.max.ms определяет, как долго фолловер может не догонять лидера, прежде чем лидер исключит его из ISR. По умолчанию 30000 мс - увеличено с 10 секунд в Kafka 2.5 (KIP-537).
Направление компромисса здесь часто описывают наоборот. Быстрое исключение отставших реплик повышает доступность на запись: сжавшийся ISR перестаёт ждать медленного участника, и подтверждения при acks=all выдаются быстрее. Платой является снижение прочности - сообщение теперь считается подтверждённым меньшим числом копий, то есть окно возможной потери при последующем отказе лидера расширяется. Длинный таймаут, наоборот, дольше удерживает медленные реплики в ISR и сохраняет число копий, но при acks=all продюсер ждёт самого медленного участника, и это напрямую бьёт по p99 записи.
Поэтому replica.lag.time.max.ms нельзя настраивать в отрыве от min.insync.replicas: последний служит предохранителем, не позволяя ISR сжаться до опасного уровня незаметно - вместо тихой потери прочности вы получаете явную ошибку записи.
Слишком малое значение к тому же даёт «дребезг» ISR из-за кратковременных сетевых задержек и пауз GC. В подавляющем большинстве случаев дефолт оптимален, если вы его меняете, следите за IsrShrinksPerSec и IsrExpandsPerSec - регулярные сокращения и расширения означают, что таймаут подобран неудачно либо репликация не успевает (см. num.replica.fetchers).
Fsync, page cache и видимость данных
Kafka по умолчанию не выполняет fsync на каждую запись, полагаясь на фоновый сброс средствами ОС - значит, подтверждённые сообщения какое-то время существуют только в page cache. Прочность здесь обеспечивает не диск, а репликация.
Механику видимости при этом часто описывают неверно. Сообщение становится доступным консьюмеру только после того, как оно реплицировано на все реплики ISR и high watermark продвинулся за него. Это правило не зависит от acks: параметр определяет момент, когда ответ получит продюсер, но не момент, когда запись увидит читатель. При acks=1 сообщение не становится видимым раньше - снижение acks сокращает время подтверждения записи, но не сквозную задержку доставки. При сбое лидера записи, не преодолевшие high watermark, новому лидеру не отдаются и просто усекаются.
Даже без принудительного fsync репликация даёт высокую сохранность: потеря возможна только при одновременном отказе всех реплик партиции - например, при обесточивании целой стойки, против чего работает rack-awareness. Уменьшать log.flush.interval.ms или log.flush.interval.messages ради синхронной записи каждого сообщения не стоит: производительность падает резко, а выигрыш перекрывается достаточным RF, min.insync.replicas и здоровой репликацией.
Исключение - единственный брокер без репликации. Там прочность обеспечить больше нечем, и частый flush оправдан с полным пониманием его стоимости.
Транзакции и exactly-once
Транзакции Kafka опираются на идемпотентность продюсера и координатор транзакций на стороне брокера. Идемпотентность включается через enable.idempotence=true - в версиях 3.0 и новее это дефолт, продюсер нумерует записи, брокер отслеживает последовательность и отбрасывает повторы. Требования max.in.flight=1 при этом нет, вопреки распространённому мнению (разбирали выше, в разделе про продюсера).
Для полноценных транзакций продюсеру нужен уникальный и стабильный transactional.id, а также осмысленный transaction.timeout.ms (по умолчанию 60000 мс), по истечении которого координатор принудительно прервёт зависшую транзакцию.
Внутренние топики __transaction_state и __consumer_offsets Kafka создаёт сама, но с параметрами transaction.state.log.replication.factor=3, transaction.state.log.min.isr=2 и offsets.topic.replication.factor=3. Именно поэтому транзакции не запускаются на кластерах из одного-двух брокеров, и на dev-стендах эти значения приходится понижать явно.
И самое главное, о чём забывают чаще всего: транзакции на стороне записи бесполезны без соответствующей настройки чтения. Консьюмер обязан работать с isolation.level=read_committed, иначе он увидит данные прерванных транзакций, и вся схема развалится.
Для сквозного сценария «прочитал - обработал - записал» смещения исходного топика фиксируются внутри транзакции методом sendOffsetsToTransaction() - это и делает атомарной связку «запись результата плюс продвижение оффсета».
Оверхед транзакций не так велик, как принято считать: для типовых конвейеров он измеряется единицами процентов пропускной способности. Но он есть и растёт при коротких транзакциях с частыми коммитами, поэтому размер транзакционного батча подбирают по замерам, а не по интуиции.
Мониторинг и тюнинг производительности

Kafka отдаёт метрики через JMX и на брокерах, и в клиентских библиотеках. В продакшене их снимают экспортером (обычно Prometheus с JMX Exporter) и смотрят на дашборде, без клиентской половины любое изменение конфига продюсера или консьюмера остаётся непроверенным.
Клиентские метрики не менее важны брокерских: значительная часть параметров, разобранных выше, живёт в приложениях, и проверить их эффект со стороны брокера невозможно. Имена метрик ниже приведены так, как они реально называются в JMX - чтобы их можно было напрямую перенести в конфигурацию экспортера.
Метрики нагрузки и задержки
Базовые показатели пропускной способности брокера - в kafka.server:type=BrokerTopicMetrics: MessagesInPerSec, BytesInPerSec, BytesOutPerSec, плюс BytesRejectedPerSec и FailedProduceRequestsPerSec для отказов. Задержки и частота запросов - в kafka.network:type=RequestMetrics с разбивкой по типу (request=Produce, FetchConsumer, FetchFollower): RequestsPerSec даёт частоту, TotalTimeMs - полное время обработки.
Главный диагностический инструмент - не само TotalTimeMs, а его разложение на фазы, доступное там же:
RequestQueueTimeMs- ожидание в очереди перед обработкойLocalTimeMs- обработка на лидере, включая запись в логRemoteTimeMs- ожидание других реплик, то есть фактически репликация приacks=allThrottleTimeMs- задержка из-за квотResponseQueueTimeMsиResponseSendTimeMs- постановка ответа в очередь и его отправка
Именно эта разбивка отвечает на вопрос «мы тормозим на диске, в очереди, на репликации или на троттлинге», и смотреть надо не средние значения, а 95-й и 99-й процентили. Рост RemoteTimeMs при спокойном LocalTimeMs указывает на проблемы репликации, а не дисков, и лечится num.replica.fetchers, а не добавлением I/O-потоков.
Загрузку внутренних пулов показывают NetworkProcessorAvgIdlePercent и RequestHandlerAvgIdlePercent (разбирали в разделе про брокера), а размеры очередей - RequestQueueSize и ResponseQueueSize в kafka.network:type=RequestChannel. Постоянно заполненная очередь запросов означает, что брокер не справляется с потоком: либо не хватает num.io.threads, либо узким местом стали диски. Различить эти случаи помогает LocalTimeMs.
Метрики репликации и лага
UnderReplicatedPartitions (kafka.server:type=ReplicaManager) - число партиций, у которых хотя бы одна реплика отстаёт от лидера. В норме должно быть равно нулю постоянно, ненулевое значение означает либо отказ брокера, либо неспособность фолловера успевать за репликацией.
UnderMinIsrPartitions показывает партиции, где размер ISR упал ниже min.insync.replicas. Её появление серьёзнее: запись в такие партиции с acks=all уже отклоняется.
В обязательный минимум алертов входят ещё две метрики из kafka.controller:type=KafkaController. OfflinePartitionsCount - партиции вообще без лидера, прямая недоступность данных, норма 0. ActiveControllerCount - в сумме по всем узлам кластера должна быть строго равна 1: ноль означает кластер без контроллера, больше единицы - split-brain.
Полезны также IsrShrinksPerSec и IsrExpandsPerSec как индикатор нестабильности репликации и UncleanLeaderElectionsPerSec как индикатор состоявшейся потери данных.
На стороне потребителей основной показатель - лаг, отставание позиции чтения группы от конца лога. Снимается штатной kafka-consumer-groups.sh --describe --group, экспортерами вроде kafka-exporter или средствами Cruise Control. Исторически для этого применялся Burrow, но проект давно развивается слабо, и закладываться на него в новых инсталляциях не стоит. Неуклонно растущий лаг означает, что потребители не успевают - нужно либо увеличивать их число и, возможно, число партиций, либо искать узкое место в самой обработке. Резкий скачок обычно указывает на сбой потребителя или на ребаланс.
Системные ресурсы
CPU. Под высокой нагрузкой брокеры действительно должны нагружать процессор, особенно при включённом TLS и при пережатии данных. Если CPU близок к 100%, дальнейший тюнинг параметров почти бесполезен и эффективнее добавить узлы.
Память. Kafka запускают с относительно небольшим heap - типичный ориентир 5–6 ГБ, - оставляя максимум RAM под файловый кеш ОС: именно page cache обеспечивает чтение «горячих» сегментов без обращения к диску.
Сборщик мусора. По умолчанию и по рекомендации - G1GC, для больших heap имеет смысл рассматривать ZGC с его субмиллисекундными паузами. А вот ParallelGC для брокеров использовать не следует, несмотря на встречающиеся советы «взять его ради максимального throughput»: это полностью stop-the-world сборщик, и его длинные паузы приводят ровно к тем последствиям, которых мы избегаем - брокер пропускает обмен с контроллером, реплики выпадают из ISR, потребители выбрасываются из групп. Время и частоту пауз GC стоит отслеживать как самостоятельный сигнал.
Диски. Профиль нагрузки Kafka - последовательная запись, поэтому вопреки расхожему «только SSD» значительная часть крупных инсталляций успешно работает на HDD в конфигурации JBOD, что существенно дешевле при больших объёмах. SSD или NVMe оправданы при большом числе партиций (доступ становится более случайным), при интенсивных догоняющих чтениях старых сегментов и там, где важен хвост распределения задержек. Наблюдайте за IOPS, временем отклика и длиной очереди диска, распределяйте данные по устройствам через log.dirs (в KRaft-режиме JBOD поддерживается с версии 3.7, KIP-858) и не размещайте на одном физическом диске избыточное число партиций.
Настройки ОС, без которых остальное не имеет смысла: лимит файловых дескрипторов ulimit -n порядка 100 000 и выше - его нехватка одна из самых частых причин отказа брокера при росте числа партиций и соединений, vm.swappiness=1, проверенный vm.max_map_count, поднятые net.core.rmem_max и net.core.wmem_max при работе с большими TCP-буферами, монтирование с noatime и предпочтительно XFS.
В KRaft-кластерах отдельно мониторьте состояние кворума контроллеров и отставание реплик лога метаданных.
Клиентские метрики
Всё, что перечислено ниже, отдаётся клиентами через JMX и замыкает петлю обратной связи: без этих чисел изменение конфига продюсера или консьюмера остаётся непроверенной гипотезой.
Продюсер:
Метрика | На что отвечает |
|---|---|
| реальный размер батча против |
| эффект |
| упирается ли приложение в |
| задержка запроса к брокеру |
| живёт ли продюсер на повторах |
| окупается ли сжатие |
Консьюмер:
Метрика | На что отвечает |
|---|---|
| лаг по худшей партиции, главный показатель |
| эффект |
| наполняемость выборки |
| стоимость коммитов |
| насколько близко вы к |
| стабильность группы |
Практики тюнинга
Первым шагом определите, что именно оптимизируете: пропускную способность, задержку, надёжность или доступность. Эти величины взаимосвязаны, и улучшение одной почти всегда оплачивается другой, поэтому попытка «настроить всё сразу» заканчивается конфигурацией, которая ни в чём не хороша.
Для максимального throughput увеличивают размеры батчей и буферов, число потоков и партиций, включают сжатие. Для минимальной задержки, наоборот, уменьшают batch.size и linger.ms, снижают fetch.min.bytes. Но рефлекторно жертвовать acks и идемпотентностью не стоит: сначала измерьте их реальную стоимость на своём профиле, она часто оказывается меньше ожидаемой.
Влияние изменений проверяют нагрузочным тестом в среде, близкой к боевой. Штатные kafka-producer-perf-test.sh и kafka-consumer-perf-test.sh из дистрибутива позволяют быстро снять базовые цифры с контролем размера сообщений и целевого throughput, для более серьёзных сценариев подходят Trogdor из состава Kafka и OpenMessaging Benchmark. Правило простое: если у изменения нет цифры «было/стало», это не тюнинг, а гадание.
Многие брокерские параметры меняются динамически через kafka-configs.sh --alter --entity-type brokers - часть применяется ко всему кластеру (--entity-default), часть только к конкретному узлу (--entity-name <broker.id>). Те, что требуют рестарта, обновляйте по одному брокеру за раз, дожидаясь обнуления UnderReplicatedPartitions перед переходом к следующему.
Идея полностью автоматического тюнинга по метрикам звучит привлекательно, но на практике плохо применима. Во-первых, клиентские параметры вроде max.poll.records живут в приложениях, и брокер изменить их не может в принципе. Во-вторых, автоматическое изменение конфигурации по порогам без гистерезиса даёт «дребезг», который вредит сильнее исходной проблемы. Рабочая версия той же идеи - не автотюнинг, а алерты по пороговым значениям и заранее написанные runbook-и: метрика подсвечивает проблему, человек применяет заготовленное изменение.
Если кластер достигает пределов даже после тюнинга, переходите к масштабированию: добавляйте брокеров и перераспределяйте партиции - с обязательным троттлингом. А при проблемах именно с объёмом хранения рассмотрите многоуровневое хранилище (KIP-405), позволяющее выносить холодные сегменты в объектное хранилище и держать на брокерах только горячие данные.
По мере роста нагрузки параметры пересматривают: настройки, оптимальные для 100 МБ/с, не подойдут для 1 ГБ/с.
Заключение
Самая частая причина неудачного тюнинга Kafka - не неверный расчёт, а устаревший совет. acks=1 как дефолт, max.in.flight=1 ради сохранения порядка, «поднимите session.timeout до 30 секунд», «репликация встанет, если replica.fetch.max.bytes меньше message.max.bytes», «не больше нескольких сотен партиций на брокер», «только SSD», ParallelGC ради пропускной способности. Часть из этого когда-то была верна, часть не была никогда, а часть сегодня прямо вредит: значение max.in.flight больше 5 не даст продюсеру стартовать вообще.
Вторая причина - дефолты, которые ничего не защищают. min.insync.replicas=1 превращает acks=all в декорацию. default.replication.factor=1 создаёт топики без реплик. num.partitions=1 сажает новый топик на один брокер. Всё это включено из коробки и не подаёт никаких сигналов.
Отсюда порядок работы. Сверяйтесь с документацией своей версии, а не со статьёй в том числе с этой. Меняйте по одному параметру и снимайте цифру «было и стало»: без неё это не тюнинг. И держите перед глазами профиль нагрузки, с которого мы начинали - средний и максимальный размер сообщения, целевой объём в МБ/с, требуемый p99 и допустимость потерь. Конфигурация, оптимальная для 2 млн сообщений по 200 байт, не даст ничего кластеру на 50 тыс. по 100 КБ.

