Входящая очередь событий

Выберите инструмент для разработки с AI-агентом:

  • используйте Битрикс24 Вайбкод, чтобы создать приложение для Битрикс24 по описанию задачи без знания языков программирования. Агент напишет код и разместит приложение на сервере без ручной настройки хостинга
  • используйте MCP-сервер, чтобы разрабатывать интеграцию через REST API в своем проекте. Агент будет обращаться к официальной REST-документации

Входящая очередь помогает принимать события REST API и исходящие вебхуки без перегрузки обработчика. Обработчик проверяет источник запроса, сохраняет данные в очереди и отвечает Битрикс24. Тяжелую бизнес-логику асинхронно выполняют отдельные воркеры. Такой подход дополняет общие рекомендации по производительности.

Основной принцип

  1. Запросы поступают на сервер. Когда Битрикс24 отправляет HTTP-запрос, он поступает на сервер или балансировщик нагрузки. Проверьте источник запроса по значению auth.application_token:

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

    Подробнее о проверке обоих источников читайте в статье Безопасность в обработчиках. Настройка исходящего вебхука описана на странице Входящие и исходящие вебхуки.

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

  3. Запросы извлекаются из очереди. Один или несколько обработчиков (воркеров) извлекают запросы и обрабатывают их по одному или параллельно, в зависимости от доступных ресурсов.

Битрикс24 не отправляет событие повторно, если обработчик не ответил или вернул аварийный HTTP-статус. Не возвращайте успешный ответ, пока запрос не сохранен в очереди.

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

Преимущества

  1. Управление нагрузкой. Очереди позволяют обрабатывать запросы последовательно и не перегружать сервер при всплесках трафика.
  2. Повышение стабильности. Надежная очередь сохраняет запросы при высокой нагрузке и позволяет повторить обработку после сбоя.
  3. Масштабируемость. Очереди позволяют масштабировать систему: добавлять больше воркеров или серверов для обработки запросов.
  4. Быстрый ответ Битрикс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 Нужен сохраняемый поток событий для нескольких групп потребителей Повторное чтение до фиксации смещения и отдельный топик для ошибок Поддерживается в пределах срока хранения Требуется управлять топиками, смещениями и группами потребителей

Продолжите изучение