Pular para o conteúdo

Consumidores de Filas

Consumidores de filas processam mensagens de filas usando pools de workers.

flowchart LR
subgraph Consumer[Consumidor]
QD[Driver de Fila] --> DC[Canal de Entrega<br/>prefetch=10]
DC --> WP[Pool de Workers<br/>concurrency]
WP --> FH[Handler de Função]
FH --> AN[Ack/Nack]
end
OpçãoPadrãoMaxDescrição
queueObrigatório-ID do registro da fila
funcObrigatório-ID do registro da função handler
concurrency11000Quantidade de workers
prefetch1010000Tamanho do buffer de mensagens
auto_ackfalse-Fazer Ack automaticamente antes de executar o handler
driver_options{}-Opções de consumidor específicas do 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

A função handler recebe o corpo da mensagem:

-- process_order.lua
local json = require("json")
local function handler(body)
local order = json.decode(body)
-- Processa o pedido
local result, err = process_order(order)
if err then
-- Retorna erro para disparar Nack (reenfileirar)
return nil, err
end
-- Sucesso dispara Ack
return result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
ResultadoAçãoEfeito
SucessoAckMensagem removida da fila
ErroNackMensagem reenfileirada (dependente do driver)
  • Workers executam como goroutines concorrentes
  • Cada worker processa uma mensagem por vez
  • Mensagens distribuídas round-robin do canal de entrega
  • Buffer de prefetch permite driver entregar antecipadamente
concurrency: 3
prefetch: 10
Fluxo:
1. Driver entrega até 10 mensagens para o buffer
2. 3 workers pegam do buffer concorrentemente
3. Conforme workers terminam, buffer reabastece
4. Contrapressão quando todos workers ocupados e buffer cheio

Ao parar:

  1. Para de aceitar novas entregas
  2. Cancela contextos de workers
  3. Aguarda mensagens em voo (com timeout)
  4. Retorna erro de timeout se workers não terminarem
# Driver de fila (memória para dev/teste)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# Definição de fila
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # Sobrescreve nome (padrão: nome da entrada)
codec: json # Codec de payload (opcional)
dead_letter: # Tratamento dead-letter (opcional)
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # Driver de memória: tamanho limitado da fila
CampoDescrição
queue_nameSobrescreve nome da fila (padrão: nome do ID da entrada)
codecNome do codec de payload
dead_letter.queueID de registro da fila dead-letter
dead_letter.max_attemptsNúmero máximo de tentativas de entrega antes de rotear para a DLQ
driver_optionsConfigurações específicas do driver indexadas por nome do driver
O roteamento dead-letter depende do driver. AMQP respeita a configuração DLX a nível de broker; o driver de memória não roteia para a DLQ.

Fila em memória embutida para desenvolvimento/testes:

  • Tipo: queue.driver.memory
  • Mensagens armazenadas em memória
  • Nack reenfileira a mensagem no final da fila
  • Sem persistência entre reinicializações