/*
Trebuie să transferăm date dintr-o sursă către un consumator.
Sursa oferă date în loturi mici (~zece înregistrări), în timp ce consumatorul funcționează mai eficient cu loturi mari (~o mie de înregistrări).
Un exemplu real este transferul de date din cozi de tip Kafka către o bază de date Clickhouse.
Sursa:
- Practic infinită.
- Sursa nu returnează niciodată mai mult de MaxItems înregistrări într-un singur apel Next.
- În cadrul unei "sesiuni" (un apel al funcției Pipe), sursa returnează date noi la fiecare apel Next.
- Cu toate acestea, după repornire, sursa va începe de la poziția "confirmată" anterioară, definită de cookie.
Prin urmare, *fiecare* valoare de cookie returnată de Next, după salvarea datelor în receptor,
trebuie confirmată cu un apel Commit, în aceeași ordine în care au fost returnate de Next.
Receptorul:
- Nu poate procesa mai mult de MaxItems odată.
Este necesar să implementați funcția func Pipe(p Producer, c Consumer) error
care citește date din sursă, le grupează într-un buffer de dimensiune cel mult MaxItems și le salvează în receptor,
după care confirmă progresul în sursă.
*/
const MaxItems = 9999
type Producer interface {
// Next returnează:
// - un lot de items de procesat
// - cookie pentru a fi confirmat când procesarea este terminată
// - eroare
Next() (items []any, cookie int, err error)
// Commit este folosit pentru a marca lotul de date ca procesat
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}