Страницы

Поиск по вопросам

Показаны сообщения с ярлыком long-poll. Показать все сообщения
Показаны сообщения с ярлыком long-poll. Показать все сообщения

суббота, 8 февраля 2020 г.

Остановка выполнения цикла while

#comet #while #php #sleep #long_poll


Всем привет! Интересует вопрос - как можно правильно остановить выполнение while
цикла со sleep внутри. 
Дело в том, что я реализовываю "long poll" для быстрого обновления информации на
сайте (оповещения, другая информация).
СУТЬ РАБОТЫ СКРИПТА 
1. Мы отправляем GET-запрос на PHP-страницу в которой расположен скрипт.
2. Скрипт запускает цикл while.
3. В данном цикле существует 2 условия и sleep(1) в конце, который приостанавливает
выполнение цикла.
- Если если цикл работает больше 30 секунд - останавливаем выполнение цикла, соответственно
и скрипта. 
- Если на сайте появилась новая информация, скажем "оповещение" - выводим результат
пользователю и останавливаем цикл.
ПРОБЛЕМА 
Если пользователь закроет окно браузера, или же перезагрузит страницу - выполнится
новый запрос на данный скрипт, который запустит цикл повторно. Но при этом, предыдущие
выполнение  цикла while не остановятся, пока не сработает хотя бы одно условие внутри
него. Если же пользователь отошлет много запросов на страницу - память сервера начинает
сильно грузится и в результате все висит, пока все while не остановят свое выполнение.
ВОПРОС 
Как остановить выполнение цикла и скрипта в целом при повторном выполнении или же
перезагрузке страницы клиентом? 
Пока на данный момент я сделал следующее.
При выполнении скрипта мы генерируем рандомное число, которое записываем в переменную
и так же с сессию. Далее внутри while делаю условие, в котором данная переменная должна
равняться значению сессии. Если false - останавливаем цикл. 
Если мы запускаем скрипт по новой - значение сессии меняется и соответственно все
предыдущие выполнения while останавливаются. 
Можно ли использовать такой вариант? И на сколько он безопасен?
UPDATE.PHP
 // инициализируем сессию пользователя
 $session = Session::instance();

        $rand = rand();
        $session->set('user_update_key', $rand);
        $key = $rand;

        $limit = 20;
        $seconds = 0;

        set_time_limit($limit + 1);

        while (TRUE) {
            if (Session::instance()->get('user_update_key') == $key) {

                if (есть ли обновления. если есть - выводим) {
                    echo 'информация обновлений';
                    flush();
                    exit;
                }
                if ($seconds == $limit) { // завершаем выполнение цикла по истечению
времени.
                    echo 'Close';
                    flush();
                    exit;
                }
                $seconds++;
            } else {
                unset($session, $rand, $key, $notice, $mysession, $last_notice, $limit,
$notise_return, $seconds);
                flush();
                exit;
            }
            session_write_close();
            sleep(1);
            session_start();

}    


Ответы

Ответ 1



Похоже, вас интересует двусторонняя коммуникация пользователя с долгоиграющим процессом на PHP, и процессов PHP между собой во время их исполнения. PHP под это не заточен, как уже писали, но есть расширения, которые могут оказаться полезными. Без дополнительных примочек подходит механизм сессий, или же похожий собственный, задействующий, например, БД как "общую память". В первом приближении, я бы сделал точно то же, что и вы. Почему вы хотите именно long polling, а не новый разовый запрос раз в секунду? VK виджеты для комментариев, скажем, именно так работают. Если веб сервер настроен не закрывать мгновенно каждое соединение, то цепочка таких запросов не станет тормозить. (См. настройку KeepAlive для Apache) Наконец, можно посомтреть в сторону других технологий и протоколов, напр. Flash и RTMP: маленькая флешка в веб странице, которая по RTMP связывается с сервером RED5, и получает/отправляет инфу через JS коммуникацию с флешкой - настоящий риалтайм!

Ответ 2



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

Ответ 3



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

суббота, 11 января 2020 г.

Постоянное подключение на PHP со свойствами LongPolling и WebSocket

