Потребители очередей
Потребители очередей
Заголовок раздела «Потребители очередей»Потребители очередей обрабатывают сообщения с помощью пулов воркеров.
flowchart LR subgraph Consumer[Потребитель] QD[Драйвер очереди] --> DC[Канал доставки<br/>prefetch=10] DC --> WP[Пул воркеров<br/>concurrency] WP --> FH[Обработчик функции] FH --> AN[Ack/Nack] endКонфигурация
Заголовок раздела «Конфигурация»| Параметр | По умолчанию | Максимум | Описание |
|---|---|---|---|
queue | Обязательно | - | ID очереди в реестре |
func | Обязательно | - | ID функции-обработчика в реестре |
concurrency | 1 | 1000 | Количество воркеров |
prefetch | 10 | 10000 | Размер буфера сообщений |
auto_ack | false | - | Автоматически подтверждать до запуска обработчика |
driver_options | {} | - | Специфичные для драйвера опции потребителя |
Определение записи
Заголовок раздела «Определение записи»- name: order_consumer kind: queue.consumer queue: app:orders func: app:process_order concurrency: 5 prefetch: 20 lifecycle: auto_start: true depends_on: - app:ordersФункция-обработчик
Заголовок раздела «Функция-обработчик»Обработчик получает тело сообщения:
-- process_order.lualocal json = require("json")
local function handler(body) local order = json.decode(body)
-- Обработка заказа local result, err = process_order(order) if err then -- Возврат ошибки вызывает Nack (повторная постановка в очередь) return nil, err end
-- Успех вызывает Ack return resultend
return handler- name: process_order kind: function.lua source: file://process_order.lua modules: - jsonПодтверждение
Заголовок раздела «Подтверждение»| Результат | Действие | Эффект |
|---|---|---|
| Успех | Ack | Сообщение удаляется из очереди |
| Ошибка | Nack | Сообщение возвращается в очередь (зависит от драйвера) |
Пул воркеров
Заголовок раздела «Пул воркеров»- Воркеры работают как параллельные горутины
- Каждый воркер обрабатывает одно сообщение за раз
- Сообщения распределяются round-robin из канала доставки
- Буфер prefetch позволяет драйверу доставлять сообщения заранее
concurrency: 3prefetch: 10
Поток:1. Драйвер доставляет до 10 сообщений в буфер2. 3 воркера параллельно забирают из буфера3. По мере завершения воркеров буфер пополняется4. Backpressure, когда все воркеры заняты и буфер полонКорректное завершение
Заголовок раздела «Корректное завершение»При остановке:
- Прекращение приёма новых сообщений
- Отмена контекстов воркеров
- Ожидание обрабатываемых сообщений (с таймаутом)
- Возврат ошибки таймаута, если воркеры не завершились
Объявление очереди
Заголовок раздела «Объявление очереди»# Драйвер очереди (memory для разработки/тестов)- name: queue_driver kind: queue.driver.memory lifecycle: auto_start: true
# Определение очереди- name: orders kind: queue.queue driver: app:queue_driver queue_name: orders # Переопределить имя (по умолчанию: имя записи) codec: json # Кодек полезной нагрузки (опционально) dead_letter: # Обработка dead-letter (опционально) queue: app:dlq max_attempts: 5 driver_options: memory: max_length: 10000 # Memory-драйвер: ограниченный размер очереди| Поле | Описание |
|---|---|
queue_name | Переопределить имя очереди (по умолчанию: имя записи) |
codec | Имя кодека полезной нагрузки |
dead_letter.queue | ID очереди dead-letter в реестре |
dead_letter.max_attempts | Максимальное число попыток доставки до маршрутизации в DLQ |
driver_options | Специфичные для драйвера настройки, сгруппированные по имени драйвера |
Memory-драйвер
Заголовок раздела «Memory-драйвер»Встроенная in-memory очередь для разработки и тестирования:
- Тип:
queue.driver.memory - Сообщения хранятся в памяти
- Nack возвращает сообщение в конец очереди
- Не сохраняется между перезапусками
См. также
Заголовок раздела «См. также»- Очереди сообщений — справочник модуля Queue
- Конфигурация очередей — драйверы и определения записей
- Деревья супервизии — жизненный цикл потребителей
- Управление процессами — создание процессов и взаимодействие