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

Очереди

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

flowchart LR
P[Publisher] --> D[Driver]
D --> Q[Queue]
Q --> C[Consumer]
C --> W[Worker Pool]
W --> F[Function]
  • Драйвер — реализация бэкенда (память, AMQP, SQS)
  • Очередь — логическая очередь, привязанная к драйверу
  • Консьюмер — связывает очередь с обработчиком, настраивает параллелизм
  • Пул воркеров — параллельная обработка сообщений

Несколько очередей могут использовать один драйвер. Несколько консьюмеров могут обрабатывать одну очередь.

ТипОписание
queue.driver.memoryДрайвер очереди в памяти
queue.driver.amqpДрайвер AMQP (RabbitMQ)
queue.driver.sqsДрайвер AWS SQS (также LocalStack, ElasticMQ)
queue.queueОбъявление очереди с привязкой к драйверу
queue.consumerКонсьюмер для обработки сообщений

Внутрипроцессный драйвер для разработки и однонодовых развёртываний. Без внешних зависимостей.

- name: memory_driver
kind: queue.driver.memory
lifecycle:
auto_start: true

Для RabbitMQ и брокеров, совместимых с AMQP 0-9-1.

- name: amqp_driver
kind: queue.driver.amqp
url: "amqp://guest:guest@localhost:5672/"
vhost: "/"
connection_name: "wippy-service"
heartbeat: "10s"
connection_timeout: "30s"
reconnect_delay: "1s"
reconnect_max_delay: "30s"
default_message_ttl: "1h"
default_queue_expiry: "24h"
prefetch_count: 10
lifecycle:
auto_start: true
ПолеТипПо умолчаниюОписание
urlstringamqp://guest:guest@localhost:5672/URL брокера
vhoststring-Переопределение виртуального хоста
connection_namestring-Идентификатор, отображаемый в UI брокера
auth_mechanismstringPLAINPLAIN, EXTERNAL (mTLS) или AMQPLAIN
heartbeatduration-Интервал keep-alive
connection_timeoutduration-Тайм-аут подключения
reconnect_delayduration1sНачальная задержка переподключения
reconnect_max_delayduration30sМаксимальная задержка переподключения
default_message_ttlduration-TTL сообщений по умолчанию для объявленных очередей
default_queue_ttlduration-TTL по умолчанию для объявленных очередей
default_queue_expiryduration-Срок жизни по умолчанию для объявленных очередей
prefetch_countint-Лимит prefetch на уровне канала
frame_sizeint-Лимит размера AMQP-фрейма
channel_maxint-Максимум каналов на соединение
tlsobject-Настройки TLS (см. ниже)

Блок TLS:

tls:
enabled: true
server_name: "rabbit.example.com"
cert_env: "AMQP_CLIENT_CERT"
key_env: "AMQP_CLIENT_KEY"
ca_env: "AMQP_CA_CERT"
insecure_skip_verify: false

Инлайновые поля cert/key/ca содержат PEM-контент; варианты *_env разрешаются через реестр env. Эти два источника взаимоисключающие для каждого поля. insecure_skip_verify отключает проверку сертификата (только для разработки).

Для AWS SQS и SQS-совместимых эндпойнтов (LocalStack, ElasticMQ). Учётные данные, регион и другие настройки AWS SDK поступают из общего ресурса config.aws.

- name: aws_config
kind: config.aws
region: us-east-1
access_key_id_env: app:AWS_ACCESS_KEY_ID
secret_access_key_env: app:AWS_SECRET_ACCESS_KEY
- name: sqs_driver
kind: queue.driver.sqs
config: app:aws_config
endpoint: "http://localhost:9324"
message_retention_period: 345600
default_delay_seconds: 0
lifecycle:
auto_start: true
ПолеТипПо умолчаниюОписание
configRegistry IDобязательноРесурс config.aws с регионом и учётными данными
endpointstring-Кастомный URL эндпойнта (LocalStack, ElasticMQ); опустить для реального AWS
message_retention_periodint345600 (4д)Срок хранения на уровне очереди в секундах (60–1209600)
default_delay_secondsint0Задержка доставки по умолчанию при CreateQueue (0–900)
disable_message_checksum_validationboolfalseОтключить проверку контрольных сумм SQS при отправке/приёме
use_fipsboolfalseИспользовать FIPS-совместимые эндпойнты
use_dual_stackboolfalseИспользовать dual-stack эндпойнты (IPv4 + IPv6)

Очереди создаются драйвером автоматически при первом использовании. Используйте заголовки с префиксом sqs.* для адресации SQS-специфичных атрибутов при публикации; нейтральные ключи вроде correlation_id и content_type по возможности транслируются в системные атрибуты SQS.

