Перейти к содержимому

Потребители очередей

Потребители очередей обрабатывают сообщения с помощью пулов воркеров.

flowchart LR
subgraph Consumer[Потребитель]
QD[Драйвер очереди] --> DC[Канал доставки<br/>prefetch=10]
DC --> WP[Пул воркеров<br/>concurrency]
WP --> FH[Обработчик функции]
FH --> AN[Ack/Nack]
end
ПараметрПо умолчаниюМаксимумОписание
queueОбязательно-ID очереди в реестре
funcОбязательно-ID функции-обработчика в реестре
concurrency11000Количество воркеров
prefetch1010000Размер буфера сообщений
auto_ackfalse-Автоматически подтверждать до запуска обработчика
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.lua
local 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 result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
РезультатДействиеЭффект
УспехAckСообщение удаляется из очереди
ОшибкаNackСообщение возвращается в очередь (зависит от драйвера)
  • Воркеры работают как параллельные горутины
  • Каждый воркер обрабатывает одно сообщение за раз
  • Сообщения распределяются round-robin из канала доставки
  • Буфер prefetch позволяет драйверу доставлять сообщения заранее
concurrency: 3
prefetch: 10
Поток:
1. Драйвер доставляет до 10 сообщений в буфер
2. 3 воркера параллельно забирают из буфера
3. По мере завершения воркеров буфер пополняется
4. Backpressure, когда все воркеры заняты и буфер полон

При остановке:

  1. Прекращение приёма новых сообщений
  2. Отмена контекстов воркеров
  3. Ожидание обрабатываемых сообщений (с таймаутом)
  4. Возврат ошибки таймаута, если воркеры не завершились
# Драйвер очереди (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.queueID очереди dead-letter в реестре
dead_letter.max_attemptsМаксимальное число попыток доставки до маршрутизации в DLQ
driver_optionsСпецифичные для драйвера настройки, сгруппированные по имени драйвера
Маршрутизация dead-letter зависит от драйвера. AMQP учитывает конфигурацию DLX на уровне брокера; memory-драйвер не направляет сообщения в DLQ.

Встроенная in-memory очередь для разработки и тестирования:

  • Тип: queue.driver.memory
  • Сообщения хранятся в памяти
  • Nack возвращает сообщение в конец очереди
  • Не сохраняется между перезапусками