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
- O produtor adiciona a tarefa com
XADD. - O grupo é criado uma vez com
XGROUP CREATE ... MKSTREAM. - O worker lê mensagens novas com
XREADGROUP. - O worker processa a tarefa e só então confirma com
XACK. - Um processo de recuperação inspeciona
XPENDINGe usaXAUTOCLAIMpara assumir mensagens de consumidores inativos. - O stream é limitado com
MAXLENaproximado ouXTRIMpara 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 iovalkeyimport 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 valkeyimport 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
| Recurso | Melhor uso | Limitação principal |
|---|---|---|
| Streams | tarefas persistentes, vários workers e recuperação | exige consumer groups, XACK, retenção e idempotência |
| Lists | fila simples com um consumidor por item | BLPOP remove antes do processamento; prefira BLMOVE quando a perda for relevante |
| Pub/Sub | notificações ao vivo e fan-out efêmero | mensagem 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/XTRIMe 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.