- name: tasks
kind: queue.queue
driver: app.queue:memory_driver
codec: json/plain
queue_name: "app_tasks"
driver_options:
memory:
max_length: 500
dead_letter:
queue: app.queue:tasks_dlq
max_attempts: 5
ПолеТипОбязательноОписание
driverRegistry IDДаДрайвер очереди
codecstringНетКодировка тел сообщений на проводе. По умолчанию json/plain (см. Кодеки)
queue_namestringНетВнешнее имя очереди (по умолчанию имя записи)
driver_optionsobjectНетПод-набор для каждого драйвера, ключ — kind драйвера
dead_letter.queueRegistry IDНетID очереди для неуспешных сообщений
dead_letter.max_attemptsintНетКоличество попыток до маршрутизации в DLQ

Ключи под driver_options сгруппированы по имени драйвера. Драйвер читает только свой под-набор — остальные ключи неактивны, что позволяет одной записи очереди объявлять настройки для нескольких драйверов при необходимости.

memory:

КлючОписание
max_lengthОграниченный размер буфера (0 = неограниченно)

amqp:

КлючОписание
durableПереживает перезапуск брокера
auto_deleteУдаляется при отключении последнего консьюмера
message_ttlПереопределение TTL сообщений на уровне очереди
queue_expiryСрок истечения для неиспользуемой очереди
max_lengthМаксимум хранимых сообщений

codec выбирает, как тело сообщения сериализуется перед передачей брокеру. Это строка формата payload, по умолчанию json/plain:

КодекФормат
json/plainJSON (по умолчанию)
application/msgpackMessagePack

AMQP-драйвер устанавливает соответствующий content-type (application/json или application/msgpack) на публикуемых сообщениях. Неизвестный кодек приводит к ошибке при объявлении очереди, а не во время публикации.

- name: task_consumer
kind: queue.consumer
queue: app.queue:tasks
func: app.queue:task_handler
concurrency: 4
prefetch: 20
auto_ack: false
driver_options:
amqp:
consumer_tag: "worker-1"
exclusive: false
lifecycle:
auto_start: true
depends_on:
- app.queue:tasks
ПолеПо умолчаниюОписание
queueобязательноID очереди в реестре
funcобязательноID функции-обработчика в реестре
concurrency1Количество параллельных воркеров
prefetch10Размер буфера на воркер
auto_ackfalseЕсли true, рантайм не вызывает ack брокера; успех/ошибка обработчика — единственный сигнал settle
driver_options-Под-набор для каждого драйвера (та же структура, что у очереди)

Опции консьюмера amqp:

КлючОписание
exclusiveЭксклюзивный доступ к очереди для одного консьюмера
no_localОтклонять сообщения, опубликованные на том же соединении
no_waitНе ждать подтверждения брокера при подписке
consumer_tagИдентификатор для этой подписки
Консьюмеры учитывают контекст вызова и могут подчиняться политикам безопасности. Настройте актёра и политики на уровне lifecycle. См. Безопасность.

Воркеры работают как параллельные горутины:

concurrency: 3, prefetch: 10
1. Драйвер доставляет до 10 сообщений в буфер
2. 3 воркера параллельно забирают из буфера
3. По мере завершения воркеров буфер пополняется
4. Обратное давление при занятых воркерах и полном буфере

Обработчики консьюмера получают декодированное тело сообщения первым аргументом. Используйте queue.message() для доступа к метаданным доставки (id, headers).

local queue = require("queue")
local logger = require("logger")
local function main(body)
local msg = queue.message()
logger:info("processing", {
id = msg:id(),
correlation_id = msg:header("correlation_id")
})
local ok, err = process_task(body)
if err then
return false -- nack: redelivery or DLQ
end
return true -- ack: remove from queue
end
return { main = main }
- name: task_handler
kind: function.lua
source: file://task_handler.lua
method: main
modules:
- queue
- logger

Рантайм автоматически фиксирует результат на основе возврата обработчика:

Результат обработчикаДействие
true или возврат, не равный falseAck
falseNack (повторная доставка или dead-letter в зависимости от драйвера)
Брошенная ошибкаNack

Вызывайте msg:ack() или msg:nack() явно только для досрочной фиксации. Фиксация однократна: побеждает первый сработавший вызов.

Когда на очереди настроен dead_letter, сообщение, которое получает nack сверх max_attempts, маршрутизируется в DLQ с заголовками x_dead_letter_reason и x_original_queue, устанавливаемыми драйвером. Издатели не должны устанавливать никакие заголовки x_* — они зарезервированы для учёта DLQ.

Из Lua-кода:

local queue = require("queue")
queue.publish("app.queue:tasks", {
id = "task-123",
action = "process",
data = payload
})

См. Модуль Queue для полного API.

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

  1. Прекращение приёма новых сообщений
  2. Отмена контекстов воркеров
  3. Ожидание завершения обрабатываемых сообщений (с тайм-аутом)
  4. Ошибка, если воркеры не успели завершиться