Оптимизация крупных запросов на основе Flux

Оптимизация крупных запросов на основе Flux

1. Почему Flux?

Существующая бизнес-логика действовала на основе последовательной циклической структуры. Она обходила список данных по одному элементу, вызывая внешние API или БД, собирала результаты и затем выполняла следующую логику, что было очень типичной структурой.

Изначально, поскольку объем данных был небольшим, серьезных проблем не возникало. Однако по мере увеличения масштабов сервиса и роста количества обрабатываемых опций и объектов запроса, ухудшение производительности начало проявляться более явно.

В частности, в структурах, сильно зависящих от внешних сервисов или баз данных, время ожидания сетевого ввода-вывода накапливалось, что привело к резкому увеличению общего времени отклика.

Например, если вызвать API с средним временем отклика 200 мс последовательно 100 раз, то даже простого вычисления будет достаточно, чтобы затратить более 20 секунд. Проблема в том, что, несмотря на наличие достаточных ресурсов ЦП, большая часть времени уходила на ожидание ввода-вывода.

То есть это была неэффективная структура, которая не использовала системные ресурсы должным образом.

Существующий метод имел следующие характеристики.

Во-первых, следующий запрос мог быть начат только после завершения предыдущего.

Во-вторых, время ожидания ответа от сети накопилось в общем времени ответа.

В-третьих, по мере увеличения числа данных время ответа увеличивалось линейно.

Четвертое, если происходит задержка внешнего API, общее время ответа сервиса также увеличивается.

В частности, в средах, где запрашиваются большие объемы опциональных данных, эта проблема проявляется еще более серьезно. В реальной производственной среде в определенные часы происходили одновременные запросы на сотни и более, и в этом процессе регулярно возникали проблемы с задержкой ответов и тайм-аутами.

Первым делом, которое я рассматривал для решения этой проблемы, была параллельная обработка.

Поскольку это была независимая операция запроса, не было необходимости обрабатывать ее последовательно.

Особенно, поскольку большинство задач в основном зависело от ожидания сетевого ввода-вывода, можно было ожидать значительного улучшения производительности, если правильно применить параллельную обработку. В этом процессе была выбрана технология на основе Reactor - Flux. Flux не является простым инструментом для перебора коллекций, а представляет собой библиотеку на основе реактивных потоков, которая позволяет обрабатывать данные в виде потоков и естественным образом сочетать параллельную обработку и асинхронное выполнение.

Прежде всего, я выбрал Flux по следующим причинам.

Во-первых, удалось надежно обрабатывать большие объемы данных в виде потоков.

Во-вторых, с помощью параллельных рельсов легко настроить параллельную обработку.

В-третьих, удалось детально контролировать ресурсы потоков на основе Scheduler.

В-четвертых, благодаря структуре на основе Backpressure удалось предотвратить неумеренное увеличение потоков.

В-пятых, мы смогли легко интегрироваться с существующими системами на основе Java и Spring. Что наиболее важно, так это возможность «сократить общее время обработки до уровня самого длительного одиночного запроса».

В существующем последовательном процессе, если есть n запросов, общее время выполнения имело следующую структуру.

Существующий подход:

n × среднее время отклика

С другой стороны, в параллельной структуре обработки это могло быть улучшено в следующей форме.

Способ улучшения:

MAX(индивидуальное время ответа) + α(накладные расходы расписания)

То есть, если параллельных обработок достаточно, общее время ответа будет сходиться к уровню времени самой длительной одиночной задачи.

2. Процесс применения и стратегия реализации

В процессе фактического применения возникли проблемы, которые не решались простым применением асинхронного вызова. В производственной среде нужно было учитывать не только производительность, но и стабильность, целостность данных, а также управление ресурсами потоков.

Особенно важным моментом было требование, что «результаты параллельной обработки должны быть снова собраны синхронно». В реальной бизнес-логике все результаты запросов должны быть собраны, прежде чем могут быть выполнены расчеты и последующие обработчики.

То есть, запросы можно выполнять параллельно, но окончательные результаты должны обязательно объединяться синхронно.

(1) Конфигурация потока на основе Flux

Сначала данные для запроса были преобразованы в потоковый формат с помощью Flux.fromIterable(). Затем была использована .parallel(size) для создания параллельных рельсов.

Значение size здесь не определяется только на основе количества ядер CPU. Мы скорректировали количество параллельной обработки, учитывая такие факторы, как время отклика внешнего API, нагрузка на БД, задержка сети и состояние памяти сервера.

Слишком высокое количество параллельных потоков может привести к чрезмерной нагрузке на внешние сервисы. Поэтому в рабочей среде мы также учитывали следующие критерии.

• Лимит частоты внешнего API

• Размер пула подключений к базе данных

• Среднее время ответа

• Использование памяти сервера

• Использование ЦПУ

• Состояние сетевой задержки

(2) Оптимизация планировщика

Еще одним важным элементом параллельной обработки было управление потоками.

Сначала использовался простой parallel(), но в реальной рабочей среде стало очевидно, что стратегия управления потоками имеет огромное значение.

