/*
Meil tuleb andmed edastada ühest allikast ühe tarbijani.
Allikas tagastab andmeid väikestes osades (~kümme kirjet), samas kui tarbija töötab tõhusamalt suuremate osadega (~tuhat kirjet).
Reaalne näide on andmete edastamine Kafka-tüüpi järjekordadest Clickhouse andmebaasi.
Allikas:
- Peaaegu lõpmatu.
- Allikas ei tagasta kunagi rohkem kui MaxItems kirjet ühes Next-kutses.
- Ühes "sessioonis" (ühe Pipe funktsiooni kutses) tagastab allikas iga Next-kutses uusi andmeid.
- Kuid pärast taaskäivitust algab allikas eelmisest "kinnitust" saanud positsioonist, mis on määratud cookie-ga.
Seetõttu, *iga* cookie väärtus, mille Next tagastab, pärast andmete salvestamist vastuvõtjasse,
tuleb kinnitada Commit-kutsuga, samas järjekorras, nagu need on tagastatud Next poolt.
Vastuvõtja:
- Ei saa töödelda rohkem kui MaxItems korraga.
Vajalik on implementeerida funktsioon func Pipe(p Producer, c Consumer) error,
mis loeb andmeid allikast, rühmitab need mahuni, mis ei ületa MaxItems, ja salvestab vastuvõtjasse,
pärast mida kinnitab edusamme allikas.
*/
const MaxItems = 9999
type Producer interface {
// Next tagastab:
// - partii elemente töötlemiseks
// - cookie, mida tuleb kinnitada, kui töötlemine on lõppenud
// - viga
Next() (items []any, cookie int, err error)
// Commit kasutatakse andmepartii märkimiseks töödelduks
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}