Обнаружение сбоев у потребителей#
После подписки на набор топиков потребитель автоматически присоединится к группе при вызове опроса poll(Duration). API опроса разработан для поддержания активности потребителей. Пока опросы вызываются, потребитель будет оставаться в группе и продолжать получать сообщения из назначенных ему разделов. При этом потребитель периодически посылает на сервер сигнал, что с ним все в порядке. Если потребитель выходит из строя или не отправляет сигнал в течение времени, заданного конфигурационным параметром session.timeout.ms, то он будет считаться отключенным, и его разделы будут переназначены.
Иногда случается ситуация «живой блокировки»: потребитель продолжает посылать сигнал «все в порядке», но в реальности никакой деятельности не выполняет. Чтобы в этом случае потребитель не мог бесконечно удерживать свои разделы, применяется механизм обнаружения активности с использованием параметра max.poll.interval.ms.
Таким образом, если опрос не вызывается по крайней мере с частотой заданного максимального интервала, то клиент сам покинет группу, чтобы другой потребитель мог занять его разделы. Когда это происходит, возможен сбой фиксации смещения (в этом случае будет выдано исключение CommitFailedException в результате вызова метода commitSync()). Это механизм безопасности, который гарантирует, что возможность фиксировать смещения будет только у активных членов группы. Поэтому, чтобы остаться в группе, нужно продолжать вызывать опросы.
Потребитель предоставляет два параметра конфигурации для управления поведением цикла опроса:
max.poll.interval.ms— увеличивая интервал между ожидаемыми опросами, можно предоставить потребителю больше времени для обработки пакета записей, возвращенных изpoll(Duration). Недостаток данного способа заключается в том, что увеличение этого значения может задержать перебалансировку группы, поскольку потребитель присоединится к перебалансировке только внутри вызова опроса. С помощью этого параметра можно ограничить время для завершения перебалансировки, однако есть риск замедлить прогресс, если потребитель на самом деле не может достаточно часто вызыватьpoll.max.poll.records— позволяет ограничить максимальное количество записей, возвращаемых за один вызов опроса. Это может упростить прогнозирование максимума, который должен быть обработан в течение каждого интервала опроса. Настроив это значение, можно сократить интервал опроса, что уменьшит влияние перебалансировки групп.
В случаях, когда время обработки сообщений непредсказуемо варьируется, ни один из этих вариантов не сможет помочь наверняка. В таких случаях рекомендуется перенести обработку сообщений в другой поток, что позволит потребителю продолжать вызывать poll, пока процессор все еще работает. При этом необходимо строго следить, чтобы зафиксированные смещения не опережали фактическую позицию. Как правило, следует отключить автоматическую фиксацию и вручную фиксировать обработанные смещения для записей только после того, как поток завершит их обработку (в зависимости от необходимой семантики доставки).
Внимание
Чтобы новые записи из опроса не поступали до тех пор, пока поток не завершит обработку ранее возвращенных записей, приостановите раздел.