Consumers de queue
Los consumers de queue entregan mensajes de una queue a handlers de funciones mediante un pool de workers configurable.
Resumen
flowchart LR
subgraph Consumer
QD[Queue Driver] --> DC[Delivery Channel
prefetch=10]
DC --> WP[Worker Pool
concurrency]
WP --> FH[Function Handler]
FH --> AN[Ack/Nack]
end
Configuración
| Opción | Predeterminado | Máximo | Descripción |
|---|---|---|---|
queue |
Requerido | - | ID de registro de la cola |
func |
Requerido | - | ID de registro de la función handler |
concurrency |
1 | 1000 | Cantidad de workers |
prefetch |
10 | 10000 | Tamaño del buffer de mensajes |
auto_ack |
false | - | Auto-ack a nivel de driver (AMQP Consume autoAck; ignorado por el driver de memoria) |
driver_options |
{} |
- | Opciones de consumidor específicas del driver |
Definición de la entrada
- name: order_consumer
kind: queue.consumer
queue: app:orders
func: app:process_order
concurrency: 5
prefetch: 20
lifecycle:
auto_start: true
requires:
- app:orders
Función handler
La función handler recibe el body después de que el codec de la queue lo decodifique. Usa queue.message() para acceder al delivery actual y sus metadatos:
-- process_order.lua
local queue = require("queue")
local logger = require("logger")
local function main(order)
local msg, msg_err = queue.message()
if msg_err then
return nil, msg_err
end
logger:info("processing order", {
message_id = msg:id(),
order_id = order.id
})
return {processed = true, order_id = order.id}
end
return {main = main}
- name: process_order
kind: function.lua
source: file://process_order.lua
method: main
modules:
- queue
- logger
Acknowledgment
A menos que el handler liquide explícitamente el delivery, el consumer usa el resultado de invocar la función:
El handler puede resolver el mensaje él mismo con queue.message() y msg:ack() / msg:nack(); el consumidor omite entonces su propio ack/nack.
Pool de Workers
Los valores de retorno ordinarios, incluido false, no eligen el comportamiento de acknowledgment. Llama a msg:ack() o msg:nack() para liquidarlo explícitamente. La liquidación es single-shot: gana la primera. Con AMQP auto_ack: true, el broker confirma al entregar, por lo que un fallo posterior del handler no puede provocar redelivery del broker.
Pool de workers
- Los workers se ejecutan como goroutines concurrentes.
- Cada worker procesa un mensaje a la vez.
- Los workers toman mensajes de un delivery channel compartido. El siguiente worker libre recibe el siguiente mensaje, sin orden o rotación garantizados entre workers.
- El buffer de prefetch permite al driver entregar mensajes antes de procesarlos.
Ejemplo
concurrency: 3
prefetch: 10
Flow:
1. Driver delivers up to 10 messages to buffer
2. 3 workers pull from buffer concurrently
3. As workers finish, buffer refills
4. Backpressure when all workers busy and buffer full
Apagado ordenado :id=shutdown-ordenado
Durante el apagado, el consumidor:
- Deja de aceptar deliveries nuevos.
- Cancela los contextos de workers.
- Espera los handlers in-flight hasta el stop timeout.
- Devuelve un error de timeout si los workers no terminan.
Declaración de la queue
# Queue driver (memory for dev/test)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# Queue definition
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # Override name (default: entry name)
codec: json/plain # Payload codec (optional; json/plain is the default)
dead_letter: # Accepted configuration; not enforced by built-in drivers
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # Memory driver: bounded queue size
| Campo | Descripción |
|---|---|
queue_name |
Sobrescribe el nombre de la queue (default: nombre del ID de entrada) |
codec |
Nombre del codec del payload |
dead_letter.queue |
ID de registro aceptado para una dead-letter queue; los drivers integrados no lo aplican |
dead_letter.max_attempts |
Número de intentos aceptado en configuración; los drivers integrados no lo aplican |
driver_options |
Settings específicos del driver, agrupados por nombre de driver |
Driver en memoria
El driver integrado en memoria está pensado para desarrollo y pruebas:
- Su kind es
queue.driver.memory. - Los mensajes se almacenan en memoria.
- Nack intenta volver a encolar al final una copia del mensaje; ese intento puede fallar cuando la queue limitada está llena.
- Los mensajes no persisten entre reinicios.
Véase también
- Message Queue — Referencia del módulo Queue
- Configuración de queue — Drivers y definiciones de entradas
- Supervisión — Ciclo de vida del consumer
- Gestión de procesos — Creación y comunicación de procesos