#php #websocket #long_poll


LongPolling работает так, что клиент подключается к серверу по протоколу http(s)
и зависает, пока сервер не ответит. После этого создается новое подключение. Это подключение
происходит абсолютно стандартным образом с отправкой всех заголовков и cookie. Для
каждого такого подключения апач (или другой сервер) создает изолированный процесс.

WebSocket в свою очередь подключается к специальному серверу по отдельному протоколу
ws(s). Выходит, что есть всего один процесс, в котором должна проходить авторизация,
и велика вероятность утечек памяти. Но преимущества этого подхода в том, что есть постоянное
подключение, что обеспечивает моментальную реакцию на события. Помимо прочего, в случае
работы с WS не обязательно хранить сообщения где-то в БД, чтобы ничего не было потеряно
при переподключении.

Вопрос вот в чем: есть ли что-то среднее между LP и WS?


Работает по протоколу http(s) со всеми плюшками (headers и cookies)
Можно написать простой PHP скрипт, который будет запускаться по запросу и обрабатывать
все эти подключения, а не пилить полноценный WS сервер, который нужно будет еще мониторить
и следить за тем, чтобы он правильно работал. (P.S. Я в проекте единственный разработчик/дизайнер/верстальщик/сис.админ/другое
слово, поэтому не хочется добавлять себе еще хлопот)
Для каждого подключения спаунит отдельный процесс, изолированный от других подключений
Имеет постоянное подключение. И чтобы можно было отправлять сообщения хотя бы в одну
сторону (сервер => клиент)




Подробнее о том, в чем собственно проблема.

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

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

Смотрел в строну Workerman. Опять же возникает вопросы аутентификации, поддержки
сервера в рабочем состоянии при утечках и крашах (supervisord?), перезапуска при деплое 


  Если я где-то заблуждаюсь, тыкните, пожалуйста, пальцем.

    


Ответы

Ответ 1



Достойных альтернатив для WebSocket в наше время считайте что нет. Или WebSocket, или всевозможные костыли. Если без костылей, то... Первое, что нужно сделать: настроить nginx для передачи WebSocket соединений дальше по цепочке. Так, все ваши клиенты будут работать с nginx, а не напрямую с вашими другими про­граммами. Сохраняются и используются все достоинства nginx. Не нужно использовать выде­ленный порт. Не нужно настраивать HTTPS дважды. Не нужно перезапускать какие-то другие про­граммы после обновления сертификатов Let's Encrypt. Настройка делается в четыре строчки: location /mywsapp { proxy_pass http://mywsbackend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; } Дальше, нужна альтернатива самостоятельной обработке входящих соединений WebSocket. Такая программа будет брать на себя обработку входящих WS соединений во всей их слож­ности, оставляя вам лишь обработку входящих сообщений и отправку исходящих. websocketd Например, программа websocketd работает как inetd, только для процессов. При новом соеди­нении запускается ваша программа, которой нужно читать стан­дартный вход и писать в стан­дартный выход. Можно не читать, только писать. Пример использования в другом ответе.

понедельник, 16 декабря 2019 г.

NGINX: раздача дописываемого файла

#nginx #http #long_poll


В моем приложении есть файл, который по сути напоминает лог т.е. постоянно дописывается.
Хочу использовать хедер Range для как альтернативу long polling и вебсокетам для чтения
дополняющегося файла, так как на реализацию этих технологий нужно время, тем более
с учетом использования С++. Насколько эффективно NGINX справится с такой задачей, возможно
будут проблемы при работе с дополняющимися файлами? И какие параметры будут оптимальными
для реализации этого подхода? Может быть нужно использовать directio в приложении и
NGINX для оптимизации? Или стоит вообще отказаться от этой идеи?...
Сервер на базе linux, файловая система ext4. Размер файла не превышает 100МБ(одновременно
запись идет в один файл, потом создается новый).

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


Ответы

Ответ 1



