/*
Dobbiamo trasferire dati da una sorgente a un consumatore.
La sorgente fornisce i dati in piccoli batch (~dieci record), mentre il consumatore lavora meglio con batch più grandi (~mille record).
Un esempio reale è il trasferimento di dati da code di tipo Kafka a un database Clickhouse.
Sorgente:
- Quasi infinita.
- La sorgente non restituisce mai più di MaxItems record in una singola chiamata a Next.
- All'interno di una "sessione" (una chiamata alla funzione Pipe), la sorgente restituisce nuovi dati ad ogni chiamata a Next.
- Tuttavia, dopo un riavvio, la sorgente ricomincia dalla posizione "confermata" precedente, definita da cookie.
Pertanto, *ogni* valore di cookie restituito da Next, dopo aver salvato i dati nel ricevitore,
deve essere confermato con una chiamata a Commit, nello stesso ordine in cui sono stati restituiti da Next.
Ricevitore:
- Non può elaborare più di MaxItems alla volta.
È necessario implementare la funzione func Pipe(p Producer, c Consumer) error
che legge i dati dalla sorgente, li raggruppa in un buffer di dimensione non superiore a MaxItems e li salva nel ricevitore,
dopodiché conferma il progresso nella sorgente.
*/
const MaxItems = 9999
type Producer interface {
// Next restituisce:
// - un batch di elementi da elaborare
// - cookie da confermare al termine dell'elaborazione
// - errore
Next() (items []any, cookie int, err error)
// Commit viene usato per marcare il batch di dati come processato
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}