Разделяемый читатель топика: устройство и ограничения
Разделяемый (общий) читатель — модель чтения из топика, при которой за потребителем закрепляется одно или несколько сообщений, а не целая партиция. Это позволяет нескольким потребителям параллельно обрабатывать сообщения из одной партиции и использовать топики в роли очередей сообщений. Типичный сценарий — обмен сообщениями между микросервисами; чтение по Amazon SQS API работает через разделяемый читатель.
Общее описание модели, настройки читателя, политики DLQ и порядка сообщений см. в разделе Разделяемый (общий) читатель. В этой статье описано внутреннее устройство механизма: как сервер отслеживает состояние обработки, распределяет запросы на чтение по партициям, какие ограничения следует учитывать при проектировании высоконагруженных очередей и какое потребление ресурсов создаёт инфлайт.
Особенность реализации
Для каждой пары «разделяемый (общий) читатель — партиция топика» сервер поддерживает собственное состояние обработки сообщений. Состояние состоит из непрерывного блока сообщений из партиции топика — инфлайта. Первым сообщением инфлайта является сообщение, которое ещё не обработано, находится в обработке или ждёт перемещения в DLQ. Если первое сообщение успешно обработано, то начало блока перемещается на следующее сообщение. Сообщения для обработки отдаются только из сообщений, находящихся в инфлайте.
В инфлайте одной партиции может находиться до 120000 сообщений. Это означает, что если из одной партиции читает большое количество потребителей (больше 120000), то сообщения для обработки будут получать не более 120000 потребителей, даже если в партиции сообщений больше. Увеличение количества партиций в топике увеличивает общий размер инфлайта — общий инфлайт топика равен количеству партиций, умноженному на 120000 сообщений. Если необходимо обрабатывать большое количество сообщений одновременно, то следует увеличить количество партиций топика.
Если топик состоит из нескольких партиций, то запросы на чтение равномерно распределяются по всем партициям. Может случиться ситуация, что запрос попадёт в партицию, в которой отсутствуют сообщения. В этом случае потребитель получит ответ, что сообщений для обработки нет. Если в партиции долгое время нет сообщений для обработки, то она исключается из распределения запросов на чтение. Партиция вернётся в распределение сразу, как только в ней появятся сообщения для обработки.
Для разделяемых (общих) читателей с сохранением порядка сообщений это означает, что если в партицию будет записано подряд большое количество сообщений из одной message-group-id (больше 120000 сообщений), то весь инфлайт займут сообщения из одной message-group-id, и сообщения из других групп выдаваться не будут, пока в инфлайте не появятся сообщения из других групп.
Честные очереди
При выдаче сообщений на чтение сервер учитывает message-group-id сообщений, находящихся в инфлайте. Стараясь отдавать сообщения из разных групп, сервер распределяет нагрузку между писателями более равномерно.
Это позволяет сглаживать пики от отдельных писателей: если один писатель отправляет в очередь гораздо больше сообщений, чем остальные, обработка сообщений от «малых» писателей получает приоритет и не блокируется потоком от одного источника.
Для FIFO-очередей и разделяемых читателей с сохранением порядка сообщений при этом гарантируется порядок сообщений внутри одного message-group-id: потребитель всегда получает сообщения группы в том порядке, в котором они были записаны.
Потребление ресурсов
Состояние инфлайта потребляет ресурсы оперативной памяти и диска. На каждое сообщение в инфлайте сервер хранит служебную информацию объёмом около 32 байт — это не тело сообщения, а метаданные о его статусе обработки.
Оперативная память
Весь инфлайт хранится в памяти сервера. Для оценки потребления памяти удобно считать от миллиона сообщений: 1 млн сообщений в инфлайте ≈ 32 МБ оперативной памяти.
Например, топик с 10 партициями и одним разделяемым читателем при полностью заполненном инфлайте (120000 сообщений на партицию) удерживает до 1,2 млн сообщений в инфлайте — около 38 МБ памяти только на служебное состояние. На одну пару «разделяемый (общий) читатель — партиция» при максимальном инфлайте приходится около 3,8 МБ. Если на топике несколько разделяемых читателей, потребление памяти умножается на их количество.
При проектировании учитывайте:
- чем больше партиций и читателей, тем выше потолок потребления памяти;
- чем дольше сообщения находятся в обработке (большое время обработки сообщения, медленные потребители), тем дольше они остаются в инфлайте и занимают память;
- увеличение числа партиций расширяет параллелизм, но одновременно увеличивает суммарный объём инфлайта.
Диск
На диске также сохраняется информация о статусе сообщений — от 32 байта на сообщение. Хранение состоит из двух частей:
- Снапшот — периодически записывается полное состояние всех сообщений в инфлайте.
- Журнал изменений — сервер обрабатывает ожидающие запросы на чтение пачками и при обработке каждой пачки дописывает в журнал изменения статусов сообщений.
Следовательно, высокая частота чтения и большой инфлайт увеличивают объём дисковых записей. При планировании ёмкости кластера закладывайте ресурсы диска наряду с оперативной памятью — особенно для топиков с большим числом партиций, несколькими разделяемыми читателями и длительным временем обработки сообщений.