/*
Musíme přenést data z nějakého zdroje k nějakému příjemci.
Zdroj vrací data v malých dávkách (~deset záznamů), zatímco příjemce pracuje efektivněji s většími dávkami (~tisíc záznamů).
Reálný příklad je přenos dat z front typu Kafka do databáze Clickhouse.
Zdroj:
- Téměř nekonečný.
- Zdroj nikdy nevrátí více než MaxItems záznamů v jednom volání Next.
- V rámci jedné "seance" (jedno volání funkce Pipe) zdroj při každém volání Next vrací nová data.
- Nicméně, po restartu začne zdroj od předchozí "potvrzené" pozice, určené cookie.
Proto *každá* hodnota cookie, kterou Next vrátí, po uložení dat do příjemce,
musí být potvrzena voláním Commit, ve stejném pořadí, v jakém byla vrácena Next.
Příjemce:
- Nemůže zpracovat více než MaxItems najednou.
Je třeba implementovat funkci func Pipe(p Producer, c Consumer) error,
která čte data ze zdroje, seskupuje je do bufferu o velikosti nejvýše MaxItems a ukládá je do příjemce,
po čemž potvrzuje pokrok ve zdroji.
*/
const MaxItems = 9999
type Producer interface {
// Next vrací:
// - dávku položek k zpracování
// - cookie k potvrzení po dokončení zpracování
// - chybu
Next() (items []any, cookie int, err error)
// Commit slouží k označení dávky dat jako zpracované
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}