Класс Queue
Очередь — это коллекция типа «Первый вошёл, первый вышел» (First In, First Out или FIFO), которая позволяет работать только с самым первым значением. Итерация происходит от начала к концу с удалением взятого элемента.
Обзор классов
class Ds\Queue implements Ds\Collection , ArrayAccess <
const int MIN_CAPACITY = 8 ;
public allocate ( int $capacity ): void
public capacity (): int
public clear (): void
public isEmpty (): bool
public push ( mixed . $values ): void
public toArray (): array
Предопределённые константы
Ds\Queue::MIN_CAPACITY
Список изменений
| Версия | Описание |
|---|---|
| PECL ds 1.3.0 | Теперь класс реализует ArrayAccess . |
User Contributed Notes
There are no user contributed notes for this page.
- Copyright © 2001-2023 The PHP Group
- My PHP.net
- Contact
- Other PHP.net sites
- Privacy policy
Простой пример реализации очереди на PHP
В предыдущей статье я объяснил, что такое очередь, как она работает, и на абстрактном примере показал, что из себя представляет. В экосистеме PHP существует множество готовых реализаций клиентов очередей. Эта статья будет посвящена практической части работы с очередью. А использовать я буду простую, стабильную и быструю очередь. Сегодня вы увидите пример очереди задач на php, пример добавления заданий, удаления, обработки.
Для себя, вы можете представлять, что очередь является обычной базой данных, к которой можно выполнять множество запросов, к любым таблицам.
В предыдущей статье, в примере было 3 очереди: images , videos , emails . Потому, чтобы не выпадать из общей концепции статьи, реализуем эти очереди на практике.
Добавление задачи в очередь, мало чем отличается от INSERT -а в БД. Большинство очередей работают по принципу «первый прибыл, первый ушёл». То есть, сохраняя полную последовательность вызовов, сортируя по времени добавления задачи. Когда добавляется первая задача в очередь, она обрабатывается, потом, вторую, третью, и т.д. пока не выполнит все задачи из списка.
Отличие от базы данных — это то, что после выполнения, задача удаляется из очереди.
Одной из самых замечательных вещей является то, что очереди позволяют «резервировать» задание в очереди. Это значит, что, в случае, если очередь отслеживают два воркера, и один из них взял задание на выполнение, то, второй уже не может взять то же самое задание. Тут можно провести аналогию с «блокировкой файлов» — когда ставится флаг, который запрещает перезаписывать файл, пока его кто-то редактирует, или читает. Это важно понимать: одна задача выполняется одним воркером.
Что такое воркер?
Воркером может быть что угодно! Он может быть реализован на любом языке программирования (PHP, C, Ruby, Python), главное — чтобы он позволял подключаться, и общаться с самой очередью. Для PHP есть замечательная библиотека, которую я настоятельно рекомендую испробовать. Воркер — это интерфейс, который позволяет общаться с самой очередью. То есть, у вас должна быть установен драйвер очереди, и сама библиотека для работы с очередью.
Рассмотрим простой пример работы воркера:
use Pheanstalk\Pheanstalk; // подключимся к очереди $worker = new Pheanstalk('127.0.0.1'); // этот воркер работает над очередью изображений, и создаёт превью $worker->watch('images'); // поиск задач в очереди if ($job = $worker->reserve()) < // здесь будет код по обработке изображений >
Код выше просто показывает, как в PHP подключиться к очереди, используя библиотеку Pheanstalk, как настроить очередь images и доставать доступные задачи из очереди. Скрипт последовательно извлекает задачи, выполняет, и когда задач не останется — завершает своё выполнение.
Однако, в коде выше есть небольшой нюанс. Если запустить этот код, и в нём будет задание, то это задание будет циклично повторно выполняться. Эта задача никогда не будет удалена из очереди. Потому, чтобы удалить задачу из очереди изменим код:
use Pheanstalk\Pheanstalk; // подключаемся к очереди $worker = new Pheanstalk('127.0.0.1'); // этот воркер работает над очередью изображений, и создаёт превью $worker->watch('images'); // поиск задач в очереди if ($job = $worker->reserve()) < // здесь будет код по обработке изображений //после того, как задачу будет выполнена, удалить её из очереди $worker->delete($job); >
Запуск воркера
Для одноразового воркера достаточно обычного запуска скрипт (в консоли, или перехода по ссылке в браузере). Однако, этот вариант не подходит, потому что каждый раз вручную запускать скрипт — это, как минимум надоедает.
Потому, есть несколько вариантов.
Первый — можно поставить на крон (скажем, выполнять задачу каждые 5 минут). И каждые 5 минут будут обрабатываться все задачи, добавленные ранее в очередь. Этот вариант рабочий, однако, недостаток в том, что задачи будут выполняться не сразу, а только в течении 5 минут.
Второй варианта — это зациклить скрипт на бесконечное выполнение. Этот подход позволяет один раз запустить скрипт, который, в цикле, с определённым интервалом будет проверять новые задачи, и их выполнять.
Или же, более продвинутый вариант, и рекомендуемый вариант, используя supervisor. Ранее, я уже демонстрировал работу с supervisor-ом при работе с вебсокетами laravel. Его очень просто настраивать, и запускать. Проблема, которую он решает, заключается в том, что он следит за запущеными процессами, и, если какой-то из них прерывает своё выполнение, то он, тут же пытается перезапустить его. Тем самым, supervisor следит, чтобы скрипт был всегда запущен.
Как сделать бесконечный цикл while
Предыдущий пример был рабочий, но он выполнял только те задачи, которые были добавлены ранее. В этом разделе я напишу скрипт, который сможет выполнять задачи добавленные во время того, когда скрипт уже был запущен.
Теперь код будет выглядеть так:
use Pheanstalk\Pheanstalk; // подключимся к очереди $worker = new Pheanstalk('127.0.0.1'); // этот воркер работает над очередью изображений, и создаёт превью $worker->watch('images'); while (true) < // поиск задач в очереди if (!$job = $worker->reserve()) < // если задач пока нету, пропускаем итерацию // "засыпаем" на 1 секунду sleep(1); continue; >// здесь будет код по обработке изображений //после того, как задачу будет выполнена, удалить её из очереди $worker->delete($job); // возвращаемся в начало, к следующей задаче >
Это базовый пример, который должен был показать, что внутри очередей всё не так страшно, как может показаться.
Добавление задач в очередь
Предыдущие примеры были слишком абстрактными. Рассмотрим на примере, приближенном к боевым условиям. Представим, что пользователь присылает изображения, и нужно их принять и обработать:
use Pheanstalk\Pheanstalk; // подключимся к очереди $worker = new Pheanstalk('127.0.0.1'); // . // код сохранения изображений на сервер foreach ($_FILES as $image) < $data = [ 'path' =>$image['tmp_name'], 'name' => $image['name'], ]; // добавляем задачу на обработку изображения в очередь $pheanstalk ->useTube('images') //название выше созданной очереди images ->put(json_encode($data)); // полезные, данные, которые потребуются обработчику >
И добавив немного кода к базовому примеру, мы уже имеем рабочий скрипт, который сможет успешно сохранить изображения, и добавить их в очередь на обработку. Теперь, если пользователь загрузит даже 1000 изображений, то каждой из картинок будет создана задача в очереди.
А уже, внутри самого обработчика задачи, мы получим массив $data , которые был передан в метод put .
Данные в метод виде JSON-а, т.к. скрипт принимает строку. А JSON самый удобный формат передачи данных объектов, или массивов.
Получение данных, и обработка внутри воркера
Ввиду того, что мы получаем JSON в обработчике, то, его нужно декодировать. Для этого, напишем скрипт, который будет доставать переданную информацию:
use Pheanstalk\Pheanstalk; // подключимся к очереди $worker = new Pheanstalk('127.0.0.1'); // этот воркер работает над очередью изображений, и создаёт превью $worker->watch('images'); while (true) < // поиск задач в очереди if (!$job = $worker->reserve()) < // если задач пока нет, пропускаем итерацию // "засыпаем" на 1 секунду sleep(1); continue; >$data = json_decode($job->getDat(), true); // здесь будет код по обработке изображений //после того, как задачу будет выполнена, удалить её из очереди $worker->delete($job); // возвращаемся в начало, к следующей задаче >
В примере $data обрабатывается функцией json_decode() , потому что мы уверены, что данные приходят в формате JSON
Резюме
Api этой библиотеки содержит много функций, позволяющих работать с очередью, однако это уже выходит за рамки этой статьи. Основная цель, которую я преследовал — показать, как работать с очередью на PHP. Сегодня вы узнали, как работать с очередью, почему очереди стоит использовать, их преимущества, и, как добавить, и извлечь информацию из очереди.
Как всегда, для того, чтобы хорошо усвоить материал, нужно попробовать создать своими руками то, что мы проделали в этой статье. Пишите код, и понимание быстро придёт. Удачи!

