Streaming
Tópicos, partições e offsets — um log em disco, com a semântica do Kafka.
adopt Arcane.Stream as S
corrente := S.corrente("eventos/")
S.topico(corrente, "pedidos", 4)
S.publicar(corrente, "pedidos", {"id": 1, "valor": 90}, "cliente-7")
cycle e in S.consumir(corrente, "pedidos", "faturamento"):
processar(e["valor"])
S.confirmar(corrente, "pedidos", "faturamento", e)Kafka é um cluster: réplicas, eleição de líder, coordenação entre máquinas. Arcane.Stream é um log em arquivo com a mesma semântica de tópico, partição e offset — o que cabe num processo, e que é onde a maioria dos fluxos de verdade começa.
Um log, não uma fila#
É a diferença que define o módulo. A fila entrega e esquece; o log guarda, e cada consumidor lembra onde parou.
faturamento := S.consumir(corrente, "pedidos", "faturamento") // 12
auditoria := S.consumir(corrente, "pedidos", "auditoria") // 12
S.confirmar_ate(corrente, "pedidos", "faturamento", faturamento)
S.consumir(corrente, "pedidos", "faturamento") // 0 — já processou
S.consumir(corrente, "pedidos", "auditoria") // 12 — não mexeuDois grupos leem o mesmo evento sem disputar. Numa fila, o primeiro a ler tira o evento do outro.
A partição é a unidade de ordem#
S.publicar(corrente, "pedidos", evento, "cliente-7")Eventos com a mesma chave caem sempre na mesma partição, e ali a ordem é garantida — os três eventos do cliente 7 chegam na ordem em que aconteceram.
Entre partições não há ordem, e é justamente isso que permite processar quatro em paralelo. Quem quer ordem total usa uma partição só, e paga com a serialização.
Voltar e reprocessar#
S.voltar(corrente, "pedidos", "faturamento", 0)É o que uma fila não permite, e o motivo de o log guardar o evento depois de entregue: quando a regra de processamento estava errada — e vai estar — dá para rodar tudo de novo.
O atraso é a métrica que se vigia#
S.atraso(corrente, "pedidos", "faturamento")
// {"total": 40213, "por_particao": {"p0": 12000, "p1": 9800, …}}Um atraso que só cresce significa que a produção passou o consumo — e o momento de agir é antes de o disco encher, não depois.
Retenção#
S.reter(corrente, "pedidos", 10) // guarda os 10 últimos segmentosUm log que só cresce enche o disco. A retenção apaga por segmento, não por evento: apagar o meio de um arquivo exigiria reescrevê-lo inteiro, e o offset dos que sobram mudaria — quebrando a posição de todo grupo.
Janelas#
cycle j in S.janela(eventos, 60):
out $"{j["inicio"]}: {j["quantos"]} eventos"É a operação que dá sentido a um fluxo: quantos por minuto é a pergunta que se faz, e ela não existe sem janela.
O que ele não faz#
- Sem réplica. O log vive num disco só.
- Sem transação entre tópicos. Publicar em dois é duas operações.
- Sem coordenação automática entre consumidores de um mesmo grupo — cada processo lê as partições que você mandar.
- *A garantia é ao menos uma vez. Um processo que cai entre processar e confirmar reprocessa o evento; exatamente uma vez* exige transação de ponta a ponta, e ninguém a tem de graça.