Sobes.tech
Middle+

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