В серці. Назавжди.
Вчора у мене помер однокласник. А сьогодні бабуся. І хто б міг уявити, що цей рік принесе війну, смерть товариша, та смерть члена сім’ї? Це боляче. Проте це добре нагадування про те, як швидко тече час. І як його ціна збільшується кожної марно витраченої секунди. І я не скажу щось
20 мая 2022 г. 1 min read

Ось такий він, руський мир
«Руський мир» — звучить дуже сильно та виправдовуюче. Гарна обгортка виправдання слабкості, аморальності та нікчемності своїх дійсних намірів. Руський мир, який дуже солодко звучить для всіх, хто хоче закрити очі на факт повномасштабної війни. Дуже добре виправдання вбивства для купки звірів. Втім, це ж росія, в якій все виглядає логічно
16 апр. 2022 г. 3 min read

Перехват запросов и ответов JavaScript Fetch API
Перехватчики — это блоки кода, которые вы можете использовать для предварительной или последующей обработки HTTP-вызовов, помогая в обработке глобальных ошибок, аутентификации, логирования, изменения тела запроса и многом другом. В этой статье вы узнаете, как перехватывать вызовы JavaScript Fetch API. Есть два типа событий, для которых вы можете захотеть перехватить HTTP-вызовы:
Как работать с очередью в php
Блог о разработке высокопроизводительных PHP приложений
Подписка

