- RabbitMQ. Работа с очередями
- Подготовка
- Round-robin диспетчеризация
- Подтверждение сообщений
- Потерянное подтверждение
- Долговечность сообщения
- Обратите внимание на стойкость сообщений
- Справедливая отправка
- RabbitMQ tutorial 2 — Очередь задач
- Очереди задач
- Подготовка
- Циклическое распределение
- Подтверждение сообщений
- Устойчивость сообщений
- Равномерное распределение сообщений
- Ну а теперь все вместе
RabbitMQ. Работа с очередями
Перевод второго из шести уроков c официального сайта RabitMQ.
В первом уроке мы писали программу для отправки и получения сообщений из именованной очереди.
Основная идея Work Queues — рабочих очередей (так-же известные как: Task Queues)-избежать немедленного выполнения ресурсоемкой задачи и ждать ее завершения. Вместо этого мы планируем выполнить ресурсоемкую задачу позже. Мы инкапсулируем задачу в виде сообщения и отправляем его в очередь. Рабочий процесс, запущенный в фоновом режиме, будет запускать и выполнить поставленные задачи. И когда вы запустите множество воркеров, то задачи будут распределены и разделены между ними.
Эта концепция особенно полезна в web-приложениях, где невозможно обработать сложную задачу во время короткого срока HTTP запроса.
Подготовка
В предыдущей части этого урока мы отправили сообщение, содержащее » Hello World!». Теперь мы будем отправлять строки, которые будут имитировать сложные задачи. К сожалению у нас нет реально сложной задачи, такой как работа с изображениями или формирование крупного PDF файла, поэтому мы выполним имитацию сложной задачи с помощью функции sleep(). Мы будет определять сложность задачи по количеству точек. Каждая точка будет составлять одну секунду «работы«. Например: задача, которая отображает слово Hello. — займет три секунды.
Мы немного изменим код файла send.php из нашего предыдущего примера, чтобы реализовать отправку произвольных сообщений из командной строки.
Эта программа будет планировать задачи в нашей рабочей очереди. Дадим ей наименование: new_task.php.
Наш старый receive.php скрипт также требует некоторых изменений. В новый файл receive.php необходимо добавить функцию sleep() для имитации выполнения сложной задачи. Новый receive.php аналогично будет выводить сообщения из очереди и выполнять задачу. Дадим ему наименование: worker.php и добавим в него код:
Обратите внимание, что данная поддельная задача имитирует время выполнения.
Запустить задачу отправителя:
Round-robin диспетчеризация
Одним из преимуществ использования очереди задач является возможность легко распараллеливать работу.
Если мы создаем задержку в работе, то мы можем просто добавить больше воркеров и таким образом легко масштабироваться.
Давайте попробуем запустить 2 воркера worker.php. Для этого в 2х разных консолях выполните:
В третьей консоли у вас должен быть запущен отправитель new_task.php, где вы можете запустить на публикацию несколько сообщений:
После запуска публикаций сообщений откройте воркеров и посмотрите на них.
Вы должны увидеть приблизительно следующие:
- Воркер номер 1 —
По умолчанию RabbitMQ будет отправлять каждое сообщение к следующему получателю, в определенной последовательности. В среднем каждый получатель получит одинаковое количество сообщений. Такой способ распространения сообщений называется циклическим. Попробуйте выполните тоже самое, только с 3 или более отправителями.
Подтверждение сообщений
Выполнение задачи может занять несколько секунд. Вы можете задаться вопросом, что происходит в том случае, если один из получателей начинает выполнять долгую или сложную задачу и данная задача не выполняется или выполняется только частично? В текущем примере, как только RabbitMQ доставляет сообщение получателю, он сразу же помечает его для удаления. В таком случае, если вы убьете воркера, то потеряете сообщение, которое он только что обрабатывал. Так-же будут потеряны все сообщения, которые были отправлены этому конкретному воркеру, но еще не были обработаны.
Но я думаю, что мы бы не хотели потерять ни одного сообщения из очереди. И если воркер завершил свою работу, то мы бы хотели передать задачу другому воркеру.
Для того, чтобы убедиться, что сообщение никогда не будет потеряно. В RabbitMQ есть поддержка подтверждения сообщений .
Подтверждение (nowledgement) — отправляется получателем обратно, для того, чтобы сообщить RabbitMQ, что конкретное сообщение было получено, обработано и что RabbitMQ может удалить его.
Если получатель умирает (его канал закрыт, соединение закрыто или TCP-соединение потеряно) без отправки подтверждения, RabbitMQ поймет, что сообщение не было обработано полностью, и снова поставит его в очередь. Если есть доступ к другим получателям, то RabbitMQ быстро передаст задачу другому получателю. Таким образом, вы можете быть уверены, что сообщение не будет потеряно даже в том случае, если все получатели или их часть будет не доступна.
У сообщений отсутствуют тайм-ауты. И если получатель умрет, то RabbitMQ повторно добавит сообщение в очередь. Это нормально даже в том случае, если обработка сообщений занимает очень много времени.
По умолчанию подтверждения сообщений отключены. Пришло время включить их, установив четвертый параметр basic_consume в false (true означает отсутствие подтверждения) и отправить надлежащее подтверждение от воркера, как только мы закончим с задачей.
Используя этот код, мы можем быть уверены, что даже если вы отключите работника с помощью CTRL+C во время обработки сообщения, ничего не будет потеряно. Вскоре после смерти работника все неподтвержденные сообщения будут переданы повторно в очередь.
Подтверждение должно быть отправлено по тому же каналу, который получил доставку. Любая попытка подтвердить использование другого канала приведет к исключению протокола канального уровня. Для того, чтобы узнать больше см. руководство .
Потерянное подтверждение
Пропуск подтверждения, является распространенной ошибкой. Получить эту ошибку очень просто, но последствия могут быть очень серьезными. Сообщения будут повторно доставлены, когда ваш клиент завершит работу (это может выглядеть как случайная повторная поставка), но RabbitMQ будет использовать все больше и больше памяти, так как он не сможет выполнить обработку и доставку сообщений
Для отладки такого рода ошибок можно использовать rabbitmqctl для отображения информации messages_unacknowledged:
Тоже самое в Windows, Только без sudo
Долговечность сообщения
Только что мы узнали как защитить задачу от потери, после того как получатель умирает. Но наши задачи все равно будут потеряны после остановки или перезагрузки сервера RabbitMQ.
Когда RabbitMQ завершает работу или происходит сбой, то он забывает очереди и сообщения. Для того, чтобы никогда не потерять сообщения и очереди, нужно отметить очередь как долговечную.
Для этого нужно методу queue_declare третьим параметром передать значение true:
Хотя эта команда верна сама по себе,но она не будет работать при текущих настройках. Все этого потому, что очередь с именем hello уже определена как не долговечная. RabbitMQ не позволяет переопределить существующую очередь с разными параметрами и вернет ошибку любой программе, которая попытается это сделать. Но есть простой и быстрый способ обойти эту проблему. Для этого нужно переименовать текущую очередь. Давайте назовем нашу новую очередь как: task_queue.
Данный флаг, установленный в true, должен быть применен как к коду отправителя, так и к коду получателя.
На данный момент мы уверены, что очередь task_queue не будет потеряна, даже если RabbitMQ перезапустится. Теперь нам нужно отметить наши сообщения как постоянные, установив свойство delivery_mode = 2, которое AMQPMessage принимает как часть массива свойств.
Обратите внимание на стойкость сообщений
Справедливая отправка
Возможно, вы заметили, что диспетчеризация все еще работает не совсем так, как мы хотим. Например, в ситуации с двумя воркерами, один воркер будет постоянно занят, а другой почти не будет работать. Но RabbitMQ ничего об этом не знает и все равно будет отправлять сообщения равномерно в том же порядке.
Это происходит потому, что когда RabbitMQ отправляет сообщение в очередь, то не смотрит на количество неподтвержденных сообщений для получателя. Он просто слепо отправляет каждое n-е сообщение n-му потребителю.
Для того, чтобы избавиться от этой проблемы, мы должны воспользоваться методом basic_qos и передать ему параметр prefetch_count = 1. Это говорит RabbitMQ не давать более одного сообщения работнику за раз. Или, другими словам: «Не отправляйте новое сообщение работнику, пока он не обработает и не подтвердит предыдущее». Вместо этого он отправит его следующему работнику, который еще не занят.
Источник
RabbitMQ tutorial 2 — Очередь задач
В продолжение первого урока по изучению азов RabbitMQ публикую перевод второго урока с официального сайта. Все примеры, как и ранее, на python, но по-прежнему их можно реализовать на большинстве популярных ЯП.
Очереди задач
В первом уроке мы написали две программы: одна отправляла сообщения, вторая их принимала. В этом уроке мы создадим очередь, которая будет использоваться для распределения ресурсоемких задач между несколькими подписчиками.
Основная цель такой очереди ‒ не начинать выполнение задачи прямо сейчас и не ждать, пока оно завершится. Вместо этого задачи откладываются. Каждое сообщение соответствует одной задаче. Программа-обработчик, работающая в фоновом режиме, примет задачу на обработку, и через какое-то время она будет выполнена. При запуске нескольких обработчиков задачи будут разделены между ними.
Такой принцип работы особенно полезен для применения в веб-приложениях, где невозможно обработать ресурсоемкую задачу во время HTTP-запроса.
Подготовка
В предыдущем уроке мы отсылали сообщение с текстом «Hello World!». А сейчас мы будем посылать сообщения, соответствующие ресурсоемким задачам. Мы не будем выполнять реальные задачи, такие как изменение размера изображения или рендеринг pdf файла, давайте просто сделаем заглушку, используя функцию time.sleep(). Сложность задачи будет определяться количеством точек в строке сообщения. Каждая точка будет “выполняться” одну секунду. Например, задача с сообщением «Hello. » будет выполняться 3 секунды.
Мы немного изменим код программы send.py из предыдущего примера, чтобы было возможно отправлять произвольные сообщения из командной строки. Эта программа будет отправлять сообщения в нашу очередь, планируя выполнение новых задач. Назовем ее new_task.py:
Программа receive.py из предыдущего примера также должна быть изменена: необходимо симулировать выполнение полезной работы, по секунде на каждую точку текста сообщения. Программа будет получать сообщение из очереди и выполнять задачу. Назовем ее worker.py:
Циклическое распределение
Одно из преимуществ использования очереди задач ‒ возможность выполнять работу параллельно несколькими программами. Если мы не успеваем выполнять все поступающие задачи, то можем просто прибавить количество обработчиков.
Для начала давайте запустим сразу две программы worker.py. Обе они будут получать сообщения из очереди, но как именно? Сейчас увидим.
Вам необходимо открыть три окна терминала. В двух из них будет запущена программа worker.py. Это будут два подписчика ‒ C1 и C2.
В третьем окне мы будем публиковать новые задачи. После того как подписчики запущены, можно отправлять любое количество сообщений:
Давайте посмотрим, что было доставлено подписчикам:
По умолчанию, RabbitMQ будет передавать каждое новое сообщение следующему подписчику. Таким образом, все подписчики получат одинаковое количество сообщений. Такой способ распределения сообщений называется циклический [алгоритм round-robin]. Попробуйте то же самое с тремя или более подписчиками.
Подтверждение сообщений
Выполнение этих задач занимает несколько секунд. Возможно, Вы уже задались вопросом, что будет, если обработчик начал выполнение задачи, но неожиданно прекратил работу, выполнив ее лишь частично. В текущей реализации наших программ сообщение удаляется, как только RabbitMQ доставил его подписчику. Поэтому, если Вы остановите обработчик во время работы, задача не будет выполнена, а сообщение будет утеряно. Также будут утеряны доставленные сообщения, обработка которых еще не была начата.
Но мы не хотим терять какие-либо задачи. Нам нужно, чтобы в случае аварийного выхода одного обработчика сообщение передавалось другому.
Чтобы мы могли быть уверенны в отсутствии потерянных сообщений, RabbitMQ поддерживает подтверждение сообщений. Подтверждение (ack) отправляется подписчиком для информирования RabbitMQ о том, что полученное сообщение было обработано и RabbitMQ может его удалить.
Если подписчик прекратил работу и не отправил подтверждение, RabbitMQ поймет, что сообщение не было обработано, и передаст его другому подписчику. Так Вы можете быть уверены, что ни одно сообщение не будет потеряно, даже если выполнение программы-обработчика неожиданно прекратилось.
Для обработки сообщений отсутствует тайм-аут. RabbitMQ передаст их другому подписчику только если соединение с первым будет закрыто, поэтому нет никаких ограничений на время обработки сообщения.
По умолчанию используется ручное подтверждение сообщений. В предыдущем примере мы принудительно включили автоматическое подтверждение сообщений, указав no_ack=True. Теперь мы уберем этот флаг и будем отправлять подтверждение из обработчика сразу после выполнения задачи.
Теперь, даже если Вы остановите обработчика, нажав Ctrl+C во время обработки сообщения, ничто не будет утеряно. После остановки обработчика RabbitMQ заново передаст неподтвержденные сообщения.
Не забывайте подтверждать сообщения
Иногда разработчики забывают добавить в код basic_ack. Последствия этой небольшой ошибки могут быть существенными. Сообщение будет заново передано только тогда, когда программа-обработчик будет остановлена, но RabbitMQ будет потреблять все больше и больше памяти, т.к. не будет удалять неподтвержденные сообщения.
Для отладки такого рода ошибок Вы можете использовать rabbitmqctl для вывода на экран поля messages_unacknowledged (неподтвержденные сообщения):
[или воспользоваться более удобным скриптом мониторинга, который я приводил в первой части]
Устойчивость сообщений
Мы разобрались, как не потерять задачи, если подписчик неожиданно прекратил работу. Но задачи будут утеряны, если прекратит работу сервер RabbitMQ.
По умолчанию при остановке или падении сервера RabbitMQ все очереди и сообщения теряются, но это поведение можно изменить. Для того чтобы сообщения оставались в очереди после перезапуска сервера, необходимо сделать как очереди, так и сообщения устойчивыми.
Сначала убедимся, что не будет потеряна очередь. Для этого необходимо объявить ее, как устойчивую (durable):
Хотя эта команда сама по себе правильная, сейчас она не будет работать, потому что очередь hello уже объявлена, как неустойчивая. RabbitMQ не позволяет переопределить параметры для уже существующей очереди и вернет ошибку при попытке это сделать. Но есть простой обходной путь ‒ давайте объявим очередь с другим именем, например, task_queue:
Этот код необходимо исправить и для программы-поставщика, и для программы-подписчика.
Так мы можем быть уверены, что очередь task_queue не будет потеряна при перезапуске сервера RabbitMQ. Теперь необходимо пометить сообщения, как устойчивые. Для этого нужно передать свойство delivery_mode со значением 2:
Замечание по поводу устойчивости сообщений
Пометка сообщения, как устойчивого, не дает гарантии, что сообщение не будет утеряно. Не смотря на то, что это и заставляет RabbitMQ сохранять сообщение на диск, есть небольшой промежуток времени, когда RabbitMQ подтвердил принятие сообщения, но еще не записал его. Также RabbitMQ не делает fsync(2) для каждого сообщения, поэтому какие-то из них могут быть сохранены в кеш, но еще не записаны на диск. Гарантия устойчивости сообщений не полная, но ее более чем достаточно для нашей очереди задач. Если Вам требуется более высокая надежность, Вы можете оборачивать операции в транзакции.
Равномерное распределение сообщений
Вы, возможно, заметили, что распределение сообщений все еще работает не так, как нам нужно. Например, при работе двух подписчиков, если все нечетные сообщения содержат сложные задачи [требуют много времени на выполнение], а четные ‒ простые, то первый обработчик будет постоянно занят, а второй большую часть времени будет свободен. Но RabbitMQ об этом ничего не знает и все-равно будет передавать сообщения подписчикам по очереди.
Так происходит, потому что RabbitMQ распределяет сообщения в тот момент, когда они попадают в очередь, и не учитывает количество неподтвержденных сообщений у подписчиков. RabbitMQ просто отправляет каждое n-ое сообщение n-ому подписчику.
Для того чтобы изменить такое поведение, мы можем использовать метод basic_qos с опцией prefetch_count=1. Это заставит RabbitMQ не отдавать подписчику единовременно более одного сообщения. Другими словами, подписчик не получит новое сообщение, до тех пор пока не обработает и не подтвердит предыдущее. RabbitMQ передаст сообщение первому освободившемуся подписчику.
Замечание по поводу размера очереди
Если все подписчики заняты, то размер очереди может увеличиваться. Следует обращать на это внимание и, возможно, увеличить количество подписчиков.
Ну а теперь все вместе
Полный код программы new_task.py:
Полный код программы worker.py:
Используя подтверждение сообщений и prefetch_count, вы можете создать очередь задач. Настройка устойчивости позволит задачам сохраняться даже после перезапуска сервера RabbitMQ.
В третьем уроке мы разберем, как можно отправить одно сообщение нескольким подписчикам.
Источник