Не имеет смысла изобретать велосипед так как для nginx уже есть модуль для работы с вебсокетами, который берет всю сложную работу на себя. Вам остается только передавать сообщения для доставки в nginx. Устанавливается на раз: sudo apt install libnginx-mod-nchan Настройка тоже ничего сложного не представляет. Может работать через Redis в конфигурации с множеством серверов. Если вы по какой-то причине не можете использовать этот модуль, то опять же стоит изобретать очередную БД на коленке. На серьёзных нагрузках файлы на диске никогда не заменят настоящую БД, будь это MySQL или что. Хотя бы потому что типичная БД может гарантировать что горячие данные будут в оперативке, тогда как дисковый кеш вообще ничего не гарантирует. Это очень легко проверить простым бенчмарком. Из альтернатив можно рассмотреть, например, websocketd. Эта программа, очень близкая по сути inetd, существует отдельно от nginx. Вам нужно лишь перенаправить WS соединения на неё, что очень просто. location /mywsapp { proxy_pass http://mywsbackend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; } В простейшем виде, без nginx, программа, которая будет читать из файла новый строки и передавать их по WebSocket будет выглядеть так: (tail.sh) #!/bin/bash tail -q -n0 -F /tmp/websocketdata.txt Страница, которая читает данные, выглядит так: (index.html)
Если оба файла в одном каталоге, то запускаем: python -m SimpleHTTPServer 8000 & ./websocketd --port=8080 ./tail.sh & Открываем в браузере http://127.0.0.1:8000/ и в соседней консоли пишем в файл: tee -a /tmp/websocketdata.txt Печатаем текст и смотрим как он появляется в окне браузера в тот же миг. Это читающей программе будут доступны всевозможные переменные как если бы она работала CGI в окружении, например QUERY_STRING и другие. Без необходимости их проверять можно сделать даже так, убрав прослойку из bash: ./websocketd --port=8080 tail -q -n0 -F /tmp/websocketdata.txt Чтобы не усложнять программу для чтения файлов можно проверять доступ к подключению через WebSocket на стороне nginx используя директиву auth_request.

среда, 11 декабря 2019 г.

Правильная организация многопоточной работы Java

#java #многопоточность #long_poll


Предисловие

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

Проблема

Когда я получаю большой список обновлений, они, естественно, обрабатываются каждое
по порядку. Их может быть до 256 за запрос, а запросов я могу делать до 3 в секунду
как минимум, итого - 750 обновлений каждую секунду вполне себе возможный вариант. Допустим,
время обработки одного обновления (которое оповещает о входящем сообщении) занимает
десятую часть секунды. Потом вспоминаем об общем их количестве и понимаем, что, в таком
случае, до следующего сообщения очередь дойдет очень нескоро. Нас это, конечно же,
не устраивает.

Как я пытался эту проблему решить

Естественно, подумал я, логичнее всего будет отдавать обработку всего этого в новый
поток, дабы каждое новое сообщение обрабатывалось мгновенно после получения, а значит
задержек быть не должно. 
Изначально я без задней мысли написал везде new Thread(() -> handle(...)).start();
и думал, что проблем теперь не будет. Лишь потом я задумался, что в таком случае на
каждое обновление будет создаваться новый анонимный поток, и таких потоков уже может
быть по 750 штук каждую секунду. 

Изначально (первые несколько минут) все работало отлично, затем задержки стали расти
и расти, в итоге дойдя до десятков минут. Вполне логично было предположить, что где-то
происходят косяки, засоряется память, сеть и всё прочее, видимо, потоки сами не очень-то
и хотели закрываться сразу после выполнения обработки события. Я подключил профайлер,
ничего особого не увидел, как ни странно - по нагрузке на память, процессор и прочее,
мои классы и объекты были где-то далеко в низу, а в топе были char[] и String, с непонятными
кракозябрами в содержании.  А запущенных потоков было всего не более 150, и созданных
объектов не более 3 миллионов.

Дабы исправить создание огромного количества "беспризорных" потоков, я организовал
действие примерно так: я делаю запрос к longpoll-серверу, получаю кучу обновлений,
отдаю их обработчику. Обработчик - класс, наследующийся от потока, который просто на
вечном цикле берёт обновления из очереди и обрабатывает их. По сути, тут всего один
поток, работающий параллельно с главным, и я просто убрал задержки между запросами
для получения обновлений, однако они всё также обрабатываются друг за другом. Как я
это организовал (исходный код) можно увидеть здесь - в классе LongPoll само взаимодействие
с сервером, в классе UpdatesHandler непосредственно обработка.

Вопрос

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

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

Мне предложили попробовать executor service, который, вроде бы, создает пул потоков,
сам ими управляет и переиспользует, но я не знаю, поможет ли это и вообще в этом ли
проблема? И если использовать его всё-таки, то как правильнее поступить - создать 750
потоков, которые будут переиспользоваться? А не многовато ли? А если нет, то какая
разница, один поток или десять будут обрабатывать сотни обновлений, пусть будет задержка
не в минуту, а немного меньше, это не тот результат, который нужен.

Желаемый результат

Чтобы обработка каждого обновления проходила асинхронно, и чтобы каждое обновление
не ожидало окончания обработки предыдущего, также и чтобы запрос за новым списком обновлений
не ожидал окончания обработки всех предыдущих обновлений. Но необходимо, чтобы не было
потрачено огромное количество ресурсов, которые ограничены, и чтобы был учтен большой
объем обновлений, и чтобы работало раз и на всё время стабильно :) Проблема ещё в том,
что обработка одного сообщения может занять как сотую часть секунды, так и секунд 5
в некоторых случаях, поэтому сделать всего потоков 10 не будет смысла. 
    


