Sobes.tech
Middle+

/* 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 }