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

Очередь сообщений

Публикация и потребление сообщений из распределённых очередей. Поддержка нескольких бэкендов, включая RabbitMQ и другие AMQP-совместимые брокеры.

Настройку очередей см. в Queue.

local queue = require("queue")

Отправка сообщений в очередь по ID:

local ok, err = queue.publish("app:tasks", {
action = "send_email",
user_id = 456,
template = "welcome"
})
if err then
return nil, err
end
ПараметрТипОписание
queue_idstringИдентификатор очереди (формат: “namespace:name”)
dataanyДанные сообщения (таблицы, строки, числа, булевы)
headerstableОпциональные заголовки сообщения

Возвращает: boolean, error

Заголовки обеспечивают маршрутизацию, приоритезацию и трассировку:

queue.publish("app:notifications", {
type = "order_shipped",
order_id = order.id
}, {
priority = "high",
correlation_id = request_id
})

В консьюмере очереди доступ к текущему сообщению:

local msg, err = queue.message()
if err then
return nil, err
end
local msg_id = msg:id()
local priority = msg:header("priority")
local all_headers = msg:headers()

Возвращает: Message, error

Доступен только при обработке сообщений в контексте консьюмера.

МетодВозвращаетОписание
id()string, errorУникальный идентификатор сообщения
header(key)any, errorОдно значение заголовка (nil если отсутствует)
headers()table, errorВсе заголовки сообщения
ack()boolean, errorПодтвердить обработку (single-shot)
nack()boolean, errorСообщить о сбое для повторной доставки или dead-letter (single-shot)

Runtime автоматически выполняет ack при успехе обработчика и nack при ошибке. Вызывайте ack/nack только для раннего завершения.

local stats, err = queue.info("app:tasks")
-- stats может содержать: message_count, consumer_count, ready (зависит от драйвера)

Возвращает: table, error

Консьюмеры очередей определяются как точки входа, получающие payload напрямую:

entries:
- kind: queue.consumer
id: email_worker
queue: app:emails
method: handle_email
function handle_email(payload)
local msg = queue.message()
logger:info("Processing", {
message_id = msg:id(),
to = payload.to
})
local ok, err = email.send(payload.to, payload.template, payload.data)
if err then
return nil, err -- Сообщение будет возвращено в очередь или отправлено в dead-letter
end
end

Операции очереди подчиняются вычислению политики безопасности.

ДействиеРесурсОписание
queue.publish-Общее разрешение на публикацию сообщений
queue.publish.queueID очередиПубликация в конкретную очередь

Проверяются оба разрешения: сначала общее, затем для конкретной очереди.

УсловиеKindПовторяемо
Пустой ID очередиerrors.INVALIDнет
Пустые данные сообщенияerrors.INVALIDнет
Нет контекста доставкиerrors.INVALIDнет
Публикация не разрешенаerrors.INVALIDнет
Ошибка публикацииerrors.INTERNALнет

См. Обработка ошибок для работы с ошибками.