При увеличении числа потребителей очереди обработка перестала ускоряться. Как определить, что узким местом стало число партиций?
Если сообщения распределяются по партициям, максимальный параллелизм потребления обычно ограничен числом партиций: одну партицию в рамках consumer group в каждый момент обрабатывает только один потребитель. Если потребителей больше, чем партиций, лишние экземпляры простаивают, поэтому дальнейшее масштабирование не увеличивает пропускную способность.
Партиционирование появилось как способ масштабировать потоковую обработку горизонтально, не отказываясь полностью от упорядочивания сообщений. Поток делят на независимые части, а обработку каждой части передают отдельному потребителю.
Такой подход сочетает параллелизм между партициями с последовательной обработкой внутри одной партиции. Он используется в системах журналов и очередей, где важны высокая пропускная способность и возможность повторного чтения сообщений.
Допустим, очередь содержит четыре партиции, а сервис запущен в десяти экземплярах, объединённых в одну consumer group. Только четыре экземпляра смогут одновременно владеть партициями; остальные будут ждать назначения работы.
Если производительность одной партиции недостаточна, увеличение числа потребителей выше четырёх не устранит задержку. Неверный вывод о нехватке экземпляров может привести к лишним расходам, хотя реальное ограничение находится в структуре очереди или в обработке самих сообщений.
Сначала сравнивают количество активных партиций с количеством потребителей внутри одной группы. Если потребителей больше, чем партиций, это сильный признак, что достигнут потолок параллелизма на уровне распределения сообщений.
Затем проверяют метрики: загрузку и задержку каждой партиции, размер отставания потребителя, время обработки сообщения и наличие свободных потребителей. Если несколько партиций постоянно имеют большое отставание, а дополнительные потребители не получают назначений, узким местом, вероятно, является число партиций.
Увеличение числа партиций позволяет поднять потенциальный параллелизм, но не бесплатно. Оно может вызвать перераспределение данных, изменить правила маршрутизации ключей и нарушить прежнюю привязку связанных сообщений к конкретным партициям.
Если порядок нужен для сообщений одного объекта, например заказа, их обычно направляют в одну партицию по ключу объекта. Тогда порядок между разными заказами можно обрабатывать параллельно, но сообщения одного заказа сохраняют последовательность. При этом один «горячий» ключ всё равно способен ограничить производительность одной партиции.
Важно отличать ограничение числа партиций от других причин: медленной базы данных, внешних вызовов, блокировок, ограниченного пула соединений или неравномерного распределения ключей. Поэтому одного наблюдения о простое части потребителей недостаточно — нужно сопоставить его с загрузкой партиций и временем обработки.
Сервис обрабатывает события заказов из восьми партиций. После увеличения числа экземпляров с восьми до двадцати задержка обработки почти не изменилась: восемь экземпляров работали, остальные оставались без назначенных партиций.
Рассматривались два варианта. Увеличение числа потребителей было дешёвым и обратимым, но не давало дополнительного параллелизма. Увеличение числа партиций повышало пропускную способность, однако требовало проверить распределение ключей и влияние на порядок событий.
Выбрали увеличение числа партиций после анализа ключей заказов и нагрузки. События одного заказа сохранили общий ключ, а разные заказы стали распределяться по большему числу партиций. После этого добавление потребителей стало реально увеличивать параллелизм; отдельно контролировали «горячие» ключи и отставание каждой партиции.
Внутренние потоки могут параллельно обрабатывать сообщения, уже полученные потребителем, но это не увеличивает число независимых владельцев партиций. Если порядок внутри партиции должен сохраняться, параллельная обработка требует дополнительной координации и может привести к завершению более позднего сообщения раньше раннего.
Такой подход допустим, когда порядок не нужен либо его контролируют отдельно. Но он не отменяет ограничение на чтение и распределение самой партиции и может создать новые проблемы с подтверждением прогресса.
Глобальный порядок обычно отсутствует уже при наличии нескольких партиций: сообщения из разных партиций могут обрабатываться независимо. При изменении числа партиций может измениться партиция, в которую попадёт ключ, если алгоритм маршрутизации зависит от количества партиций.
Поэтому гарантировать порядок следует не глобально, а в требуемой области — например, для одного заказа. Для этого все события заказа должны стабильно маршрутизироваться в одну партицию с учётом правил миграции и особенностей используемой очереди.
Партиции могут быть загружены неравномерно. Например, один популярный ключ или сегмент данных создаёт значительно больше сообщений, чем остальные, а связанные сообщения нельзя безопасно распределить между несколькими партициями.
В результате один потребитель будет перегружен, а другие — простаивать. Решение зависит от требований к порядку: можно изменить ключ маршрутизации, разделить горячий объект на независимые подканалы либо оптимизировать обработку, но простое добавление потребителей проблему одной перегруженной партиции не устранит.