/*
Mums reikia perduoti duomenis iš šaltinio į vartotoją.
Šaltinis pateikia duomenis mažais paketais (~dešimt įrašų), o vartotojas efektyviau dirba su didesniais paketais (~tūkstantis įrašų).
Realiame pavyzdyje duomenys perduodami iš Kafka tipo eilės į Clickhouse duomenų bazę.
Šaltinis:
- Praktškai begalinis.
- Šaltinis niekada negrąžina daugiau kaip MaxItems įrašų vienu Next iškvietimu.
- Vienos "sesijos" (vieno Pipe funkcijos iškvietimo) metu šaltinis kiekvieną kartą grąžina naujus duomenis.
- Tačiau, perkraunant sistemą, šaltinis pradeda nuo ankstesnės "patvirtintos" pozicijos, nustatytos cookie.
Todėl, *kiekviena* cookie reikšmė, kurią grąžina Next, po duomenų išsaugojimo gavėjo, turi būti patvirtinta su Commit iškvietimu,
ta pačia tvarka, kaip ir grąžinta iš Next.
Gavėjas:
- Negali apdoroti daugiau kaip MaxItems vienu metu.
Reikia įgyvendinti funkciją func Pipe(p Producer, c Consumer) error,
kuri skaito duomenis iš šaltinio, juos grupuoja į buferį, kurio dydis ne didesnis kaip MaxItems, ir saugo jį gavėjui,
po to patvirtina pažangą šaltinyje.
*/
const MaxItems = 9999
type Producer interface {
// Next grąžina:
// - paketą elementų apdorojimui
// - cookie, kurį reikia patvirtinti, kai apdorojimas baigtas
// - klaidą
Next() (items []any, cookie int, err error)
// Commit naudojamas žymėti duomenų paketą kaip apdorotą
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}