Filas

Filas gerenciadas

O componente Filas da Zenifra usa Valkey atualmente e é indicado para tarefas assíncronas e workers. O padrão recomendado é Streams com consumer groups: as mensagens permanecem no stream, cada worker recebe uma parte do trabalho e o produtor não precisa conhecer os consumidores.

Fluxo recomendado com Streams

  1. O produtor adiciona a tarefa com XADD.
  2. O grupo é criado uma vez com XGROUP CREATE ... MKSTREAM.
  3. O worker lê mensagens novas com XREADGROUP.
  4. O worker processa a tarefa e só então confirma com XACK.
  5. Um processo de recuperação inspeciona XPENDING e usa XAUTOCLAIM para assumir mensagens de consumidores inativos.
  6. O stream é limitado com MAXLEN aproximado ou XTRIM para evitar crescimento sem limite.

Streams normalmente oferece entrega at-least-once. O processamento deve ser idempotente: uma tarefa pode ser entregue novamente depois de timeout, falha ou recuperação de um consumidor.

Node.js e TypeScript

npm install iovalkey
import Redis from 'iovalkey'
import { randomUUID } from 'node:crypto'

const url = new URL(process.env.VALKEY_URL ?? '')
const client = new Redis({
  host: url.hostname,
  port: Number(url.port),
  username: decodeURIComponent(url.username),
  password: decodeURIComponent(url.password),
  tls: { servername: url.hostname },
})

const stream = 'jobs:emails'
const group = 'email-workers'
const consumer = `worker-${randomUUID()}`

async function processEmail(fields) {
  console.log('process message', fields)
}

try {
  await client.xgroup('CREATE', stream, group, '$', 'MKSTREAM')
} catch (error) {
  if (!String(error).includes('BUSYGROUP')) throw error
}

await client.xadd(stream, 'MAXLEN', '~', 10000, '*', 'type', 'welcome', 'user_id', '42')
const batches = await client.xreadgroup('GROUP', group, consumer, 'COUNT', 10, 'BLOCK', 5000, 'STREAMS', stream, '>')

for (const [, messages] of batches ?? []) {
  for (const [id, fields] of messages) {
    await processEmail(fields)
    await client.xack(stream, group, id)
  }
}

await client.xtrim(stream, 'MAXLEN', '~', 10000)
await client.quit()

processEmail() representa o trabalho da aplicação e deve ser idempotente. Um processo separado deve consultar XPENDING e executar XAUTOCLAIM quando um consumidor ficar inativo.

Python

python -m pip install valkey
import os
import uuid
from urllib.parse import urlparse
from valkey import Valkey
from valkey.exceptions import ResponseError

def process_email(fields):
    print("process message", fields)

url = urlparse(os.environ["VALKEY_URL"])
client = Valkey(
    host=url.hostname, port=url.port, username=url.username, password=url.password,
    ssl=True, ssl_check_hostname=True, decode_responses=True,
)

stream = "jobs:emails"
group = "email-workers"
consumer = f"worker-{uuid.uuid4()}"

try:
    client.xgroup_create(stream, group, id="$", mkstream=True)
except ResponseError as error:
    if "BUSYGROUP" not in str(error):
        raise

client.xadd(stream, {"type": "welcome", "user_id": "42"}, maxlen=10000, approximate=True)
for _, messages in client.xreadgroup(group, consumer, {stream: ">"}, count=10, block=5000):
    for message_id, fields in messages:
        process_email(fields)
        client.xack(stream, group, message_id)

client.xtrim(stream, maxlen=10000, approximate=True)
client.close()

process_email() representa o trabalho da aplicação e deve ser idempotente. Um processo separado deve consultar XPENDING e usar XAUTOCLAIM para recuperar mensagens de consumidores inativos.

Pub/Sub, Lists e Streams

RecursoMelhor usoLimitação principal
Streamstarefas persistentes, vários workers e recuperaçãoexige consumer groups, XACK, retenção e idempotência
Listsfila simples com um consumidor por itemBLPOP remove antes do processamento; prefira BLMOVE quando a perda for relevante
Pub/Subnotificações ao vivo e fan-out efêmeromensagem não persistida; desconexão do consumidor perde o evento

Pub/Sub é suportado pelo serviço, mas não é a implementação principal do perfil Filas. Use PUBLISH/SUBSCRIBE para presença, invalidação ao vivo ou notificações que possam ser perdidas. Para tarefas de negócio, use Streams.

Limites e expectativas

  • Não há garantia exactly-once; use uma chave de idempotência na aplicação.
  • Defina retry, expiração, recuperação de pendências e uma estratégia de dead-letter usando outro stream quando necessário.
  • Limite o tamanho do stream com MAXLEN/XTRIM e monitore pendências.
  • Não confunda alta disponibilidade do plano com confirmação de cada mensagem em qualquer falha.
  • O produto não fornece exchanges, routing keys, dead-letter queue nativa ou particionamento automático de RabbitMQ/Kafka.

Próximos passos

Nessa página