Ir al contenido

Consumidores de Cola

Los consumidores de cola procesan mensajes de colas usando pools de workers.

flowchart LR
subgraph Consumer[Consumidor]
QD[Driver de Cola] --> DC[Canal de Entrega<br/>prefetch=10]
DC --> WP[Pool de Workers<br/>concurrency]
WP --> FH[Handler de Función]
FH --> AN[Ack/Nack]
end
OpciónPor DefectoMáxDescripción
queueRequerido-ID de registro de la cola
funcRequerido-ID de registro de la función handler
concurrency11000Cantidad de workers
prefetch1010000Tamaño del buffer de mensajes
auto_ackfalse-Hacer Ack automáticamente antes de ejecutar el handler
driver_options{}-Opciones de consumidor específicas del driver
- 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

La función handler recibe el cuerpo del mensaje:

-- process_order.lua
local json = require("json")
local function handler(body)
local order = json.decode(body)
-- Procesar la orden
local result, err = process_order(order)
if err then
-- Retornar error dispara Nack (requeue)
return nil, err
end
-- Éxito dispara Ack
return result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
ResultadoAcciónEfecto
ÉxitoAckMensaje removido de la cola
ErrorNackMensaje reencolado (dependiente del driver)
  • Los workers se ejecutan como goroutines concurrentes
  • Cada worker procesa un mensaje a la vez
  • Los mensajes se distribuyen round-robin desde el canal de entrega
  • El buffer de prefetch permite que el driver entregue por adelantado
concurrency: 3
prefetch: 10
Flujo:
1. El driver entrega hasta 10 mensajes al buffer
2. 3 workers extraen del buffer concurrentemente
3. A medida que los workers terminan, el buffer se rellena
4. Backpressure cuando todos los workers están ocupados y el buffer lleno

Al detener:

  1. Dejar de aceptar nuevas entregas
  2. Cancelar contextos de workers
  3. Esperar mensajes en vuelo (con timeout)
  4. Retornar error de timeout si los workers no terminan
# Driver de cola (memoria para dev/test)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# Definición de cola
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # Sobrescribir nombre (por defecto: nombre de entrada)
codec: json # Códec de payload (opcional)
dead_letter: # Manejo dead-letter (opcional)
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # Driver de memoria: tamaño acotado de cola
CampoDescripción
queue_nameSobrescribir nombre de cola (por defecto: nombre del ID de entrada)
codecNombre del códec de payload
dead_letter.queueID de registro de la cola dead-letter
dead_letter.max_attemptsMáximo de intentos de entrega antes de enrutar al DLQ
driver_optionsConfiguración específica del driver indexada por nombre del driver
El enrutamiento dead-letter depende del driver. AMQP respeta la configuración DLX a nivel de broker; el driver de memoria no enruta al DLQ.

Cola en memoria incorporada para desarrollo/pruebas:

  • Tipo: queue.driver.memory
  • Mensajes almacenados en memoria
  • Nack reencola el mensaje al final de la cola
  • Sin persistencia a través de reinicios