Очереди сообщений, AMQP, RabbitMQ
Использование очередей сообщений в архитектуре распределенных систем довольно распространенная техника разделения большой системы на компоненты.
Очереди, простой и в тот же час масштабируемый инструмент, позволяющий «подружить» независимые системы и научить их работать совместно.
Их задача предоставить возможность различным подсистемам обмениваться сообщениями обеспечивая маршрутизацию, гарантированную доставку и масштабирование.
Ниже пойдет речь о самых простых очередях сообщений построенных на основе БД, стандарте AMQP и отличной системе управления очередями RabbitMQ.
Кому нужны эти очереди?
Проблема соединения независимых компонентов в одно логическое целое появляется одновременно с ростом количества компонентов. При небольших нагрузках можно обойтись табличкой в базе данных или даже возможностями файловой системы. Но в какой-то момент сделанные напильником методы перестают работать или же масштабирование приносит больше головной боли чем результатов.
Давайте посмотрим на примере как можно создать примитивную очередь используя табличку в базе данных и проанализируем какие проблемы возникнут при необходимости увеличить активность раз так в 1000000.
Итак, поставим задачу: написать индексатор веб сайтов — систему которая будет ходить по ссылкам указанного сайта и сохранять заголовок/url страницы, а также картинки размером больше чем 10х10 учитывая возможность дубликатов.
Исходя из задачи систему можно разделить на два основных компонента: вебсканер и анализатор ресурсов, а последний в свою очередь может быть разбит на анализатор страниц и анализатор картинок.
Вебсканер — на входе получает url страницы, сканирует ее и возвращает внутренние ссылки и картинки, при этом передавая текущую страницу анализатору страниц, а картинки анализатору картинок соответственно.
Анализатор страницы — проверят нету ли страницы с аналогичным url в индексе, и если нет, получает заголовок (ограничимся содержимым тега title) и сохраняет.
Анализатор картинки — проверят нету ли картинки в индексе и соответствует ли размер требованиям. Если все критерии удовлетворительны изображение помещается в индекс.

