Pular para o conteúdo

Streaming

Tópicos, partições e offsets — um log em disco, com a semântica do Kafka.

dataforge
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.

dataforge
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 mexeu

Dois 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#

dataforge
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#

dataforge
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#

dataforge
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#

dataforge
S.reter(corrente, "pedidos", 10)     // guarda os 10 últimos segmentos

Um 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#

dataforge
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.