Ответы

Ответ 1



Как я указал в самом вопросе, мне посоветовали попробовать ExecutorService, и он действительно решил все проблемы. Логика решения была примерно такой: простая итерация по списку полученных событий времени не занимает, а вот обработка некоторых (сообщений) занимает, и порой немало. Значит, нужно пробегаться по массиву событий, и, если текущее событие оповещает о сообщении, отдавать его обработчику в новом потоке. Можно было бы использовать FixedThreadPool, но тогда мне было не очень понятно, сколько потоков было бы необходимо создать и сколько их вообще нужно. Как подсказала данная статья, есть ещё один хороший вариант — использовать CachedThreadPool. Он создаёт новые потоки, только если текущие заняты, и затем сам их подчищает. Таким образом, большое количество потоков будет создано только в случае большого количества сообщений, которые не успевали бы обрабатываться, но по окончании их обработки, все новые потоки будут "убиты". Конкретная реализация: Используем: ExecutorService service = Executors.newCachedThreadPool(); Отдаём наше "задание", которое требует обработки в новом потоке: service.submit(() -> { // код, который будет выполнен в новом потоке }); Всё, вот так просто всё оказалось. Проверил на практике, среднее количество запущенных потоков стало равно 7-8, а в максимуме доходило до 30, а это вполне приемлемо.

понедельник, 4 марта 2019 г.

Постоянное подключение на PHP со свойствами LongPolling и WebSocket

LongPolling работает так, что клиент подключается к серверу по протоколу http(s) и зависает, пока сервер не ответит. После этого создается новое подключение. Это подключение происходит абсолютно стандартным образом с отправкой всех заголовков и cookie. Для каждого такого подключения апач (или другой сервер) создает изолированный процесс.
WebSocket в свою очередь подключается к специальному серверу по отдельному протоколу ws(s). Выходит, что есть всего один процесс, в котором должна проходить авторизация, и велика вероятность утечек памяти. Но преимущества этого подхода в том, что есть постоянное подключение, что обеспечивает моментальную реакцию на события. Помимо прочего, в случае работы с WS не обязательно хранить сообщения где-то в БД, чтобы ничего не было потеряно при переподключении.
Вопрос вот в чем: есть ли что-то среднее между LP и WS?
Работает по протоколу http(s) со всеми плюшками (headers и cookies) Можно написать простой PHP скрипт, который будет запускаться по запросу и обрабатывать все эти подключения, а не пилить полноценный WS сервер, который нужно будет еще мониторить и следить за тем, чтобы он правильно работал. (P.S. Я в проекте единственный разработчик/дизайнер/верстальщик/сис.админ/другое слово, поэтому не хочется добавлять себе еще хлопот) Для каждого подключения спаунит отдельный процесс, изолированный от других подключений Имеет постоянное подключение. И чтобы можно было отправлять сообщения хотя бы в одну сторону (сервер => клиент)

Подробнее о том, в чем собственно проблема.
В проекте используется framework Yii2 с местной системой авторизации и со всеми его плюшками. Описанный мною подход позволил бы мне аутентифицировать пользователя встроенными в framework методами, а саму логику внедрить в action контроллера.
В противном случае мне нужно будет написать отдельный сервер или использовать готовое решение и как-то интегрировать его с фреймворком.
Смотрел в строну Workerman. Опять же возникает вопросы аутентификации, поддержки сервера в рабочем состоянии при утечках и крашах (supervisord?), перезапуска при деплое
Если я где-то заблуждаюсь, тыкните, пожалуйста, пальцем.


Ответ

Достойных альтернатив для WebSocket в наше время считайте что нет. Или WebSocket, или всевозможные костыли. Если без костылей, то...
Первое, что нужно сделать: настроить nginx для передачи WebSocket соединений дальше по цепочке. Так, все ваши клиенты будут работать с nginx, а не напрямую с вашими другими про­граммами. Сохраняются и используются все достоинства nginx. Не нужно использовать выде­ленный порт. Не нужно настраивать HTTPS дважды. Не нужно перезапускать какие-то другие про­граммы после обновления сертификатов Let's Encrypt.
Настройка делается в четыре строчки:
location /mywsapp { proxy_pass http://mywsbackend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; }
Дальше, нужна альтернатива самостоятельной обработке входящих соединений WebSocket. Такая программа будет брать на себя обработку входящих WS соединений во всей их слож­ности, оставляя вам лишь обработку входящих сообщений и отправку исходящих.
websocketd
Например, программа websocketd работает как inetd, только для процессов. При новом соеди­нении запускается ваша программа, которой нужно читать стан­дартный вход и писать в стан­дартный выход. Можно не читать, только писать. Пример использования в другом ответе.
"; }
Всевозможные заголовки, которые отправляет ваш клиент, доступны через переменные окру­жения. Например, куки видны в $_ENV['HTTP_COOKIE'] и так далее.
Вашим требованиям в части запуска процесса по факту подключения этот подход соответ­ствует. Есть недостатки: если вам повезёт вырасти до сотен или тысяч одномоментных под­ключений, то это будет означать сотни и тысячи висящих в памяти процессов PHP. Впрочем, об этом можно будет подумать потом.
Кроме ограничений по памяти, постоянно работающие процессы PHP создают ещё одну проблему: если вам захочется обновить их логику, то вам потребуется механизм отключения всех таких процессов.
nchan
Если забыть про некоторые требования, то модуль nchan для nginx может ещё больше аб­стра­гировать работу с сообщениями. Для каждого подключающегося клиента вы выделяете какой-то длинный и секретный ID (что-то типа номера почтового ящика), по которому тот подключается к nchan и ждёт сообщения. Вы, со своей стороны, по этому же ID отправляете сообщения по простому HTTP. Если у вас появится больше одного сервера, то забирать сообщения nchan может из Redis.
Установка проще не может быть:
sudo apt install libnginx-mod-nchan
Пример конфигурации.
Недостатки: увеличенная сложность конфигурации nginx и всей архитектуры системы, необходимость хранить идентификаторы клиентов за пределами $_SESSION.
Достоинства: не нужно для каждого соединения держать процесс PHP.

Список будет пополняться.

суббота, 27 октября 2018 г.

NGINX: раздача дописываемого файла

В моем приложении есть файл, который по сути напоминает лог т.е. постоянно дописывается. Хочу использовать хедер Range для как альтернативу long polling и вебсокетам для чтения дополняющегося файла, так как на реализацию этих технологий нужно время, тем более с учетом использования С++. Насколько эффективно NGINX справится с такой задачей, возможно будут проблемы при работе с дополняющимися файлами? И какие параметры будут оптимальными для реализации этого подхода? Может быть нужно использовать directio в приложении и NGINX для оптимизации? Или стоит вообще отказаться от этой идеи?... Сервер на базе linux, файловая система ext4. Размер файла не превышает 100МБ(одновременно запись идет в один файл, потом создается новый).
PS: В файл придется писать в любом случае и разбивка на мелкие файлы не связана с реализацией этого подхода.


Ответ

Не имеет смысла изобретать велосипед так как для nginx уже есть модуль для работы с вебсокетами, который берет всю сложную работу на себя. Вам остается только передавать сообщения для доставки в nginx. Устанавливается на раз:
sudo apt install libnginx-mod-nchan
Настройка тоже ничего сложного не представляет. Может работать через Redis в конфигурации с множеством серверов.
Если вы по какой-то причине не можете использовать этот модуль, то опять же стоит изобретать очередную БД на коленке. На серьёзных нагрузках файлы на диске никогда не заменят настоящую БД, будь это MySQL или что. Хотя бы потому что типичная БД может гарантировать что горячие данные будут в оперативке, тогда как дисковый кеш вообще ничего не гарантирует. Это очень легко проверить простым бенчмарком.

Из альтернатив можно рассмотреть, например, websocketd. Эта программа, очень близкая по сути inetd, существует отдельно от nginx. Вам нужно лишь перенаправить WS соединения на неё, что очень просто.
location /mywsapp { proxy_pass http://mywsbackend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; }
В простейшем виде, без nginx, программа, которая будет читать из файла новый строки и передавать их по WebSocket будет выглядеть так: (tail.sh)
#!/bin/bash tail -q -n0 -F /tmp/websocketdata.txt
Страница, которая читает данные, выглядит так: (index.html)


Если оба файла в одном каталоге, то запускаем:
python -m SimpleHTTPServer 8000 & ./websocketd --port=8080 ./tail.sh &
Открываем в браузере http://127.0.0.1:8000/ и в соседней консоли пишем в файл:
tee -a /tmp/websocketdata.txt
Печатаем текст и смотрим как он появляется в окне браузера в тот же миг.
Это читающей программе будут доступны всевозможные переменные как если бы она работала CGI в окружении, например QUERY_STRING и другие. Без необходимости их проверять можно сделать даже так, убрав прослойку из bash:
./websocketd --port=8080 tail -q -n0 -F /tmp/websocketdata.txt
Чтобы не усложнять программу для чтения файлов можно проверять доступ к подключению через WebSocket на стороне nginx используя директиву auth_request

вторник, 16 октября 2018 г.

Правильная организация многопоточной работы Java

Предисловие
Работаю с технологией longpoll-инга. То есть суть работы заключается в том, что я делаю запросы к определенному серверу, и, в случае наличия обновлений, сервер возвращает мне их список, если нет - через определенное время мне возвращается пустой ответ. Если сказать немного конкретнее - я работаю с VK API и личными сообщениями, поэтому время обработки каждого обновления из списка крайне критично.
Проблема
Когда я получаю большой список обновлений, они, естественно, обрабатываются каждое по порядку. Их может быть до 256 за запрос, а запросов я могу делать до 3 в секунду как минимум, итого - 750 обновлений каждую секунду вполне себе возможный вариант. Допустим, время обработки одного обновления (которое оповещает о входящем сообщении) занимает десятую часть секунды. Потом вспоминаем об общем их количестве и понимаем, что, в таком случае, до следующего сообщения очередь дойдет очень нескоро. Нас это, конечно же, не устраивает.
Как я пытался эту проблему решить
Естественно, подумал я, логичнее всего будет отдавать обработку всего этого в новый поток, дабы каждое новое сообщение обрабатывалось мгновенно после получения, а значит задержек быть не должно. Изначально я без задней мысли написал везде new Thread(() -> handle(...)).start(); и думал, что проблем теперь не будет. Лишь потом я задумался, что в таком случае на каждое обновление будет создаваться новый анонимный поток, и таких потоков уже может быть по 750 штук каждую секунду.
Изначально (первые несколько минут) все работало отлично, затем задержки стали расти и расти, в итоге дойдя до десятков минут. Вполне логично было предположить, что где-то происходят косяки, засоряется память, сеть и всё прочее, видимо, потоки сами не очень-то и хотели закрываться сразу после выполнения обработки события. Я подключил профайлер, ничего особого не увидел, как ни странно - по нагрузке на память, процессор и прочее, мои классы и объекты были где-то далеко в низу, а в топе были char[] и String, с непонятными кракозябрами в содержании. А запущенных потоков было всего не более 150, и созданных объектов не более 3 миллионов.
Дабы исправить создание огромного количества "беспризорных" потоков, я организовал действие примерно так: я делаю запрос к longpoll-серверу, получаю кучу обновлений, отдаю их обработчику. Обработчик - класс, наследующийся от потока, который просто на вечном цикле берёт обновления из очереди и обрабатывает их. По сути, тут всего один поток, работающий параллельно с главным, и я просто убрал задержки между запросами для получения обновлений, однако они всё также обрабатываются друг за другом. Как я это организовал (исходный код) можно увидеть здесь - в классе LongPoll само взаимодействие с сервером, в классе UpdatesHandler непосредственно обработка.
Вопрос
Как лучше в данном случае организовать работу или хотя бы отследить, из-за чего случаются проблемы? В логе всё хорошо, веду полное логгирование через log4j.
Основная проблема заключается в том, что спустя несколько минут, обработка сообщений начинает занимать неприлично большое количество времени, хотя изначально при запуске всё работает как часы. Даже нагрузка не так влияет на это, дело тут в чём-то другом, и в чём - я не могу пока понять.
Мне предложили попробовать executor service, который, вроде бы, создает пул потоков, сам ими управляет и переиспользует, но я не знаю, поможет ли это и вообще в этом ли проблема? И если использовать его всё-таки, то как правильнее поступить - создать 750 потоков, которые будут переиспользоваться? А не многовато ли? А если нет, то какая разница, один поток или десять будут обрабатывать сотни обновлений, пусть будет задержка не в минуту, а немного меньше, это не тот результат, который нужен.
Желаемый результат
Чтобы обработка каждого обновления проходила асинхронно, и чтобы каждое обновление не ожидало окончания обработки предыдущего, также и чтобы запрос за новым списком обновлений не ожидал окончания обработки всех предыдущих обновлений. Но необходимо, чтобы не было потрачено огромное количество ресурсов, которые ограничены, и чтобы был учтен большой объем обновлений, и чтобы работало раз и на всё время стабильно :) Проблема ещё в том, что обработка одного сообщения может занять как сотую часть секунды, так и секунд 5 в некоторых случаях, поэтому сделать всего потоков 10 не будет смысла.


Ответ

Как я указал в самом вопросе, мне посоветовали попробовать ExecutorService, и он действительно решил все проблемы.
Логика решения была примерно такой: простая итерация по списку полученных событий времени не занимает, а вот обработка некоторых (сообщений) занимает, и порой немало.
Значит, нужно пробегаться по массиву событий, и, если текущее событие оповещает о сообщении, отдавать его обработчику в новом потоке.
Можно было бы использовать FixedThreadPool, но тогда мне было не очень понятно, сколько потоков было бы необходимо создать и сколько их вообще нужно.
Как подсказала данная статья, есть ещё один хороший вариант — использовать CachedThreadPool. Он создаёт новые потоки, только если текущие заняты, и затем сам их подчищает. Таким образом, большое количество потоков будет создано только в случае большого количества сообщений, которые не успевали бы обрабатываться, но по окончании их обработки, все новые потоки будут "убиты".
Конкретная реализация:
Используем: ExecutorService service = Executors.newCachedThreadPool(); Отдаём наше "задание", которое требует обработки в новом потоке:
service.submit(() -> { // код, который будет выполнен в новом потоке });
Всё, вот так просто всё оказалось. Проверил на практике, среднее количество запущенных потоков стало равно 7-8, а в максимуме доходило до 30, а это вполне приемлемо.