Consumidores de Filas
Consumidores de filas processam mensagens de filas usando pools de workers.
Visão Geral
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
Configuração
| Opção | Padrão | Max | Descrição |
|---|---|---|---|
queue |
Obrigatório | - | ID do registro da fila |
func |
Obrigatório | - | ID do registro da função handler |
concurrency |
1 | 1000 | Quantidade de workers |
prefetch |
10 | 10000 | Tamanho do buffer de mensagens |
auto_ack |
false | - | Auto-ack no nível do driver (AMQP Consume autoAck; ignorado pelo driver de memória) |
driver_options |
{} |
- | Opções de consumidor específicas do driver |
Definição de 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
Função Handler
A função handler recebe o corpo depois que o codec da fila o decodifica. Use queue.message() para acessar a entrega atual e seus metadados:
-- 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
Confirmação
A menos que o handler conclua explicitamente a entrega, o consumidor usa o resultado da invocação da função:
| Resultado do handler | Ação | Efeito |
|---|---|---|
| Conclui sem erro de invocação | Ack | Mensagem removida da fila |
| Retorna ou gera um erro de invocação | Nack | A reentrega depende do driver |
Valores comuns de retorno, inclusive false, não determinam o comportamento de confirmação. Chame msg:ack() ou msg:nack() para concluir explicitamente. A conclusão acontece uma única vez: a primeira vence. Com auto_ack: true no AMQP, o broker confirma no momento da entrega; portanto, uma falha posterior do handler não pode causar reentrega pelo broker.
O handler pode resolver a mensagem por conta própria com queue.message() e msg:ack() / msg:nack(); o consumer então pula seu próprio ack/nack.
Pool de Workers
- 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
Exemplo
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
Encerramento Gracioso
Ao parar:
- Para de aceitar novas entregas
- Cancela contextos de workers
- Aguarda mensagens em voo (com timeout)
- Retorna erro de timeout se workers não terminarem
Declaração de Fila
# 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 | Descrição |
|---|---|
queue_name |
Sobrescreve nome da fila (padrão: nome do ID da entrada) |
codec |
Nome do codec de payload |
dead_letter.queue |
ID de registro aceito para uma fila dead-letter; não aplicado pelos drivers integrados |
dead_letter.max_attempts |
Contagem de tentativas aceita na configuração; não aplicada pelos drivers integrados |
driver_options |
Configurações específicas do driver indexadas por nome do driver |
Driver de Memória
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
Veja Também
- Fila de Mensagens - Referência do módulo de filas
- Configuração de Filas - Drivers de fila e definições de entrada
- Árvores de Supervisão - Ciclo de vida do consumidor
- Gerenciamento de Processos - Criação e comunicação de processos