Pular para o conteúdo

Carga incremental

Marca d'água, reprocessamento e a idempotência que separa um pipeline de um script.

Ler tudo toda vez funciona até o dado crescer. A partir daí, o pipeline precisa saber até onde já leu — e precisa aguentar rodar duas vezes sem contar o mesmo dia duas vezes.

dataforge
adopt Arcane.Pipeline as Pipe
adopt Arcane.OS as OS
adopt Arcane.IO as IO

pasta := $"{OS.temp_dir()}/df-inc-{randint(100000, 999999)}"
IO.mkdir(pasta)
defer:
    IO.remove_tree(pasta)

// Um fluxo com ESTADO lembra entre execuções.
fluxo := Pipe.fluxo("vendas", $"{pasta}/estado.json")

// Na primeira execução, a marca é void — e isso é o "carregue tudo".
assert Pipe.marca(fluxo, "ultima_venda") is void

TODAS := [
    {"id": 1, "dia": "2026-01-05"},
    {"id": 2, "dia": "2026-01-06"},
    {"id": 3, "dia": "2026-01-07"},
]

action carregar(ctx):
    desde := Pipe.marca(fluxo, "ultima_venda") ?? ""
    novas := [v cycle v in TODAS given v["dia"] > desde]
    given len(novas) > 0:
        // A marca só vai ao disco no FIM da execução: se a etapa
        // seguinte falhar, o próximo run relê o mesmo lote.
        Pipe.marcar(fluxo, "ultima_venda", novas[len(novas) - 1]["dia"])
    yield len(novas)

Pipe.etapa(fluxo, "carregar", carregar)
r := Pipe.rodar(fluxo)
assert r["ok"] is yes
out $"primeira execução: {r['resultados']['carregar']} linha(s)"

// A segunda execução não relê nada.
r2 := Pipe.rodar(fluxo)
out $"segunda execução: {r2['resultados']['carregar']} linha(s)"
assert r2["resultados"]["carregar"] is 0

A marca só vale se ela for gravada DEPOIS#

Gravar a marca antes de a etapa seguinte terminar é o defeito clássico: a carga vai até o dia 7, a transformação falha, e o próximo run começa do dia 8 — o dia 7 nunca é processado, e ninguém descobre. Por isso ela só vai ao disco no fim da execução.

Reprocessar é um comando, e não um acidente#

dataforge
adopt Arcane.Pipeline as Pipe
adopt Arcane.OS as OS
adopt Arcane.IO as IO

pasta := $"{OS.temp_dir()}/df-inc2-{randint(100000, 999999)}"
IO.mkdir(pasta)
defer:
    IO.remove_tree(pasta)

fluxo := Pipe.fluxo("vendas", $"{pasta}/estado.json")
Pipe.etapa(fluxo, "ler", lambda ctx => 1)
Pipe.marcar(fluxo, "ate", "2026-01-07")
Pipe.rodar(fluxo)
assert Pipe.marca(fluxo, "ate") is "2026-01-07"

// 'esquecer_marca' é o botão de reprocessar — explícito, e não um
// efeito colateral de apagar um arquivo.
Pipe.esquecer_marca(fluxo, "ate")
assert Pipe.marca(fluxo, "ate") is void
out "a próxima execução relê tudo"

Idempotência: a partição inteira, e não o acréscimo#

A forma que funciona para reprocessar um dia é apagar a partição daquele dia e gravar de novo, e não acrescentar. Acrescentar duplica; sobrescrever a partição é idempotente por construção:

dataforge
adopt Arcane.Lago as Lago
adopt Arcane.OS as OS
adopt Arcane.IO as IO

raiz := $"{OS.temp_dir()}/df-lago-{randint(100000, 999999)}"
IO.mkdir(raiz)
defer:
    IO.remove_tree(raiz)

lago := Lago.lago(raiz)

action gravar_dia(dia, linhas):
    // Apagar a partição ANTES é o que torna o reprocessamento
    // idempotente: rodar duas vezes dá o mesmo resultado.
    Lago.remover_particao(lago, "vendas", {"dia": dia})
    Lago.acrescentar(lago, "vendas", linhas, ["dia"])
    yield len(linhas)

gravar_dia("2026-01-05", [{"dia": "2026-01-05", "v": 10}])
gravar_dia("2026-01-05", [{"dia": "2026-01-05", "v": 10}])   // de novo

lidas := Lago.ler(lago, "vendas", {"dia": "2026-01-05"})
assert len(lidas) is 1
out "rodou duas vezes, e há uma linha"
EstratégiaReprocessar éQuando
acrescentarduplicarsó para evento imutável com id próprio
sobrescrever a partiçãoidempotenteo padrão para dado por período
upsert por chaveidempotentequando a linha muda de valor depois
truncar e recarregaridempotente, e carotabela pequena de apoio