Входящая очередь событий
Выберите инструмент для разработки с AI-агентом:
- используйте Битрикс24 Вайбкод, чтобы создать приложение для Битрикс24 по описанию задачи без знания языков программирования. Агент напишет код и разместит приложение на сервере без ручной настройки хостинга
- используйте MCP-сервер, чтобы разрабатывать интеграцию через REST API в своем проекте. Агент будет обращаться к официальной REST-документации
Входящая очередь помогает принимать события REST API и исходящие вебхуки без перегрузки обработчика. Обработчик проверяет источник запроса, сохраняет данные в очереди и отвечает Битрикс24. Тяжелую бизнес-логику асинхронно выполняют отдельные воркеры. Такой подход дополняет общие рекомендации по производительности.
Основной принцип
-
Запросы поступают на сервер. Когда Битрикс24 отправляет HTTP-запрос, он поступает на сервер или балансировщик нагрузки. Проверьте источник запроса по значению
auth.application_token:- для события приложения сравните параметр с токеном, который приложение сохранило при установке
- для исходящего вебхука сравните параметр со значением поля Токен приложения в настройках вебхука
Подробнее о проверке обоих источников читайте в статье Безопасность в обработчиках. Настройка исходящего вебхука описана на странице Входящие и исходящие вебхуки.
-
Запросы добавляются в очередь. Сохраните данные запроса в базе данных или очереди сообщений. Отправляйте успешный HTTP-ответ только после того, как данные надежно сохранены. После ответа выполняйте тяжелую бизнес-логику асинхронно.
-
Запросы извлекаются из очереди. Один или несколько обработчиков (воркеров) извлекают запросы и обрабатывают их по одному или параллельно, в зависимости от доступных ресурсов.
Битрикс24 не отправляет событие повторно, если обработчик не ответил или вернул аварийный HTTP-статус. Не возвращайте успешный ответ, пока запрос не сохранен в очереди.
Внутренняя очередь может передать запрос повторно после сбоя воркера. Сделайте обработку идемпотентной, чтобы повторный запуск не изменил данные дважды.
Преимущества
- Управление нагрузкой. Очереди позволяют обрабатывать запросы последовательно и не перегружать сервер при всплесках трафика.
- Повышение стабильности. Надежная очередь сохраняет запросы при высокой нагрузке и позволяет повторить обработку после сбоя.
- Масштабируемость. Очереди позволяют масштабировать систему: добавлять больше воркеров или серверов для обработки запросов.
- Быстрый ответ Битрикс24. Обработчик подтверждает прием запроса сразу после сохранения в очереди и не ждет завершения бизнес-логики.
Варианты реализации
Выбор способа зависит от требований к надежности, повторному чтению и эксплуатации. Для входящей очереди можно использовать базу данных, Redis или брокер сообщений.
1. Сохранение в базу данных
База данных подходит, если не нужна отдельная инфраструктура для очереди. Обработчик сохраняет тело запроса и метаданные в таблице, а воркер выбирает записи со статусом pending, обрабатывает их и обновляет статус. При высокой нагрузке производительность может снижаться из-за конкуренции за доступ к таблице.
Пример на PHP:
В примерах $requestData содержит проверенные данные входящего запроса, а $db — настроенное подключение PDO. Структуру таблицы и обработку исключений адаптируйте к своему проекту.
// Сохраняем запрос в базу данных и только после этого отправляем успешный HTTP-ответ
$statement = $db->prepare(
'INSERT INTO requests (data, status) VALUES (:data, :status)'
);
$statement->execute([
'data' => json_encode($requestData, JSON_THROW_ON_ERROR),
'status' => 'pending',
]);
// Резервируем один запрос, чтобы другой воркер не взял его одновременно
$db->beginTransaction();
$request = $db->query(
"SELECT * FROM requests WHERE status = 'pending' " .
"ORDER BY id LIMIT 1 FOR UPDATE SKIP LOCKED"
)->fetch();
if ($request) {
$statement = $db->prepare(
'UPDATE requests ' .
'SET status = :status, processing_started_at = CURRENT_TIMESTAMP ' .
'WHERE id = :id'
);
$statement->execute([
'status' => 'processing',
'id' => $request['id'],
]);
}
$db->commit();
if ($request) {
processRequest(json_decode($request['data'], true));
$statement = $db->prepare(
'UPDATE requests ' .
'SET status = :status, processing_started_at = NULL ' .
'WHERE id = :id'
);
$statement->execute([
'status' => 'processed',
'id' => $request['id'],
]);
}
Конструкция FOR UPDATE SKIP LOCKED требует поддержки со стороны системы управления базами данных. Добавьте в таблицу поле processing_started_at. Если запрос остается в статусе processing дольше допустимого времени обработки, отдельная фоновая задача должна вернуть его в статус pending.
2. Использование Redis
Для повышения производительности можно использовать Redis как простую очередь.
- Входящие запросы ставятся в очередь в Redis через команду
LPUSH(добавить элемент в список). - Обработчики (воркеры) с помощью
BRPOPLPUSHатомарно переносят запросы в отдельный список обработки. После успешной обработки запрос удаляется из этого списка.
Redis быстро обрабатывает большие объемы данных и подходит для высокой нагрузки.
Пример на PHP:
В примере $redis — настроенное подключение к Redis, а $requestData содержит проверенные данные входящего запроса.
// Сохраняем запрос в Redis и только после этого отправляем успешный HTTP-ответ
$redis->lPush('request_queue', json_encode($requestData));
// Атомарно переносим сообщение в список обработки
while ($request = $redis->brPopLPush(
'request_queue',
'request_processing',
5
)) {
processRequest(json_decode($request, true));
$redis->lRem('request_processing', $request, 1);
}
Если воркер завершится до удаления запроса из request_processing, верните зависший запрос в основную очередь отдельной фоновой задачей. Храните время переноса в обработку отдельно, чтобы задача не вернула запрос, который еще выполняется.
3. RabbitMQ или Kafka
RabbitMQ и Kafka требуют отдельной инфраструктуры, но предоставляют встроенные механизмы доставки и повторного чтения.
- RabbitMQ. Выбирайте для очереди задач, когда сообщения нужно распределять между воркерами и подтверждать их обработку
- Kafka. Выбирайте для потока событий, который нужно хранить и независимо читать несколькими группами потребителей
RabbitMQ повторно доставляет сообщение, если соединение с воркером закрылось до подтверждения. Подтверждайте сообщение только после успешного выполнения processRequest(). Ограничьте число повторов, а сообщения с постоянной ошибкой переносите в очередь недоставленных сообщений.
При работе с Kafka фиксируйте смещение только после успешной обработки. Ограничьте число повторов, а сообщения с постоянной ошибкой переносите в отдельный топик.
Пример использования RabbitMQ на PHP:
Пример использует классы AMQPStreamConnection и AMQPMessage из библиотеки php-amqplib. Переменная $requestData содержит проверенные данные входящего запроса.
// Отправляем сообщение в очередь и только после подтверждения публикации возвращаем успешный HTTP-ответ
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
$connection = new AMQPStreamConnection('localhost', 5672, 'user', 'password');
$channel = $connection->channel();
$channel->queue_declare('request_queue', false, true, false, false);
$channel->confirm_select();
$message = new AMQPMessage(
json_encode($requestData, JSON_THROW_ON_ERROR),
['delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT]
);
$channel->basic_publish($message, '', 'request_queue');
$channel->wait_for_pending_acks_returns();
// Обрабатываем сообщения из очереди
use PhpAmqpLib\Connection\AMQPStreamConnection;
$connection = new AMQPStreamConnection('localhost', 5672, 'user', 'password');
$channel = $connection->channel();
$channel->queue_declare('request_queue', false, true, false, false);
$callback = function($message) {
processRequest(json_decode($message->body, true));
$message->ack();
};
$channel->basic_consume('request_queue', '', false, false, false, false, $callback);
while($channel->is_consuming()) {
$channel->wait();
}
Как выбрать вариант
Выбирайте очередь по требованиям к обработке и эксплуатации, а не только по объему запросов.
| Вариант | Когда выбирать | Обработка сбоя | Повторное чтение | Эксплуатация |
|---|---|---|---|---|
| База данных | Не нужна отдельная инфраструктура, запросы обрабатывает один или несколько воркеров | Статусы записей и возврат зависших запросов в pending |
Возможно, пока записи хранятся в таблице | Требуется контролировать блокировки и очищать обработанные записи |
| Redis | Нужна быстрая внутренняя очередь с простой моделью обработки | Список обработки и возврат зависших запросов в основную очередь | Нужно реализовать самостоятельно | Требуется отдельный сервер Redis и контроль памяти |
| RabbitMQ | Нужны очередь задач, маршрутизация и подтверждение обработки | Повторная доставка неподтвержденных сообщений и очередь недоставленных сообщений | Ограничено политикой хранения очереди | Требуется настроить брокер, подтверждения и правила повторов |
| Kafka | Нужен сохраняемый поток событий для нескольких групп потребителей | Повторное чтение до фиксации смещения и отдельный топик для ошибок | Поддерживается в пределах срока хранения | Требуется управлять топиками, смещениями и группами потребителей |