Так как мы разрабатываем систему которая в будущем будет масштабироваться, нужно заложить возможность работы сервисов (вебсканнера и анализаторов) на разных серверах и возможно несколько инстансов одновременно (например 3 сканера и 5 анализаторов).
Начнем с написания сканера. Для получения ссылок и картинок пойдем простым путем и воспользуемся evaluate методом DOMXPath, который возвращает все теги по xpath путю.
class WebScanner extends Web < protected function getPageElements($tagName, $attributeName) < $elements = $this->getXPath()->evaluate('/html/body//'.$tagName); $elementsList = array(); for ($i = 0; $i < $elements->length; $i++) < $elementsList[] = $elements->item($i)->getAttribute($attributeName); > return $elementsList; > public function getAllLinks() < return $this->getPageElements('a', 'href'); > public function getAllImages() < return $this->getPageElements('img', 'src'); > >
Таким образом метод getPageElements получает все теги указанного типа и добавляет значение интересующего нас атрибута в возвращаемый массив. В результате мы можем получить все значения href ссылок и значения src картинок.
Так как классы-анализаторы тоже возможно будут работать с dom деревом, я вынес метод getXPath в базовый класс, который все остальные классы будут наследовать и таким образом унаследуют метод получения DOMXPath.
Перейдем к анализатору страницы. Его задача проверить наличие ссылки в базе данных и если если нужно добавить ссылку и title страницы.
class PageAnalyzer extends Web < public function analyze() < if (!Storage::getInstance()->containsPage($this->url)) < $title = $this->getDom()->getElementsByTagName("title")->item(0)->textContent; Storage::getInstance()->addPage($title, $this->url); > > >
Метод analyze проверяет есть ли текущая страница в индексе, если нет, то используя DomDocument получает title страницы и добавляет в базу данных.
Кроме анализатора появился класс для работы с базой данных Storage — простой синглтон-обертка mysqli.
class Storage < protected static $instance = null; protected $mysqli; protected static $PAGES_TABLE = 'pages'; protected static $IMAGES_TABLE = 'images'; public static function getInstance() < if (static::$instance == null) < static::$instance = new static(); >return static::$instance; > private function __construct() < $this->mysqli = new mysqli('127.0.0.1', 'root', '', 'blog_test'); $this->mysqli->set_charset('utf8'); > protected function containsUrl($tableName, $url) < return ($this->mysqli->query("SELECT COUNT(*) as cnt FROM ".$tableName." WHERE url='".$url."'")->fetch_object()->cnt > 0); > public function containsPage($url) < return $this->containsUrl(static::$PAGES_TABLE, $url); > public function containsImage($url) < return $this->containsUrl(static::$IMAGES_TABLE, $url); > public function addPage($title, $url) < $this->mysqli->query("INSERT INTO ".static::$PAGES_TABLE." (title,url)VALUES('".$title."','".$url."')"); > public function addImage($url) < $this->mysqli->query("INSERT INTO ".static::$IMAGES_TABLE." (url)VALUES('".$url."')"); > >
Работа с базой данных не является целью поста, поэтому его доводить до совершенства не будем, для нас достаточно возможности читать и писать в базу.
База данных (она же в нашем случае индекс) будет состоять из двух таблиц: pages — индекс страниц, images — индекс картинок. Для создания таблиц выполним:
CREATE TABLE `pages` ( `url` VARCHAR(255) NOT NULL, `title` VARCHAR(255) NOT NULL, PRIMARY KEY (`url`), KEY `title` (`title`) ) ENGINE=INNODB DEFAULT CHARSET=utf8; CREATE TABLE `images` ( `url` VARCHAR(255) NOT NULL, PRIMARY KEY (`url`) ) ENGINE=INNODB DEFAULT CHARSET=utf8;
И наконец анализатор картинок с единственным критерием проверки — картинка должна быть 10х10 или больше.
class ImageAnalyzer extends Web < protected static $ALLOWED_WIDTH = 10; protected static $ALLOWED_HEIGHT = 10; protected function isValid() < $size = getimagesize($this->url); return ($size[0] >= static::$ALLOWED_WIDTH && $size[1] >= static::$ALLOWED_HEIGHT); > public function analyze() < if (!Storage::getInstance()->containsImage($this->url)) < if ($this->isValid()) < Storage::getInstance()->addImage($this->url); > > > >
Если картинки нету в индексе, метод isValid проверяет размер и добавляет в базу при удовлетворении критерия.
Как вы наверно заметили оба класса наследуют класс Web, который содержит базовые методы для работы с данными:
abstract class Web < protected $url; protected $xPath = null; protected $dom = null; public function __construct($url) < $this->url = $url; > protected function getDom() < if ($this->dom == null) < $dom = new DOMDocument(); $dom->loadHTML(file_get_contents($this->url)); > return $dom; > protected function getXPath() < if ($this->xPath == null) < $this->xPath = new DOMXPath($this->getDom()); > return $this->xPath; > >
Все эти 3 компонента могут работать независимо на разных серверах общаясь между собой. Cканер страниц в процессе работы будет передавать ссылки анализатору страниц и картинки анализатору картинок, которые в свою очередь будут добавлять их в индекс.
Непосредственно момент «передачи» данных рассмотрим более детально.
Самый простой вариант который приходит в голову — создать таблицу queue, которая будет содержать два типа ссылок — картинки и страницы, что-то наподобие:
CREATE TABLE `queue` ( `url` varchar(255) NOT NULL, `type` enum('page','image') DEFAULT NULL, PRIMARY KEY (`url`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
При попадании на новую ссылку сканер будет добавлять ее в таблицу, откуда она позже будет считана анализаторами (и если нужно сканером), таким образом формируя «хаб» для объединения сервисов.
Но тут же получаем узкое, плохо масштабируемое место, которое может стать сильной проблемой при увеличении количества запросов. Шардинг решения построенного на основе БД будет довольно сложным из-за трудности реализации балансировки данных между шардами. А так как MySQL не поддерживает шардинг из коробки, эту логику прийдется включать в класс Storage, таким образом сильно усложняя его. И все это только для того чтобы научить компоненты общаться между собой.
Но не мы первые столкнулись с подобной проблемой) Именно для задач такого рода удобно использовать очереди сообщений.
Начнем с теории: AMQP в PHP
AMQP — протокол предназначенный для систем обмена сообщениями и описывающий характеристики и функции сообщений, очередей, роутинга, доступности и безопастности, а также схемы поведения сервера сообщений (брокера) и клиента (аналогично протоколам HTTP, FTP и др.).
Таким образом это некая спецификация того как должны общаться между собой брокер и клиент — формат и тип сообщений, метод передачи данных и т.д. AMQP совместимость избавляет от привязки клиента к конкретному брокеру, так как клиент поддерживающий AMQP может общаться з любым совместимым брокером. Поэтому разочаровавшись в RabbitMQ можно легко попробовать счастья с ActiveMQ, изменив только конфигурацию соединения с брокером.
В PHP функции для работы с AMQP совместимыми брокерами реализованы в виде 5 классов (расширение amqp, которое мы установили):
- AMQPConnection — класс для реализации соединения с брокером
- AMQPChannel — работа с каналами передачи данных
- AMQPExchange — отправка сообщений
- AMQPQueue — получение сообщений
- AMQPEnvelope — класс сообщения
Установим RabbitMQ
На убунте все делается как всегда просто:
-
Подключаем репозиторий (детально описано тут) и выполняем:
apt-get install rabbitmq-server
apt-get install pkg-config automake autoconf libsigc++-2.0-dev git clone git://github.com/alanxz/rabbitmq-c.git cd rabbitmq-c # Enable and update the codegen git submodule git submodule init git submodule update # Configure, compile and install autoreconf -i && ./configure && make && sudo make install
Если хотите установить web UI, выполняем:
rabbitmq-plugins enable rabbitmq_management
После чего установим само расширение:
pecl install amqp
Чтобы проверить что все работает напишем два тестовых скрипта: один будет отправлять сообщение в очередь, а другой считывать.
rabbitpost.php
$rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $testChannel = new AMQPChannel($rabbit); $testExchange = new AMQPExchange($testChannel); $testExchange->setName('amq.direct'); $testExchange->publish('Hello buddy!', 'route_to_everybody'); $rabbit->disconnect();
Итак, что же мы сделали? Сначала открыли соединение с брокером создав объект AMQPConnection и передав конструктору параметры соединения. Затем открыли новый канал(AMQPChannel), подключились к точке обмена сообщениями «amq.direct» и отправили сообщение с текстом «Hello buddy!» и ключом роутинга «route_to_everybody». Ключи используются для маршрутизации сообщений в разные очереди для обработки различными подсистемами.
rabbitread.php
$rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $channel = new AMQPChannel($rabbit); $q = new AMQPQueue($channel); $q->setName('direct_messages'); $q->declare(); $q->bind('amq.direct', 'route_to_everybody'); $envelope = $q->get(); if ($envelope) < print_r($envelope); $q->ack($envelope->getDeliveryTag()); > $rabbit->disconnect();
Аналогично предыдущему скрипту, мы создали новые соединение и канал, после чего создали новую очередь, которой указали проверять точку обмена «amq.direct» и доставлять все сообщения с ключом маршрутизации «route_to_everybody». Дальше мы пробуем получить сообщение из очереди, и в случае успеха выводим его в консоль.
Важным нюансом при работе с очередями сообщений является отправка acknowledge сообщений ack и nack. Ack с идентификатором доставленного сообщения отправляется в случае успешной обработки сообщения и значит, что сообщение может быть удалено с очереди, тогда как nack отправляется в случае когда обработчик по какой-то причине не справился и сообщение должно остаться в очереди (возможно будет успешно обработано другим обработчиком). До тех пор пока ack или nack не отправлен сообщение блокируется и переводится в состояние обработки — уже не отправляется другим клиентам, но еще и не удалено из очереди.
php rabbitpost.php
php rabbitread.php
Получим обьект сообщения:
AMQPEnvelope Object ( [body] => Hello buddy! [content_type] => text/plain [routing_key] => route_to_everybody [delivery_tag] => 1 [delivery_mode] => 0 [exchange_name] => amq.direct [is_redelivery] => 0 [content_encoding] => [type] => [timestamp] => 0 [priority] => 0 [expiration] => [user_id] => [app_id] => [message_id] => [reply_to] => [correlation_id] => [headers] => Array() )
Надеюсь тепер механика работы с очередями сообщений стала простой и понятной, попробуем реализовать на нашем примере с индексатором.
Для каждого типа сообщений создадим отдельную очередь:
- Ссылки которые нужно проиндексировать WebScanner’у
- Страницы которые нужно проанализировать PageAnalyzer’у
- Картинки которые нужно проанализировать ImageAnalyzer’у
Таким образом сканер при обнаружении непроиндексированной ссылки на другую страницу будет добавлять ее в очередь с ключом маршрутизации «scan», текущую, уже приндесированную страницу в «analyze_page», а если на странице найдены картинки, то все одни будут отправлены в «analyze_image».
Посмотрим что получилось.
include 'web.php'; include 'webscanner.php'; include 'storage.php'; error_reporting(E_ERROR); libxml_use_internal_errors(true); /** * Домен за рамки которого не нужно выходить (мы же не хотим весь интернет сканировать) */ $domain = 'phphighload.com'; $imageDomain = 'dropbox.com'; /** * Подключаемся к брокеру и точке обмена сообщениями */ $rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $channel = new AMQPChannel($rabbit); $queue = new AMQPExchange($channel); $queue->setName('amq.direct'); /** * Добавляем очередь откуда будем брать страницы, которые нужно проиндексировать */ $q = new AMQPQueue($channel); $q->setName('pages_to_scan'); $q->declare(); $q->bind('amq.direct', 'scan'); /** * Индексируем, пока в очереди не закончатся сообщения */ while ($page = $q->get()) < if (!is_object($page)) < continue; >$url = $page->getBody(); echo "Scanning: $url\n"; $scanner = new WebScanner($url); $links = $scanner->getAllLinks(); foreach ($links as $link) < /** * Если страница относится к указанному домену и еще не была проиндексирована -- добавляем ее в очередь на индексацию */ if (strpos($link, $domain) !== FALSE && !Storage::getInstance()->containsPage($link)) < $queue->publish($link, 'scan'); > > /** * Также если на странице есть картинки добавляем их для анализа */ $images = $scanner->getAllImages(); foreach ($images as $image) < /** * Если картинка относится к указанному домену и еще не была проанализирована -- добавляем ее в очередь */ if (strpos($image, $imageDomain) !== FALSE && !Storage::getInstance()->containsImage($image)) < $queue->publish($image, 'analyze_image'); > > /** * Текущую страницу тоже в очередь для анализации */ $queue->publish($url, 'analyze_page'); $q->ack($page->getDeliveryTag()); > $rabbit->disconnect();
Сканер подключается к очереди содержащей страницы ожидающие индексации и начинает их обработку. Если на странице обнаружена ссылка или картинка в рамках текущего домена, они добавляются в соответствующие очереди.
Теперь посмотрим как будут выглядеть PageAnalyzer и ImageAnalyzer :
analyze_page.php
include 'web.php'; include 'page_analyzer.php'; include 'storage.php'; error_reporting(E_ERROR); libxml_use_internal_errors(true); /** * Подключаемся к брокеру и точке обмена сообщениями */ $rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $channel = new AMQPChannel($rabbit); $queue = new AMQPExchange($channel); $queue->setName('amq.direct'); /** * Добавляем очередь откуда будем брать страницы */ $q = new AMQPQueue($channel); $q->setName('pages_to_scan'); $q->declare(); $q->bind('amq.direct', 'analyze_page'); /** * Обрабатываем пока в очереди не закончатся сообщения */ while (true) < $page = $q->get(); if ($page) < $url = $page->getBody(); echo "Parsing: $url\n"; $analyzer = new PageAnalyzer($url); /** * Если страница еще не была проанализирована, обрабатываем и добавляем в индекс */ $analyzer->analyze(); $q->ack($page->getDeliveryTag()); > else sleep(1); > $rabbit->disconnect();
analyze_image.php
include 'web.php'; include 'image_analyzer.php'; include 'storage.php'; /** * Подключаемся к брокеру и точке обмена сообщениями */ $rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $channel = new AMQPChannel($rabbit); $queue = new AMQPExchange($channel); $queue->setName('amq.direct'); /** * Добавляем очередь откуда будем брать страницы */ $q = new AMQPQueue($channel); $q->setName('images_to_scan'); $q->declare(); $q->bind('amq.direct', 'analyze_image'); /** * Обрабатываем пока в очереди не закончатся сообщения */ while (true) < $image = $q->get(); if ($image) < $url = $image->getBody(); echo "Checking: $url\n"; $analyzer = new ImageAnalyzer($url); /** * Если картинка еще не была проанализирована, обрабатываем и добавляем в индекс */ $analyzer->analyze(); $q->ack($image->getDeliveryTag()); > else sleep(1); > $rabbit->disconnect();
Аналогично сканеру — подключаемся к очереди и обрабатываем сообщения используя соответствующие классы.
Так как наш сканер обрабатывает только ссылки уже находящиеся в очереди, для того чтобы начать индексацию сайта необходимо записать в очередь стартовую страницу:
if (!isset($argv[1])) < die('Page url should be specified'); >$url = $argv[1]; /** * Подключаемся к брокеру и точке обмена сообщениями */ $rabbit = new AMQPConnection(array('host' => '127.0.0.1', 'port' => '5672', 'login' => 'guest', 'password' => 'guest')); $rabbit->connect(); $channel = new AMQPChannel($rabbit); $queue = new AMQPExchange($channel); $queue->setName('amq.direct'); /** * Добавляем стартовую страницу в очередь для индексирования */ $q = new AMQPQueue($channel); $q->setName('pages_to_scan'); $q->declare(); $q->bind('amq.direct', 'scan'); $queue->publish($url, 'scan');
Ну что, по-тестируем
Запустим в разных терминалах анализаторы:
php analyze_page.php
php analyze_image.php
Которые будут ждать появления новых элементов для сканирования в очереди.
И выполнив в третьем:
php init.php http://phphighload.com && php scanner.php
мы запустим цепную реакцию индексирования сайта phphighload.com и в результате получим список всех страниц и картинок в рамках домена.
Для тех кто устанавливал из официального репозитория и не поленился настроить UI, можно наблюдать за работой очереди перейдя по ссылке: http://rabbit_server_ip:15672/ (стандартные логин/пароль guest).

Если немного подождать, то проверив таблицу pages получим страницы сайта с названиями, а в images все картинки которые лежат на Дропбоксе размером больше 10х10.
Конечно же сам индексатор можно сделать более умным, чтобы он не зацикливался, не добавлял дубликаты в очередь для индексирования, но это тема другой статьи).
Такая распределенная архитектура хороша тем, что поддерживает горизонтальное масштабирование практически неограниченного размера. При необходимости всегда можно добавить новые сервера занимающиеся только анализом изображений или только сканированием страниц, соединенные в одну систему при помощи очереди сообщений, которая в свою очередь тоже легко масштабируется и с коробки умеет работать на нескольких серверах, быстро обрабатывая огромное количество сообщений.
Надеюсь у меня получилось показать насколько удобны и универсальны очереди для построения распределенных систем. Если у вас есть дополнения, замечания или коррекции, всегда рад пообщаться в комментариях).
P.S. Кому интересно запустить «индексатор» у себя, исходники примера выложил на github: https://github.com/kooler/phphighload.com/tree/master/rabbitmq
В тему:
- Стратегии масштабирования MySQL
- Горизонтальное масштабирование — балансировка нагрузки
- Шардинг в MySQL
Если пост понравился — нажмите на +1 — мне будет приятно.
Как сделать очередь на PHP по работе с CRM системой?
Всем привет!
В одной из популярных CRM систем, есть публичное приложение, которым пользуются 400+ аккаунтов. При изменении сущностей в CRM (сделка, контакт и .т.д) на наш сервер приходят вебхуки, при обработке которых нам необходимо опрашивать CRM аккаунт клиента через API. Ограничение API — 7 запросов в секунду. При превышении этого лимита, аккаунт блокируется для получения данных по API. Перед блокировкой CRM система конечно отправляет в заголовке ответа, что начинается превышение лимита запросов, но этот ответ ничем не помогает, так как во первых запрос уже сделан, во-вторых запросы от CRM системы все равно обрабатываются параллельно.
Проблема: если клиент импортирует в CRM сразу 10000 контактов, то CRM будет посылать к нам на сервер ~10-20 запросов в секунду о создании контакта. Так как нам нужно каждый запрос обработать, наш сервер тоже делает ~10-20 запросов в секунду по API уже в CRM. Так как ограничение запросов в 7 шт, наше приложение блокируется.
Вопрос: как можно решить проблему? Напрашивается очевидное решение — очереди, но познания в них у меня небольшие. На сколько знаю воркеров может быть фиксированное количество, например 10. А что если 30 аккаунтов CRM системы решат импортировать одновременно по 10000 контактов? А воркеров только 10. Остальные 20 будут ждать?
И есть ли в готовых решениях по работе с очередями динамическое создание воркеров? Например, когда есть процесс, который сканирует очереди. Если в очередях есть сообщения, запускает воркер на эту очередь. Как только воркер справился с очередью, процесс воркера убивается. Новый воркер будет запущен, если в очереди появились новые сообщения.
- Вопрос задан более трёх лет назад
- 508 просмотров