Особенно в задачах, сосредоточенных на сетевом I/O, существует возможность блокировки, поэтому использование только обычного пула потоков, ориентированного на CPU, было нецелесообразно. Чтобы решить эту проблему, мы использовали Schedulers.boundedElastic(). boundedElastic - это эластичный пул потоков, предоставляемый Reactor. Он позволяет увеличивать количество потоков при необходимости, но накладывает ограничения, чтобы предотвратить неограниченное увеличение.

То есть, была предложена структура, подходящая для задач, ориентированных на ожидание ввода-вывода, при этом сохраняющая стабильность системы.

В частности, имелись следующие преимущества.

Во-первых, при необходимости потоки могли быть гибко расширены.

Во-вторых, неиспользуемые потоки были автоматически очищены.

В-третьих, удалось сократить риск OOM, вызванный неограниченным созданием потоков.

В-четвёртых, работал стабильно и в средеBlocking I/O.

В реальной операционной среде «надежное использование ресурсов» было гораздо важнее, чем простая производительность. boundedElastic был очень подходящим выбором с точки зрения операционной стабильности.

(3) Синхронизация на основе CountDownLatch

Самая большая проблема заключалась в том, как сочетать асинхронность и синхронность. Благодаря асинхронной обработке скорость запроса значительно улучшилась, но в конечном итоге все результаты запросов должны были быть собраны, чтобы можно было выполнить следующую бизнес-логику.

Таким образом, процесс запроса асинхронный, но конечный поток должен управляться синхронно. Для решения этой задачи мы использовали CountDownLatch. Каждый поток асинхронно вызывает API и сохраняет результаты в responseBodyList.

И каждую раз, когда работа завершается, я вызываю countDown(). Основной поток настроен так, чтобы ожидать завершения всех работ через await().

Если подвести итоги структуры, она выглядит следующим образом.

• Асинхронный : Параллельный вызов API в каждом потоке и сохранение результатов в responseBodyList

• Синхронный : Ожидание завершения всех потоков через CountDownLatch.await()

Главное преимущество этого метода заключается в том, что удалось безопасно внедрить параллельную обработку, не изменяя значительно существующую структуру бизнес-логики.

Таким образом, мы смогли параллелизовать только внутреннюю логику запросов, сохранив при этом существующую синхронную бизнес-структуру.

3. Результаты решения проблем и улучшения производительности

После применения параллельной обработки на основе Flux производительность значительно улучшилась. В предыдущем последовательном методе время ответа увеличивалось линейно с увеличением количества данных.

Например, если выполнить 100 заданий со средним временем ответа 300 мс, общее время ответа может возрасти до около 30 секунд.

Однако в параллельной архитектуре большинство запросов выполнялись одновременно, поэтому общее время ответа свелось к уровню времени самого длительного одиночного запроса. В реальной рабочей среде были подтверждены следующие улучшения.

• Существующий способ:

n × среднее время ответа

• Способ улучшения:

MAX(индивидуальное время ответа) + α(затраты на планирование и объединение)

Особенно эффект параллельной обработки оказался очень значительным в средах с задержками сети. В прежнем подходе, если определенный ответ API задерживался, вся последовательность также задерживалась.

С другой стороны, в структуре параллельной обработки, даже если некоторые запросы были медленными, другие запросы могли обрабатываться одновременно, что значительно повысило общую воспринимаемую производительность.

Мы также смогли получить гораздо более эффективные результаты с точки зрения использования ЦП. Ранее большая часть времени тратилась на ожидание ввода-вывода, но после параллельной обработки время простоя значительно сократилось.

С точки зрения управления также произошли положительные изменения.

Во-первых, частота возникновения тайм-аутов при массовом запросе опций снизилась.

Во-вторых, стабильность ответов в периоды пиковых нагрузок улучшилась.

В-третьих, процент случаев, приводящих к полному сбою сервиса из-за задержек внешних API, снизился.

В-четвертых, возможность регулировать количество параллельных процессов в зависимости от среды увеличила операционную гибкость.

Особенно полезно то, что с помощью только регулирования параллельных чисел можно легко настроить баланс между производительностью и стабильностью.

4. Заключение

В ходе этого опыта оптимизации я еще раз осознал, что гораздо важнее, чем просто использование новейших технологий, — это найти "оптимальное сочетание инструментов, соответствующее бизнес-ограничениям". Сначала я думал, что, просто применив асинхронную обработку, все проблемы будут решены.

Но на самом деле, в рабочей среде необходимо учитывать не только производительность, но и надежность, целостность данных и поддерживаемость эксплуатации.

Особенно я смог понять, что полностью менять все операции на асинхронную структуру не всегда является правильным решением.

В реальном проекте, поскольку для выполнения следующего этапа логики все результаты были необходимы, сбор синхронных результатов после асинхронной обработки был абсолютно необходимым.

В этом процессе мы убедились, что смешанная модель, использующая такие инструменты синхронизации, как CountDownLatch, может стать очень эффективным решением.

Также было очень значимым опытом подойти к Flux не только с точки зрения «современных технологий», но и как к исполняемой модели, подходящей для обработки больших объемов ввода-вывода.

В конце концов, важным является не сама технология, а выбор наиболее подходящего способа, соответствующего текущей архитектуре системы и операционной среде.

Этот опыт имел значение, превышающее простое улучшение производительности. Я считаю, что это был практический опыт оптимизации, который также учитывал стабильность и поддерживаемость в операционной среде.

Джек

Site footer